fix: quiesce cancelled session reconciliation

This commit is contained in:
Hypatia May
2026-07-24 21:27:07 +08:00
parent cbc6d81fc3
commit 224e00b2bd
20 changed files with 359 additions and 31 deletions

View File

@@ -2,5 +2,5 @@
# side as of the last confirmed-consistent state. Both languages carry equal authority;
# after editing either side, bring the other along and re-record with:
# pnpm run verify-translation-pairing --write
2026-07-23-unified-session-query-service.md: 0a466e1c36ff1796c858666b0eb36bbd0f480bb0
2026-07-23-unified-session-query-service.zh.md: 448122b8e6951058b9f633cd56112b0391e1912e
2026-07-23-unified-session-query-service.md: 676a42017ca42f9e649f6529f84787e7162faac0
2026-07-23-unified-session-query-service.zh.md: d4449a415840d61cbb10f88def1062a13e556749

View File

@@ -16,6 +16,8 @@ The interface package already owns the shared record, filter, trace, search-requ
`SessionQuerySqlite` extends that service and is the sole concrete backend. One mounted instance therefore exposes every operation through `ctx.sessionQuery`; its inherited exact operations use the shared corpus implementation, while its SQLite-owned lifecycle observes sources, reconciles the derived FTS index, ranks matches, and owns cursor generations. The interface package has no standalone concrete plugin, search-provider registry, or second context key.
SQLite reconciliation is one quiescent serialized state machine. It passes the caller's exact abort signal into durable snapshot listing and inspection, awaits each started backend operation itself, and checks cancellation after every await and before starting the next source or index operation. Cancellation therefore cannot release the serializer while an ignored or cooperative backend call is still cleaning up, and it cannot start a subsequent listing, inspection, reconciliation, or query after the signal is observed.
Backend configuration includes the inherited `readWindowMax` setting alongside its own index path, journal mode, page limits, and snippet limit. First-party apps that need session queries mount the SQLite backend and place its disposable index beside their configured persistence root.
This service topology supersedes the separate-key portion of the [exact query decision](../feature/2026-07-10-session-query-service.md) and [SQLite search decision](../feature/2026-07-10-sqlite-session-query-provider.md); their corpus, query, tokenizer, reconciliation, and safety decisions remain in force.
@@ -32,4 +34,6 @@ Consumers inject one service and can combine exact and full-text operations with
The unified object deliberately retains two internal observation strategies: exact operations read authoritative live/persisted sources per call, while full-text operations reconcile a disposable index. Sharing the context key does not make the derived index authoritative or couple exact-read availability to an FTS query.
Queued cancellation remains prompt. Cancellation during active asynchronous source observation waits for that started operation to settle, which makes rejection a quiescence boundary and preserves single-file execution for a following search. Synchronous SQLite statements remain non-preemptible and are bracketed by signal checks.
Unit coverage pins inherited and abstract behavior on one key, SQLite coverage exercises both operation families on the concrete backend, and the real Loader path verifies that one exported plugin registers the combined service.

View File

@@ -16,6 +16,8 @@ Status: implemented
`SessionQuerySqlite` 扩展该服务,并且是唯一的具体后端。因此,一个挂载实例便可通过 `ctx.sessionQuery` 暴露全部操作;其继承的精确操作使用共享的语料库实现,而由 SQLite 管理的生命周期负责观察数据源、对齐派生 FTS 索引、对匹配项排序并管理游标代际。接口包不提供独立的具体插件、搜索提供方注册表或第二个上下文键。
SQLite 的对齐过程是一个具备静止性保证的串行状态机。它将调用方的原始中止信号传给持久化快照列表与检查操作,直接等待每个已经启动的后端操作,并在每次等待后以及启动下一个数据源或索引操作前检查是否已取消。因此,即使后端忽略取消或正在配合清理,串行器也不会提前释放;观察到中止信号后,也不会再启动后续的列表、检查、对齐或查询操作。
后端配置除了自身的索引路径、日志模式、分页限制与文本片段长度上限外,还包含继承的 `readWindowMax` 设置。需要会话查询的第一方应用挂载 SQLite 后端,并将其可丢弃索引放在已配置的持久化根目录旁。
这一服务拓扑取代了[精确查询决策](../feature/2026-07-10-session-query-service.md)和 [SQLite 搜索决策](../feature/2026-07-10-sqlite-session-query-provider.md)中关于分离上下文键的部分;其中关于语料库、查询、分词器、对齐与安全性的决策仍然有效。
@@ -32,4 +34,6 @@ Status: implemented
统一后的对象有意保留两种内部观察策略:精确操作在每次调用时读取权威的实时源或持久化源,全文操作则使可丢弃索引与数据源对齐。共用上下文键不会让派生索引成为权威来源,也不会使精确读取的可用性依赖 FTS 查询。
排队阶段的取消仍会及时生效。在异步数据源观察已经开始后取消时,调用方会等待该操作完成清理后才收到拒绝;因此拒绝本身构成静止边界,并保证后续搜索仍按单一串行流程执行。同步 SQLite 语句无法在执行中被抢占,服务会在其前后检查中止信号。
单元测试在同一个键上同时固定继承实现与抽象方法的契约SQLite 测试在具体后端上覆盖两类操作,真实 Loader 路径则验证单个导出的插件能够注册组合后的服务。

View File

@@ -32,13 +32,13 @@ Both persistent and live FTS5 tables use `unicode61`. The implementation experim
The shared extractor includes message text, reasoning, nested tool-call/result content, tool names and arguments, blocked-prompt reasons, todo status/content, and error or terminal status detail. Structural boundaries, stream chunks, request headers, successful completion markers, and unknown declaration-merged event/content variants produce no document. Surface classification reuses `foldSurface()` so search agrees with model-history derivation.
One serialized operation reads the provider-neutral `SessionPersistence` snapshot listing, compares each source-qualified opaque revision with the revision stored beside the indexed session, loads only new or changed logs, reconciles rows in one transaction, and executes the query. It never calls the backend's mutating `load()` for an id currently owned by `ctx.sessions`; the TEMP overlay records persisted availability, and the durable base refreshes after the live owner detaches. A revision identifies its backing persistence store as well as the backend-local log revision, so reopening against the same store reuses indexed rows while switching to an independent store cannot collide on a session id and local counter. Observation repeats when listing changes during a load; this incorporates a mutating load repair's refreshed revision before commit. Repeated queries and unchanged reopen load no full persisted logs. New, changed, and deleted sessions update on the next stable search. A source or extraction failure cannot mark a row current, and a transaction failure rolls back so a later search retries.
One serialized operation reads the provider-neutral `SessionPersistence` snapshot listing, compares each source-qualified opaque revision with the revision stored beside the indexed session, loads only new or changed logs, reconciles rows in one transaction, and executes the query. It passes the caller's exact abort signal into snapshot listing and non-mutating inspection, directly awaits every started backend operation, and checks cancellation after each await and before starting more work. Cancellation therefore rejects only after active backend work is quiescent, starts no subsequent observation or reconciliation step, and keeps a following search serialized behind cleanup even if a backend ignores the signal. The operation never calls the backend's mutating `load()` for an id currently owned by `ctx.sessions`; the TEMP overlay records persisted availability, and the durable base refreshes after the live owner detaches. A revision identifies its backing persistence store as well as the backend-local log revision, so reopening against the same store reuses indexed rows while switching to an independent store cannot collide on a session id and local counter. Observation repeats when listing changes during a load; this incorporates a mutating load repair's refreshed revision before commit. Repeated queries and unchanged reopen load no full persisted logs. New, changed, and deleted sessions update on the next stable search. A source or extraction failure cannot mark a row current, and a transaction failure rolls back so a later search retries.
Persisted documents survive restarts. Live sessions use connection-local TEMP tables, shadow the persisted base for the same id, and reveal that base on detach. Closing the database drops live rows. Unmounting persistence hides durable rows without treating absence as authoritative deletion; remounting observes and reconciles the backend again. Conflicting immutable live and durable headers fail rather than combining sources.
The derived schema has its own application id and monotonic schema version. A recognized incompatible version resets only this derived database. A database with a foreign application id or unrecognized user tables is refused before journal-mode mutation, which prevents an accidentally configured canonical session database from being changed. On POSIX filesystems, missing directories and database files are created owner-only so new SQLite sidecars inherit that mode; existing modes are preserved. One service in one process exclusively owns a derived-index path; cross-process writers are unsupported because generations and live TEMP shadow state are connection-owned.
Cancellation rejects queued operations and caller waits around asynchronous source observation without committing an aborted observation. Node's synchronous `DatabaseSync` MATCH call cannot be interrupted once it is executing on the JavaScript thread, so the service checks the signal at serialized boundaries but does not promise mid-statement preemption.
Cancellation rejects queued operations promptly. Once asynchronous source observation starts, the caller waits for that backend promise to settle before rejection, without committing an aborted observation or starting more source/index work. Node's synchronous `DatabaseSync` metadata and MATCH calls cannot be interrupted once executing on the JavaScript thread, so the service checks the signal around those calls but does not promise mid-statement preemption.
## Alternatives considered
@@ -52,6 +52,6 @@ Cancellation rejects queued operations and caller waits around asynchronous sour
Search has a small provider-neutral API while its only backend owns every derived-index state transition. The separate database adds configuration and a lightweight snapshot read before queries, but index corruption, reset, and tokenizer changes cannot endanger canonical logs. Durable revisions avoid full-log reads and rewrites for unchanged sessions; TEMP live overlays preserve current-session truth without making uncheckpointed events durable.
The chosen tokenizer supports short tokens with a smaller index but does not promise substring recall. Literal phrases make query syntax safe and predictable at the cost of excluding boolean/full MATCH expressions. Cancellation is effective while queued or awaiting sources, but synchronous SQLite execution remains a non-preemptible section.
The chosen tokenizer supports short tokens with a smaller index but does not promise substring recall. Literal phrases make query syntax safe and predictable at the cost of excluding boolean/full MATCH expressions. Cancellation is prompt while queued and quiescent while awaiting sources; synchronous SQLite execution remains a non-preemptible section bracketed by signal checks.
Unit coverage pins extraction, filters, both search scopes, all default surfaces, metadata-before-ranking, snippets, literal escaping, deterministic ties, complete pagination, scoped cursor invalidation, dynamic persistence mount/unmount, restart reconciliation, live shadow/reveal/reopen, schema safety, rollback retry, and queued/in-flight source-wait cancellation. A keyless real-Loader-path test combines the package with the real SQLite persistence backend.

View File

@@ -960,9 +960,10 @@ abstract list(signal?: AbortSignal): Promise<SessionHeader[]>
* successful mutating {@link load} repair changes the next listed revision.
* Revisions also distinguish independently backed stores so backend-local
* counters cannot compare equal across different persistence sources.
* @param signal - optional cancellation for backend snapshot-listing work.
* @returns one header and opaque revision per materialized session without loading full logs.
*/
abstract listSnapshots(): Promise<SessionPersistenceSnapshot[]>
abstract listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]>
```
Types: [SessionEvent](../core-data-structures/core.md) · [SessionHeader](../core-data-structures/persistence.md) · [SessionId](../core-data-structures/core.md) · [SessionLocation](../core-data-structures/persistence.md) · [SessionPersistenceSnapshot](../core-data-structures/persistence.md)

View File

@@ -126,7 +126,7 @@ interface SessionPersistenceSnapshot {
## The backends
Both implement the same abstract `SessionPersistence` (locate/create/append/load/inspect/list/listSnapshots over `SessionEvent`) and pass `runPersistenceContract`, proving the seam is genuinely backend-agnostic:
Both implement the same abstract `SessionPersistence` (locate/create/append/load/inspect/list/listSnapshots over `SessionEvent`, with optional cancellation on observation methods) and pass `runPersistenceContract`, proving the seam is genuinely backend-agnostic:
- **[dsh-session-persistence-jsonl](../../packages/session-persistence/session-persistence-jsonl)** — an append-only logical JSONL log per session, stored as checksummed concatenated Zstandard frames by default or raw lines by configuration, with crash-safe atomic writes, interrupted-turn recovery, and a read/replay path.
- **[dsh-session-persistence-sqlite](../../packages/session-persistence/session-persistence-sqlite)** — `node:sqlite`, one row per `SessionEvent`. The row shape `(session_id, seq, type, time, data, source_event_seqs, surface_op)` maps 1:1 onto the event, including optional surface metadata, so there is no parallel persisted schema to keep in sync.

View File

@@ -477,8 +477,8 @@ export const SERVICE_API: readonly ServiceApiEntry[] = [
jsDoc: '/**\n * Lightweight listing from metadata, without a full-log parse.\n * @param signal - optional cancellation for backend listing work.\n * @returns one header per materialized session.\n */',
},
{
signature: 'abstract listSnapshots(): Promise<SessionPersistenceSnapshot[]>',
jsDoc: '/**\n * List materialized sessions with cheap per-log change tokens.\n *\n * Repeated observations of an unchanged log return the same revision. A\n * successful mutating {@link load} repair changes the next listed revision.\n * Revisions also distinguish independently backed stores so backend-local\n * counters cannot compare equal across different persistence sources.\n * @returns one header and opaque revision per materialized session without loading full logs.\n */',
signature: 'abstract listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]>',
jsDoc: '/**\n * List materialized sessions with cheap per-log change tokens.\n *\n * Repeated observations of an unchanged log return the same revision. A\n * successful mutating {@link load} repair changes the next listed revision.\n * Revisions also distinguish independently backed stores so backend-local\n * counters cannot compare equal across different persistence sources.\n * @param signal - optional cancellation for backend snapshot-listing work.\n * @returns one header and opaque revision per materialized session without loading full logs.\n */',
},
],
},

View File

@@ -39,7 +39,7 @@ A root belongs to one encoding. Startup discovery and targeted lookup reject the
- **Crash recovery — preserve valid tail work.** `load` validates every complete compressed frame and scans their decompressed JSONL. If the last frame is structurally incomplete, the reader keeps its complete decoded records, truncates from that frame's start, and re-encodes those records with the synthetic tool, step, and turn closers required by the shared [persistence contract](../../../.agents/notes/implemented/architecture/2026-06-14-session-persistence.md). Raw mode truncates from its first incomplete line. A checksum/decompression failure in a complete frame, or a defect at or before the last committed `turn/end`, is corruption and rejects.
- **Non-mutating inspection.** `inspect()` returns the detached valid prefix without truncating an incomplete tail or closing an interrupted turn, and leaves the lightweight revision unchanged.
- **Contiguous-seq.** `append` rejects a batch whose first `seq` does not continue the stored log, and rejects non-JSON-serializable `event.data` naming the offending event type.
- **Lightweight revisions.** `listSnapshots()` identifies a log by its device, inode, size, and nanosecond timestamps, avoiding a full-log parse while changing after append, repair, replacement, or store changes.
- **Lightweight revisions.** `listSnapshots(signal?)` identifies a log by its device, inode, size, and nanosecond timestamps, avoiding a full-log parse while changing after append, repair, replacement, or store changes. It forwards the exact signal through artifact discovery and checks cancellation around every `stat`; because filesystem `stat` is not interruptible, cancellation waits for the active call to settle, then rejects without starting another.
## Write path

View File

@@ -281,11 +281,13 @@ export class SessionPersistenceJsonl extends SessionPersistence implements Persi
}
/** List metadata plus a stat-derived identity for each append-only log. */
async listSnapshots(): Promise<SessionPersistenceSnapshot[]> {
async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
const snapshots: SessionPersistenceSnapshot[] = []
for (const artifact of await this.listArtifacts()) {
for (const artifact of await this.listArtifacts(signal)) {
signal?.throwIfAborted()
try {
const identity = await stat(artifact.path, { bigint: true })
signal?.throwIfAborted()
snapshots.push({
header: artifact.header,
revision: SessionPersistenceRevision([
@@ -297,9 +299,11 @@ export class SessionPersistenceJsonl extends SessionPersistence implements Persi
].join(':')),
})
} catch (error: unknown) {
signal?.throwIfAborted()
if (!isENOENT(error)) throw error
}
}
signal?.throwIfAborted()
return snapshots
}

View File

@@ -265,6 +265,56 @@ describe('SessionPersistenceJsonl: durability and crash semantics', () => {
discovery.mockRestore()
})
it('forwards snapshot-list cancellation and awaits in-flight discovery cleanup', async () => {
const persistence = ctx.sessionPersistence as unknown as {
listArtifacts(signal?: AbortSignal): Promise<Array<{ header: SessionHeader; path: string }>>
}
const started = Promise.withResolvers<AbortSignal>()
const cleanup = Promise.withResolvers<undefined>()
vi.spyOn(persistence, 'listArtifacts').mockImplementation(async (signal) => {
if (signal === undefined) throw new Error('expected snapshot-list signal')
started.resolve(signal)
await cleanup.promise
return []
})
const reason = new Error('JSONL snapshot discovery cancelled')
const controller = new AbortController()
const pending = ctx.sessionPersistence.listSnapshots(controller.signal)
expect(await started.promise).toBe(controller.signal)
let settled = false
void pending.then(
() => { settled = true },
() => { settled = true },
)
controller.abort(reason)
await Promise.resolve()
expect(settled).toBe(false)
cleanup.resolve(undefined)
await expect(pending).rejects.toBe(reason)
})
it('checks cancellation after an uncancellable snapshot stat settles', async () => {
const m = meta('snapshot-stat-cancellation')
await ctx.sessionPersistence.create(m)
await ctx.sessionPersistence.append(m.id, oneTurnLog())
const persistence = ctx.sessionPersistence as unknown as {
listArtifacts(signal?: AbortSignal): Promise<Array<{ header: SessionHeader; path: string }>>
}
const discovery = vi.spyOn(persistence, 'listArtifacts').mockResolvedValue([{
header: m,
path: rawLogPath(root, m.cwd, m.id),
}])
const reason = new Error('JSONL snapshot stat cancelled')
const controller = new AbortController()
const pending = ctx.sessionPersistence.listSnapshots(controller.signal)
queueMicrotask(() => { controller.abort(reason) })
await expect(pending).rejects.toBe(reason)
expect(discovery).toHaveBeenCalledWith(controller.signal)
})
it('rejects a stored v0 log containing a legacy request/header-delta event', async () => {
const m = meta('legacy-header-delta', '/legacy')
const path = rawLogPath(root, m.cwd, m.id)

View File

@@ -20,7 +20,7 @@ On filesystems with POSIX modes, the backend requests mode `0700` for missing di
- **Lazy materialization.** `create()` records intent in memory only — no row is written until the first `append`. A created-but-never-appended session has no `sessions` row, so it is absent from `list()` (which reports exactly the sessions that have a row).
- **Interrupted-turn close on load.** `load()` implements the shared [crash-recovery contract](../../../.agents/notes/implemented/architecture/2026-06-14-session-persistence.md): preserve the valid interrupted turn, append its synthetic closing events in one transaction, and remove only a torn tail row. Committed parse errors or sequence gaps make the session unloadable. Because recovery mutates stored rows, the next append starts from a balanced log and accurate cursor.
- **Non-mutating inspection.** `inspect()` returns the detached valid row prefix without deleting a torn tail row or appending recovery closers, and leaves the lightweight revision unchanged.
- **Lightweight revisions.** `listSnapshots()` combines the immutable store and database-file identity, a per-materialization incarnation id, and a per-session counter incremented in each mutating transaction. This keeps unchanged observations stable without parsing event rows and distinguishes independent stores and recreated same-id logs.
- **Lightweight revisions.** `listSnapshots(signal?)` combines the immutable store and database-file identity, a per-materialization incarnation id, and a per-session counter incremented in each mutating transaction. This keeps unchanged observations stable without parsing event rows and distinguishes independent stores and recreated same-id logs. It checks cancellation before and after shared readiness and the synchronous metadata query; the query itself is non-preemptible.
## Configuration (schemastery)

View File

@@ -266,9 +266,12 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
}
/** List metadata with a source-qualified monotonic revision per session. */
async listSnapshots(): Promise<SessionPersistenceSnapshot[]> {
async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
signal?.throwIfAborted()
await this.ready
signal?.throwIfAborted()
const rows = this.db.prepare('SELECT * FROM sessions').all() as unknown as SessionRow[]
signal?.throwIfAborted()
return rows.map(row => ({
header: rowToMeta(row),
revision: SessionPersistenceRevision(

View File

@@ -441,6 +441,31 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
await second.dispose()
})
it('awaits in-flight readiness before surfacing snapshot-list cancellation', async () => {
const b = await backend()
const internals = b.ctx.sessionPersistence as unknown as { ready: Promise<void> }
const originalReady = internals.ready
const readiness = Promise.withResolvers<undefined>()
internals.ready = readiness.promise
const reason = new Error('SQLite snapshot readiness cancelled')
const controller = new AbortController()
const pending = b.ctx.sessionPersistence.listSnapshots(controller.signal)
let settled = false
void pending.then(
() => { settled = true },
() => { settled = true },
)
controller.abort(reason)
await Promise.resolve()
expect(settled).toBe(false)
readiness.resolve(undefined)
await expect(pending).rejects.toBe(reason)
internals.ready = originalReady
await b.dispose()
})
it('exposes the schema version constant', () => {
expect(SCHEMA_VERSION).toBe(8)
})

View File

@@ -14,7 +14,7 @@ The persisted unit IS the existing `SessionEvent` (event-sourced model — the l
| `load(id): Promise<{ meta; events }>` | Return a stored header plus a balanced contiguous log. A live load first flushes its snapshot and rejects while its turn is open; a cold load preserves an interrupted final turn and closes it with synthetic `tool/result`/`step/end?`/`turn/end {interrupted}` events. Only a torn tail fragment is dropped; committed corruption and unknown `version` reject. |
| `inspect(id, signal?): Promise<{ meta; events }>` | Return a detached valid stored prefix without truncating a torn tail, synthesizing recovery closers, or publishing coordinator state. Serialized with same-id writes; the optional signal promptly rejects a queued caller, prevents that queued backend read from starting, and cancels active backend read work. Intended for read models and other observers that must never recover a log. |
| `list(signal?): Promise<SessionHeader[]>` | Lightweight listing from metadata, no full-log parse. The optional signal cancels backend listing work. A zero-event lazily-materialized session is absent from `list`. |
| `listSnapshots(): Promise<SessionPersistenceSnapshot[]>` | Lightweight metadata plus an opaque branded per-log revision, without loading event logs. A revision stays equal while that log and its backing store are unchanged, changes after append or mutating load repair, and cannot collide solely because two stores use the same local counter. |
| `listSnapshots(signal?): Promise<SessionPersistenceSnapshot[]>` | Lightweight metadata plus an opaque branded per-log revision, without loading event logs. A revision stays equal while that log and its backing store are unchanged, changes after append or mutating load repair, and cannot collide solely because two stores use the same local counter. The optional signal requests cancellation of backend discovery work; first-party backends settle any started listing work before rejecting so an awaited call is quiescent. |
## Invariants every backend must honor
@@ -33,7 +33,7 @@ Crash repair is cold-only. For a live id, `load(id)` snapshots the authoritative
When a live session emits `session/disposed`, the coordinator waits for its controller, serializes a final drain, then releases state owned by that exact `Session` object. Failed retirement leaves the controller in the live-session map, so backend teardown can retry it. Backend teardown stops event admission first, flushes every remaining controller, awaits per-id operations, and only then closes the storage handle.
The side-effect-free `locate` and lightweight `listSnapshots` queries remain backend-owned because they describe storage topology and revision identity rather than write orchestration.
The side-effect-free `locate` and lightweight `listSnapshots` queries remain backend-owned because they describe storage topology and revision identity rather than write orchestration. `listSnapshots(signal?)` passes the caller's exact signal into backend discovery so observers can cancel that work without detaching it.
The `PersistenceBackend<TornMarker>` hooks (the only seam between the coordinator and storage):

View File

@@ -122,9 +122,10 @@ export abstract class SessionPersistence extends Service {
* successful mutating {@link load} repair changes the next listed revision.
* Revisions also distinguish independently backed stores so backend-local
* counters cannot compare equal across different persistence sources.
* @param signal - optional cancellation for backend snapshot-listing work.
* @returns one header and opaque revision per materialized session without loading full logs.
*/
abstract listSnapshots(): Promise<SessionPersistenceSnapshot[]>
abstract listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]>
}
export default SessionPersistence

View File

@@ -227,9 +227,11 @@ export function runPersistenceContract(name: string, make: () => Promise<Contrac
try {
const reason = new Error('persistence observation cancelled')
const controller = new AbortController()
await expect(persistence.listSnapshots(controller.signal)).resolves.toEqual([])
controller.abort(reason)
await expect(persistence.list(controller.signal)).rejects.toBe(reason)
await expect(persistence.listSnapshots(controller.signal)).rejects.toBe(reason)
await expect(persistence.inspect(SessionId('cancelled-inspect'), controller.signal))
.rejects.toBe(reason)
} finally {

View File

@@ -137,7 +137,8 @@ class MemoryPersistence extends SessionPersistence implements PersistenceBackend
return [...this.store.values()].map(e => structuredClone(e.meta))
}
async listSnapshots(): Promise<SessionPersistenceSnapshot[]> {
async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
signal?.throwIfAborted()
return [...this.store.values()].map(entry => ({
header: structuredClone(entry.meta),
revision: SessionPersistenceRevision(`events:${entry.events.length}`),

View File

@@ -34,7 +34,7 @@ The database is disposable but reset is guarded: every recognized schema version
The index uses FTS5 `unicode61`. In the implementation experiment it supported the two-character query `AI` and produced an index about 2.1× smaller than the trigram alternative. The trade-off is token/phrase recall rather than arbitrary substring recall: `AI` does not match the token `BRAID`. Use `ctx.sessionQuery.filterEvents()` with a `text` clause when a literal whitespace-flexible substring scan is required. NUL is rejected in queries; reserved highlight markers and NUL in documents are normalized before indexing so presentation markers cannot collide with source text.
Abort signals stop queued work and caller waits around asynchronous source observation. Node's synchronous `DatabaseSync` API cannot interrupt a MATCH statement already executing on the JavaScript thread; the signal is checked immediately before and after the serialized observation/reconciliation boundary.
Abort signals stop queued work and flow unchanged through snapshot listing and non-mutating inspection. Once source work starts, the serialized state machine awaits that backend promise itself—even when a backend ignores cancellation—then checks the signal before starting any further listing, inspection, reconciliation, or query work. The caller therefore observes cancellation only after started backend work is quiescent, and a later search cannot enter the serializer while that cleanup is pending. Node's synchronous `DatabaseSync` API cannot interrupt a metadata or MATCH statement already executing on the JavaScript thread; signals are checked immediately before and after those non-preemptible calls.
## Model Experience

View File

@@ -352,6 +352,7 @@ export class SessionQuerySqlite extends SessionQueryService {
}
private async _reconcile(signal: AbortSignal | undefined): Promise<PersistenceBinding> {
assertNotAborted(signal)
const db = this._requireDb()
const persistedRows = db.prepare(
'SELECT id, revision, generation FROM persisted_sessions',
@@ -452,7 +453,8 @@ export class SessionQuerySqlite extends SessionQueryService {
try {
const canReuseIndexed = this._lastPersistenceIdentity === undefined
|| this._lastPersistenceIdentity === persistenceBinding.identity
const before = await waitWithAbort(persistence.listSnapshots(), signal)
const before = await persistence.listSnapshots(signal)
assertNotAborted(signal)
persisted = materializePersistenceSnapshots(before)
for (const entry of persisted.values()) {
if (canReuseIndexed && indexed.get(entry.header.id)?.revision === entry.revision) continue
@@ -461,13 +463,16 @@ export class SessionQuerySqlite extends SessionQueryService {
// crash-repair side effects; the live-membership retry below makes
// the returned observation live-preferred.
if (initiallyLive.has(entry.header.id) || this.ctx.sessions.get(entry.header.id) !== undefined) continue
const loaded = await waitWithAbort(persistence.inspect(entry.header.id), signal)
assertNotAborted(signal)
const loaded = await persistence.inspect(entry.header.id, signal)
assertNotAborted(signal)
assertSessionHeadersCompatible(entry.header, loaded.meta)
entry.loaded = observeSession(loaded.meta, loaded.events)
}
const after = materializePersistenceSnapshots(
await waitWithAbort(persistence.listSnapshots(), signal),
)
assertNotAborted(signal)
const afterSnapshots = await persistence.listSnapshots(signal)
assertNotAborted(signal)
const after = materializePersistenceSnapshots(afterSnapshots)
if (!samePersistenceSnapshots(persisted, after)) continue
if (this._persistenceBinding !== persistenceBinding) continue
} catch (error: unknown) {

View File

@@ -69,11 +69,16 @@ class TestPersistence extends SessionPersistence {
static nextRevision = 0
static loads = new Map<SessionIdType, number>()
static inspections = new Map<SessionIdType, number>()
static inspectSignals: Array<AbortSignal | undefined> = []
static snapshotSignals: Array<AbortSignal | undefined> = []
static loadEffect: ((entry: { meta: SessionHeader; events: SessionEvent[] }) => void) | undefined
static inspectEffect: ((entry: { meta: SessionHeader; events: SessionEvent[] }) => void | Promise<void>) | undefined
static inspectEffect: ((
entry: { meta: SessionHeader; events: SessionEvent[] },
signal?: AbortSignal,
) => void | Promise<void>) | undefined
static listGate: Promise<void> | undefined
static listStarted: (() => void) | undefined
static snapshotEffect: (() => void | Promise<void>) | undefined
static snapshotEffect: ((signal?: AbortSignal) => void | Promise<void>) | undefined
static snapshotOverride: (() => SessionPersistenceSnapshot[]) | undefined
static failure: unknown
@@ -86,6 +91,8 @@ class TestPersistence extends SessionPersistence {
this.revisions = new Map()
this.loads = new Map()
this.inspections = new Map()
this.inspectSignals = []
this.snapshotSignals = []
this.loadEffect = undefined
this.inspectEffect = undefined
for (const entry of entries) this.set(entry)
@@ -128,12 +135,13 @@ class TestPersistence extends SessionPersistence {
return structuredClone(entry)
}
async inspect(id: SessionIdType): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
async inspect(id: SessionIdType, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
TestPersistence.inspections.set(id, (TestPersistence.inspections.get(id) ?? 0) + 1)
TestPersistence.inspectSignals.push(signal)
if (TestPersistence.failure !== undefined) throw TestPersistence.failure
const entry = TestPersistence.entries.get(id)
if (entry === undefined) throw new Error('missing test session')
await TestPersistence.inspectEffect?.(entry)
await TestPersistence.inspectEffect?.(entry, signal)
TestPersistence.inspectEffect = undefined
return structuredClone(entry)
}
@@ -146,7 +154,8 @@ class TestPersistence extends SessionPersistence {
}
async listSnapshots(): Promise<SessionPersistenceSnapshot[]> {
async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
TestPersistence.snapshotSignals.push(signal)
TestPersistence.listStarted?.()
await TestPersistence.listGate
if (TestPersistence.failure !== undefined) throw TestPersistence.failure
@@ -155,7 +164,7 @@ class TestPersistence extends SessionPersistence {
header: structuredClone(entry.meta),
revision: SessionPersistenceRevision(`test:${TestPersistence.revisions.get(entry.meta.id)}`),
}))
await TestPersistence.snapshotEffect?.()
await TestPersistence.snapshotEffect?.(signal)
return snapshots
}
}
@@ -1209,6 +1218,167 @@ describe('SQLite schema, cancellation, and real persistence integration', () =>
}
})
it.each(['sessions', 'events'] as const)(
'forwards one exact reconciliation signal through both snapshot lists and persisted inspection for %s search',
async (scope) => {
const durable = header(`signal-${scope}`)
TestPersistence.reset([{ meta: durable, events: messageEvents('signal needle') }])
const ctx = await liveContext()
await ctx.plugin(TestPersistence)
const controller = new AbortController()
const result = scope === 'sessions'
? await ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
: await ctx.sessionQuery.searchEvents(
{ sessionId: durable.id, query: 'needle' },
{ signal: controller.signal },
)
expect(result.items).toHaveLength(1)
expect(TestPersistence.snapshotSignals).toEqual([controller.signal, controller.signal])
expect(TestPersistence.inspectSignals).toEqual([controller.signal])
},
)
it.each(['sessions', 'events'] as const)(
'starts no persistence observation for a pre-aborted %s search',
async (scope) => {
const durable = header(`pre-aborted-${scope}`)
TestPersistence.reset([{ meta: durable, events: messageEvents('needle') }])
const ctx = await liveContext()
await ctx.plugin(TestPersistence)
const controller = new AbortController()
controller.abort(new Error(`pre-aborted ${scope}`))
const pending = scope === 'sessions'
? ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
: ctx.sessionQuery.searchEvents(
{ sessionId: durable.id, query: 'needle' },
{ signal: controller.signal },
)
await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
expect(TestPersistence.snapshotSignals).toEqual([])
expect(TestPersistence.inspectSignals).toEqual([])
},
)
it('awaits cooperative snapshot-list cancellation cleanup without starting another observation step', async () => {
const durable = header('cooperative-list-abort')
TestPersistence.reset([{ meta: durable, events: messageEvents('needle') }])
const ctx = await liveContext()
await ctx.plugin(TestPersistence)
const started = Promise.withResolvers<AbortSignal>()
const abortObserved = Promise.withResolvers<undefined>()
const cleanup = Promise.withResolvers<undefined>()
TestPersistence.snapshotEffect = async (signal) => {
TestPersistence.snapshotEffect = undefined
if (signal === undefined) throw new Error('expected reconciliation signal')
started.resolve(signal)
await new Promise<void>((resolve) => {
signal.addEventListener('abort', () => { resolve() }, { once: true })
})
abortObserved.resolve(undefined)
await cleanup.promise
signal.throwIfAborted()
}
const controller = new AbortController()
const pending = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
expect(await started.promise).toBe(controller.signal)
let settled = false
void pending.then(
() => { settled = true },
() => { settled = true },
)
controller.abort(new Error('cooperative list cancellation'))
await abortObserved.promise
expect(settled).toBe(false)
expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
expect(TestPersistence.inspectSignals).toEqual([])
cleanup.resolve(undefined)
await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
})
it('keeps a second search serialized while an abort-ignoring snapshot list finishes', async () => {
const durable = header('serialized-list-abort')
TestPersistence.reset([{ meta: durable, events: messageEvents('needle') }])
const ctx = await liveContext()
await ctx.plugin(TestPersistence)
const cleanup = Promise.withResolvers<undefined>()
const started = Promise.withResolvers<undefined>()
TestPersistence.listGate = cleanup.promise
TestPersistence.listStarted = () => {
TestPersistence.listStarted = undefined
started.resolve(undefined)
}
const controller = new AbortController()
const first = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
await started.promise
let firstSettled = false
let secondSettled = false
void first.then(
() => { firstSettled = true },
() => { firstSettled = true },
)
controller.abort(new Error('ignored list cancellation'))
const second = ctx.sessionQuery.searchEvents({ sessionId: durable.id, query: 'needle' })
void second.then(
() => { secondSettled = true },
() => { secondSettled = true },
)
await Promise.resolve()
expect(firstSettled).toBe(false)
expect(secondSettled).toBe(false)
expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
expect(TestPersistence.inspectSignals).toEqual([])
cleanup.resolve(undefined)
await expect(first).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
await expect(second).resolves.toMatchObject({ items: [{ sessionId: durable.id }] })
})
it('awaits an abort-ignoring inspection and starts neither another inspection nor the after-list', async () => {
const first = header('ignored-inspect-first')
const second = header('ignored-inspect-second')
TestPersistence.reset([
{ meta: first, events: messageEvents('first needle') },
{ meta: second, events: messageEvents('second needle') },
])
const ctx = await liveContext()
await ctx.plugin(TestPersistence)
const started = Promise.withResolvers<AbortSignal>()
const cleanup = Promise.withResolvers<undefined>()
TestPersistence.inspectEffect = async (_entry, signal) => {
TestPersistence.inspectEffect = undefined
if (signal === undefined) throw new Error('expected reconciliation signal')
started.resolve(signal)
await cleanup.promise
}
const controller = new AbortController()
const pending = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
expect(await started.promise).toBe(controller.signal)
let settled = false
void pending.then(
() => { settled = true },
() => { settled = true },
)
controller.abort(new Error('ignored inspect cancellation'))
await Promise.resolve()
expect(settled).toBe(false)
expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
expect(TestPersistence.inspections.get(first.id)).toBe(1)
expect(TestPersistence.inspections.get(second.id)).toBeUndefined()
cleanup.resolve(undefined)
await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
expect(TestPersistence.inspections.get(second.id)).toBeUndefined()
})
it('cancels both queued and in-flight source waits without committing them', async () => {
TestPersistence.reset()
const ctx = await liveContext()
@@ -1262,8 +1432,15 @@ describe('SQLite schema, cancellation, and real persistence integration', () =>
const active = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: activeController.signal })
await activeStarted
activeController.abort()
await expect(active).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
let activeSettled = false
void active.then(
() => { activeSettled = true },
() => { activeSettled = true },
)
await Promise.resolve()
expect(activeSettled).toBe(false)
releaseActive()
await expect(active).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
const db = (ctx.sessionQuery as unknown as { _db: DatabaseSync })._db
expect(db.prepare('SELECT COUNT(*) AS count FROM persisted_sessions').get()).toEqual({ count: 0 })
@@ -1271,6 +1448,57 @@ describe('SQLite schema, cancellation, and real persistence integration', () =>
.resolves.toMatchObject({ items: [{ header: { id: SessionId('uncommitted') } }] })
})
it.each([
[new Error('ready error'), 'ready error'],
['non-error ready failure', 'session-search dependency rejected with a non-Error value'],
])('normalizes a rejected readiness wait before mapping it to an index error', async (failure, detail) => {
TestPersistence.reset()
const ctx = await liveContext()
const internals = ctx.sessionQuery as unknown as {
_ready: Promise<void>
_ensureReady(signal: AbortSignal): Promise<void>
}
internals._ready = Promise.resolve().then(() => {
throw failure
})
await expect(internals._ensureReady(new AbortController().signal))
.rejects.toThrow(`session-search SQLite index failed to open: ${detail}`)
})
it('checks cancellation after readiness before reconciliation accesses SQLite', async () => {
TestPersistence.reset()
const ctx = await liveContext()
const internals = ctx.sessionQuery as unknown as {
_db: DatabaseSync
_ready: Promise<void>
_ensureReady(signal: AbortSignal | undefined): Promise<void>
}
const readiness = Promise.withResolvers<undefined>()
internals._ready = readiness.promise
const readyWaitStarted = Promise.withResolvers<undefined>()
const ensureReady = internals._ensureReady.bind(internals)
vi.spyOn(internals, '_ensureReady').mockImplementation(async (signal) => {
const pending = ensureReady(signal)
readyWaitStarted.resolve(undefined)
return pending
})
const prepare = vi.spyOn(internals._db, 'prepare')
const reason = new Error('cancelled after readiness')
const controller = new AbortController()
const pending = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
await readyWaitStarted.promise
const queueBoundaryAbort = readiness.promise.then(() => {
queueMicrotask(() => { controller.abort(reason) })
})
readiness.resolve(undefined)
await queueBoundaryAbort
await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
expect(prepare).not.toHaveBeenCalled()
})
it('rejects queued and future work when close waits for an accepted operation', async () => {
TestPersistence.reset()
let release!: () => void