refactor(session-query): narrow phase one to exact reads

This commit is contained in:
Hypatia May
2026-07-11 12:20:35 +08:00
parent 8fd68731ba
commit ad32c57e72
35 changed files with 396 additions and 3036 deletions

View File

@@ -1,9 +1,9 @@
# session-query/ — session retrieval capability family
Trusted read-model infrastructure over live and durable session logs. The interface package owns `ctx.sessionQuery`, logical-corpus resolution, filters, traces, text extractors, and the full-text provider contract. A search backend is a separate implementation package; a model tool or UI remains a separate consumer.
Trusted exact reads over live and durable session logs. Phase one contains one interface package that owns `ctx.sessionQuery`, logical-corpus precedence, surface classification, and bounded event reads.
| Package | Role | ctx key |
|---|---|---|
| [`session-query/`](session-query/README.md) | Retrieval service and provider contract | `ctx.sessionQuery` |
| [`session-query/`](session-query/README.md) | Logical-corpus and exact-event read service | `ctx.sessionQuery` |
The family is independent of the [compaction capability](../compact/README.md): it reads compaction provenance from the canonical session log but does not participate in compaction policy or execution. The provider-neutral decision is recorded in the [session-query RFC](../../docs/rfc/implemented/feature/2026-07-10-session-query-service.md); the first proposed backend is specified separately in the [SQLite provider RFC](../../docs/rfc/proposed/feature/2026-07-10-sqlite-session-query-provider.md).
The family is independent of compaction: it reads the canonical session log but does not participate in compaction policy or execution. Full-text search remains proposed as a phase-two SQLite package rather than a speculative provider seam in this interface package.

View File

@@ -1,52 +1,23 @@
# @deepseek-ai/dsh-session-query
Provider-neutral session-history retrieval (`ctx.sessionQuery`). The service presents live `ctx.sessions` state and, when mounted, `ctx.sessionPersistence` state as one logical corpus. A matching id produces one record: live events win, while independent `live` and `persisted` flags report both source availabilities. Conflicting immutable headers fail with `SESSION_QUERY_SOURCE_CONFLICT` instead of silently merging unrelated histories.
Exact session-history retrieval through `ctx.sessionQuery`. The service presents live `ctx.sessions` and an optional, dynamically mounted `ctx.sessionPersistence` as one logical corpus. Matching ids produce one record: live events win, while `live` and `persisted` report both source availabilities. Conflicting immutable headers fail with `SESSION_QUERY_SOURCE_CONFLICT`.
This is trusted context-wide infrastructure. It performs no caller authorization; a future model tool or UI must constrain which sessions its caller may inspect.
## Reads and traces
## Reads
- `listSessions()` returns cloned lightweight records in deterministic newest-first order.
- `listEvents(sessionId)` classifies each raw event as `current`, `shadowed`, or `log-only` using the shared `dsh-session` surface fold.
- `readEvent(request)` returns the cloned target and a bounded raw-seq window. `before` and `after` default to zero and may not exceed `readWindowMax` (default 50).
- `traceSession(sessionId)` returns nearest-first parents, a known root or explicit unresolved parent id, and the complete deterministic descendant tree. A connected lineage cycle fails with `SESSION_QUERY_INVALID_LINEAGE`.
- `traceEvent(sessionId, seq)` returns direct provenance references and reverse references, direct shadows, the immediate replacer, and the transitive replacement chain toward the current surface node. Related nodes stay seq links; callers use `readEvent()` for content.
- `listSessions()` reads current persistence metadata, merges live records with live precedence, and returns cloned records in deterministic newest-first order.
- `listEvents(sessionId)` loads the live-preferred raw log and classifies each event as `current`, `shadowed`, or `log-only` with the shared `dsh-session` surface fold.
- `readEvent(request)` returns a cloned header, the full target event, and a bounded raw-seq window. `before` and `after` default to zero and may not exceed `readWindowMax`.
An installed persistence backend is optional and may mount or unmount dynamically. Cross-session operations fail with `SESSION_QUERY_PERSISTENCE_FAILED` while installed persistence is unreadable. A read targeting a known live session never depends on persistence health. Provider-side persisted rows are deactivated rather than deleted when persistence is absent.
Persistence is optional and may mount or unmount dynamically. A cross-corpus list fails with `SESSION_QUERY_PERSISTENCE_FAILED` while mounted persistence is unreadable. A read targeting a known live session does not consult persistence, so durable backend health cannot make current in-memory history unreadable. Persisted exact reads list before loading, and reject a metadata mismatch rather than combining inconsistent observations.
## Filters
`filterSessionResults()` and `filterEventResults()` are pure generic transforms over records or richer hits. Each discriminated filter is serializable. Values within one filter are OR alternatives; filters in the supplied array are an AND chain. The functions preserve order and item identity and return a fresh array.
Session filters cover id, exact cwd, inclusive creation time, parent id/root, and live/persisted availability. Event filters cover inclusive seq/time, event type, and surface status. Search requests accept the same specs as pre-ranking filters. Applying the pure functions to a materialized provider page is a post-filter: it never fetches replacement hits to refill the page.
## Full-text providers
`registerSearchProvider(provider)` is effect-scoped and ids are unique. Its async disposer removes the provider from selection immediately, lets already accepted transactions finish, and settles after they drain. Without `searchProvider`, exactly one locally available provider must be registered; explicit selection fails loudly when the named provider is missing or unavailable. Search pages default to 20 hits and reject limits above 100; a provider returning more hits than the normalized request limit fails with a typed provider error rather than silently dropping cursor-addressable results. Provider scores never cross the public API: event hits carry a plain snippet, while each session hit carries exactly one best matching event.
The service feeds providers two independent layers: a durable persisted base (`persistedInventory`, `replacePersisted`, `removePersisted`, `setPersistedActive`) and an ephemeral live override (`replaceLive`, `removeLive`). A search waits for the relevant source state observed before its call: the whole corpus for session search, only the target for a live event search. Failed derived updates do not fail session writes; affected searches receive `SESSION_QUERY_INDEX_FAILED`, and a later search retries the dirty state. `AbortSignal` lets a caller stop waiting and is also passed to provider search.
Persisted snapshots carry a SHA-256 fingerprint over canonical header/events plus the versions of relevant extractors. Reconciliation still loads and hashes canonical logs, but a provider replacement occurs only for a new or changed fingerprint; stale durable inventory entries are removed only while persistence is active and authoritative.
Providers receive resolved `SessionSearchSpec` and `SessionEventSearchSpec` values whose `limit` is required after service defaulting and validation. Public service callers use `SessionSearchRequest` and `SessionEventSearchRequest`, where `limit` remains optional.
## Errors
`SessionQueryError.code` is the closed `SessionQueryErrorCode` union: `SESSION_QUERY_ABORTED`, `SESSION_QUERY_DUPLICATE_EXTRACTOR`, `SESSION_QUERY_DUPLICATE_PROVIDER`, `SESSION_QUERY_EVENT_NOT_FOUND`, `SESSION_QUERY_INDEX_FAILED`, `SESSION_QUERY_INVALID_CONFIG`, `SESSION_QUERY_INVALID_EXTRACTOR`, `SESSION_QUERY_INVALID_FILTER`, `SESSION_QUERY_INVALID_LIMIT`, `SESSION_QUERY_INVALID_LINEAGE`, `SESSION_QUERY_INVALID_QUERY`, `SESSION_QUERY_INVALID_SURFACE`, `SESSION_QUERY_INVALID_WINDOW`, `SESSION_QUERY_PERSISTENCE_FAILED`, `SESSION_QUERY_PROVIDER_AMBIGUOUS`, `SESSION_QUERY_PROVIDER_CONFIGURED_MISSING`, `SESSION_QUERY_PROVIDER_CONFIGURED_UNAVAILABLE`, `SESSION_QUERY_PROVIDER_ERROR`, `SESSION_QUERY_PROVIDER_UNAVAILABLE`, `SESSION_QUERY_SESSION_NOT_FOUND`, and `SESSION_QUERY_SOURCE_CONFLICT`.
## Text extractors
Core extraction indexes semantic message text and reasoning, tool names/arguments/results, blocked prompts, context and steering, todos, and error/status detail. Stream chunks, request headers, and structural-only events contribute no document. Unknown event and content-block types contribute no text until their owner registers a versioned extractor with `registerEventTextExtractor()` or `registerContentTextExtractor()`.
Extractor registrations are unique per discriminant and effect-scoped. Their stable versions participate in fingerprints, so changing extraction semantics invalidates only sessions whose indexed source uses that extractor.
`SessionQueryError.code` is a closed union: `SESSION_QUERY_EVENT_NOT_FOUND`, `SESSION_QUERY_INVALID_CONFIG`, `SESSION_QUERY_INVALID_SURFACE`, `SESSION_QUERY_INVALID_WINDOW`, `SESSION_QUERY_PERSISTENCE_FAILED`, `SESSION_QUERY_SESSION_NOT_FOUND`, and `SESSION_QUERY_SOURCE_CONFLICT`.
## Configuration
| Key | Default | Contract |
|---|---:|---|
| `searchProvider` | omitted | Explicit provider id; omission requires exactly one available provider. |
| `defaultLimit` | `20` | Search page size when the request omits `limit`. |
| `maxLimit` | `100` | Maximum accepted search page size; must be at least `defaultLimit`. |
| `readWindowMax` | `50` | Maximum `before` or `after` raw-event count. |
The package ships no full-text backend and no model-facing tool. The proposed SQLite implementation is a later, independent phase described in the [SQLite provider RFC](../../../docs/rfc/proposed/feature/2026-07-10-sqlite-session-query-provider.md).
This phase deliberately has no filters, lineage/provenance traversal, extraction registry, search-provider protocol, index synchronization, or model-facing tool. Full-text search belongs beside its first real implementation; the proposed SQLite package and its single transaction/reconciliation owner are described in the [phase-two RFC](../../../docs/rfc/proposed/feature/2026-07-10-sqlite-session-query-provider.md).

View File

@@ -1,6 +1,6 @@
{
"name": "@deepseek-ai/dsh-session-query",
"description": "Provider-neutral live and persisted session retrieval service (ctx.sessionQuery)",
"description": "Live-preferred exact session-history retrieval service (ctx.sessionQuery)",
"version": "0.0.1",
"private": true,
"type": "module",

View File

@@ -1,51 +1,23 @@
/**
* Public configuration, defaults, and typed failures for session-query.
*
* @module @deepseek-ai/dsh-session-query/config
*/
/** Public configuration and typed failures for session-query. */
import { HarnessError } from '@deepseek-ai/dsh-llm'
/** Default page size for provider-backed search. */
export const SESSION_QUERY_DEFAULT_LIMIT = 20
/** Maximum page size accepted by provider-backed search. */
export const SESSION_QUERY_MAX_LIMIT = 100
/** Default maximum `before`/`after` raw-event window. */
export const SESSION_QUERY_READ_WINDOW_MAX = 50
/** Configuration for the provider-neutral session-query service. */
/** Configuration for exact session-query reads. */
export interface Config {
/** Explicit provider id; omitted auto-selects exactly one usable provider. */
searchProvider?: string
/** Default search result page size. Defaults to 20. */
defaultLimit?: number
/** Maximum accepted search page size. Defaults to 100. */
maxLimit?: number
/** Maximum accepted raw read context on either side. Defaults to 50. */
readWindowMax?: number
}
/** Complete stable machine-routable failure taxonomy for session-query. */
/** Stable machine-routable failure taxonomy for exact session reads. */
export type SessionQueryErrorCode =
| 'SESSION_QUERY_ABORTED'
| 'SESSION_QUERY_DUPLICATE_EXTRACTOR'
| 'SESSION_QUERY_DUPLICATE_PROVIDER'
| 'SESSION_QUERY_EVENT_NOT_FOUND'
| 'SESSION_QUERY_INDEX_FAILED'
| 'SESSION_QUERY_INVALID_CONFIG'
| 'SESSION_QUERY_INVALID_EXTRACTOR'
| 'SESSION_QUERY_INVALID_FILTER'
| 'SESSION_QUERY_INVALID_LIMIT'
| 'SESSION_QUERY_INVALID_LINEAGE'
| 'SESSION_QUERY_INVALID_QUERY'
| 'SESSION_QUERY_INVALID_SURFACE'
| 'SESSION_QUERY_INVALID_WINDOW'
| 'SESSION_QUERY_PERSISTENCE_FAILED'
| 'SESSION_QUERY_PROVIDER_AMBIGUOUS'
| 'SESSION_QUERY_PROVIDER_CONFIGURED_MISSING'
| 'SESSION_QUERY_PROVIDER_CONFIGURED_UNAVAILABLE'
| 'SESSION_QUERY_PROVIDER_ERROR'
| 'SESSION_QUERY_PROVIDER_UNAVAILABLE'
| 'SESSION_QUERY_SESSION_NOT_FOUND'
| 'SESSION_QUERY_SOURCE_CONFLICT'

View File

@@ -1,45 +1,32 @@
/** Live/persisted logical-corpus resolution for session-query. */
import type { Context } from 'cordis'
import type { Session, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
import type { Session, SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
import type SessionPersistence from '@deepseek-ai/dsh-session-persistence'
import type { SessionRecord } from './types.ts'
import type { LoadedSession } from './extraction.ts'
import { canonicalJson } from './extraction.ts'
import { SessionQueryError } from './config.ts'
interface PersistenceBinding {
token: symbol
service: SessionPersistence
headers: Map<SessionId, SessionHeader>
/** Notifications retained until a list that began after them completes. */
observations: Map<SessionId, PersistedObservation>
observationGeneration: number
error?: unknown
refreshing: Promise<void> | undefined
}
interface PersistedObservation {
generation: number
/** Detached source selected for one exact read. */
export interface LogicalSession {
/** Cloned source header. */
header: SessionHeader
/** Cloned raw event log. */
events: SessionEvent[]
}
/** Active persistence view used by provider reconciliation. */
export interface PersistenceView {
/** Canonical headers in deterministic creation order. */
headers: SessionHeader[]
/** Load one canonical persisted source. */
load(id: SessionId): Promise<LoadedSession>
}
/** Resolves one live-preferred corpus while containing optional persistence lifecycle. */
/** Resolves a live-preferred corpus against the persistence service mounted now. */
export class SessionCorpus {
private _persistence: PersistenceBinding | undefined
private _persistence: SessionPersistence | undefined
constructor(private readonly _ctx: Context) {
_ctx.effect(() => {
const fiber = _ctx.inject(['sessionPersistence'], (childCtx: Context) => {
this._attachPersistence(childCtx, childCtx.sessionPersistence)
const service = childCtx.sessionPersistence
this._persistence = service
childCtx.effect(() => () => {
/* v8 ignore next -- a stale optional-service disposer cannot clear a replacement */
if (this._persistence === service) this._persistence = undefined
}, 'sessionQuery.persistenceBinding')
})
return () => void fiber.dispose()
}, 'sessionQuery.optionalPersistence')
@@ -47,23 +34,22 @@ export class SessionCorpus {
/**
* List the complete logical corpus with live precedence and cloned headers.
* @returns logical records in deterministic newest-first order.
* @returns records in deterministic newest-first order.
*/
async listSessions(): Promise<SessionRecord[]> {
const binding = await this._ensurePersistence()
const persistence = this._persistence
const persisted = persistence === undefined ? [] : await listPersisted(persistence)
const records = new Map<SessionId, SessionRecord>()
if (binding !== undefined) {
for (const header of binding.headers.values()) {
records.set(header.id, { header: structuredClone(header), live: false, persisted: true })
}
for (const header of persisted) {
records.set(header.id, { header: structuredClone(header), live: false, persisted: true })
}
for (const session of this._ctx.sessions.list()) {
const persisted = binding?.headers.get(session.id)
if (persisted !== undefined) this._assertCompatibleHeaders(session.header, persisted)
const durable = records.get(session.id)
if (durable !== undefined) assertCompatibleHeaders(session.header, durable.header)
records.set(session.id, {
header: structuredClone(session.header),
live: true,
persisted: persisted !== undefined,
persisted: durable !== undefined,
})
}
return [...records.values()].sort(compareSessions)
@@ -71,156 +57,69 @@ export class SessionCorpus {
/**
* Load one logical source, preferring a detached live snapshot.
*
* A known live target never consults persistence, so an optional backend's
* failure cannot make current in-memory history unreadable.
* @param sessionId - session to resolve.
* @returns detached live-preferred metadata and events.
* @returns detached live-preferred header and events.
*/
async loadLogical(sessionId: SessionId): Promise<LoadedSession> {
async load(sessionId: SessionId): Promise<LogicalSession> {
const live = this._ctx.sessions.get(sessionId)
if (live !== undefined) return this.snapshotLive(live)
const binding = await this._ensurePersistence()
if (binding === undefined || !binding.headers.has(sessionId)) {
throw new SessionQueryError(`session "${sessionId}" not found`, 'SESSION_QUERY_SESSION_NOT_FOUND')
}
return this._loadPersisted(binding, sessionId)
}
/**
* Return a detached live source with current availability flags.
* @param session - live session to snapshot.
* @returns detached metadata and events.
*/
snapshotLive(session: Session): LoadedSession {
const persistedHeader = this._persistence?.headers.get(session.id)
if (persistedHeader !== undefined) this._assertCompatibleHeaders(session.header, persistedHeader)
return {
record: {
header: structuredClone(session.header),
live: true,
persisted: persistedHeader !== undefined,
},
events: session.events.map(event => structuredClone(event)),
}
}
/**
* Get one live session without consulting persistence.
* @param sessionId - live id to resolve.
* @returns current store object, or undefined.
*/
getLive(sessionId: SessionId): Session | undefined {
return this._ctx.sessions.get(sessionId)
}
/**
* List live sessions in store order.
* @returns fresh array of current store objects.
*/
listLive(): Session[] {
return this._ctx.sessions.list()
}
/**
* Resolve an authoritative persisted view.
* @returns cloned headers and loader, or undefined while unmounted.
*/
async persistenceView(): Promise<PersistenceView | undefined> {
const binding = await this._ensurePersistence()
if (binding === undefined) return undefined
return {
headers: [...binding.headers.values()].map(header => structuredClone(header)).sort(compareHeadersAscending),
load: id => this._loadPersisted(binding, id),
}
}
private _attachPersistence(ctx: Context, service: SessionPersistence): void {
const binding: PersistenceBinding = {
token: Symbol('session-query-persistence'),
service,
headers: new Map(),
observations: new Map(),
observationGeneration: 0,
refreshing: undefined,
}
this._persistence = binding
void this._refreshPersistence(binding)
ctx.on('session/persisted', (header) => {
/* v8 ignore next -- a stale notification can race optional-service disposal */
if (this._persistence?.token !== binding.token) return
const snapshot = structuredClone(header)
const observation = { generation: ++binding.observationGeneration, header: snapshot }
binding.headers.set(header.id, snapshot)
binding.observations.set(header.id, observation)
})
ctx.effect(() => () => { this._detachPersistence(binding) }, 'sessionQuery.persistenceBinding')
}
private _detachPersistence(binding: PersistenceBinding): void {
/* v8 ignore next -- duplicate optional-service disposal is a Cordis teardown edge */
if (this._persistence?.token !== binding.token) return
this._persistence = undefined
}
private _refreshPersistence(binding: PersistenceBinding): Promise<void> {
if (binding.refreshing !== undefined) return binding.refreshing
const startGeneration = binding.observationGeneration
const refresh = binding.service.list().then((headers) => {
/* v8 ignore next -- a list completion can race optional-service disposal */
if (this._persistence?.token !== binding.token) return
const nextHeaders = new Map(headers.map(header => [header.id, structuredClone(header)]))
for (const [id, observation] of binding.observations) {
// A notification newer than this list's snapshot is the authoritative
// read-your-writes layer; older ones must already be present in list().
if (observation.generation > startGeneration) {
nextHeaders.set(id, structuredClone(observation.header))
} else {
binding.observations.delete(id)
}
}
binding.headers = nextHeaders
binding.error = undefined
}).catch((error: unknown) => {
/* v8 ignore next -- a failed list can race optional-service disposal */
if (this._persistence?.token !== binding.token) return
binding.error = error
}).finally(() => {
/* v8 ignore next -- a newer refresh may already own the slot */
if (binding.refreshing === refresh) binding.refreshing = undefined
})
binding.refreshing = refresh
return refresh
}
private async _ensurePersistence(): Promise<PersistenceBinding | undefined> {
const binding = this._persistence
if (binding === undefined) return undefined
await this._refreshPersistence(binding)
if (binding.error !== undefined) {
const cause = binding.error
throw new SessionQueryError(`session persistence listing failed: ${errorMessage(cause)}`, 'SESSION_QUERY_PERSISTENCE_FAILED', { cause })
}
return binding
}
private async _loadPersisted(binding: PersistenceBinding, sessionId: SessionId): Promise<LoadedSession> {
if (live !== undefined) return snapshotLive(live)
const persistence = this._persistence
if (persistence === undefined) throw notFound(sessionId)
const listed = (await listPersisted(persistence)).find(header => header.id === sessionId)
if (listed === undefined) throw notFound(sessionId)
let loaded: Awaited<ReturnType<SessionPersistence['load']>>
try {
const loaded = await binding.service.load(sessionId)
const listed = binding.headers.get(sessionId)
/* v8 ignore else -- every internal persisted load starts from a listed header */
if (listed !== undefined) this._assertCompatibleHeaders(loaded.meta, listed)
return {
record: { header: structuredClone(loaded.meta), live: false, persisted: true },
events: loaded.events.map(event => structuredClone(event)),
}
loaded = await persistence.load(sessionId)
} catch (error: unknown) {
if (error instanceof SessionQueryError) throw error
throw new SessionQueryError(`failed to load session "${sessionId}": ${errorMessage(error)}`, 'SESSION_QUERY_PERSISTENCE_FAILED', { cause: error })
throw new SessionQueryError(
`failed to load session "${sessionId}": ${errorMessage(error)}`,
'SESSION_QUERY_PERSISTENCE_FAILED',
{ cause: error },
)
}
assertCompatibleHeaders(loaded.meta, listed)
return {
header: structuredClone(loaded.meta),
events: loaded.events.map(event => structuredClone(event)),
}
}
}
private _assertCompatibleHeaders(a: SessionHeader, b: SessionHeader): void {
if (canonicalJson(a) !== canonicalJson(b)) {
throw new SessionQueryError(`live and persisted headers conflict for session "${a.id}"`, 'SESSION_QUERY_SOURCE_CONFLICT')
}
async function listPersisted(persistence: SessionPersistence): Promise<SessionHeader[]> {
try {
return await persistence.list()
} catch (error: unknown) {
throw new SessionQueryError(
`session persistence listing failed: ${errorMessage(error)}`,
'SESSION_QUERY_PERSISTENCE_FAILED',
{ cause: error },
)
}
}
function snapshotLive(session: Session): LogicalSession {
return {
header: structuredClone(session.header),
events: session.events.map(event => structuredClone(event)),
}
}
function assertCompatibleHeaders(a: SessionHeader, b: SessionHeader): void {
if (
a.version !== b.version
|| a.id !== b.id
|| a.createdAt !== b.createdAt
|| a.cwd !== b.cwd
|| a.parentSession !== b.parentSession
|| a.seedLength !== b.seedLength
) {
throw new SessionQueryError(
`live and persisted headers conflict for session "${a.id}"`,
'SESSION_QUERY_SOURCE_CONFLICT',
)
}
}
@@ -228,11 +127,10 @@ function compareSessions(a: SessionRecord, b: SessionRecord): number {
return b.header.createdAt - a.header.createdAt || a.header.id.localeCompare(b.header.id)
}
function compareHeadersAscending(a: SessionHeader, b: SessionHeader): number {
return a.createdAt - b.createdAt || a.id.localeCompare(b.id)
function notFound(sessionId: SessionId): SessionQueryError {
return new SessionQueryError(`session "${sessionId}" not found`, 'SESSION_QUERY_SESSION_NOT_FOUND')
}
function errorMessage(error: unknown): string {
/* v8 ignore next -- persistence service contracts reject Error instances */
return error instanceof Error ? error.message : 'unknown error'
}

View File

@@ -1,254 +0,0 @@
/** Semantic text extraction and stable provider snapshot fingerprints. */
import { createHash } from 'node:crypto'
import type { Context } from 'cordis'
import type { ContentBlock, ContentBlockMap, ContentBlockType } from '@deepseek-ai/dsh-llm'
import type { SessionEvent, SessionEventType } from '@deepseek-ai/dsh-session'
import type {
SessionContentTextExtractor,
SessionEventTextExtractor,
SessionIndexDocument,
SessionIndexSnapshot,
SessionRecord,
} from './types.ts'
import { SessionQueryError } from './config.ts'
import { eventRecords } from './tracing.ts'
/** Canonical session source consumed by extraction and provider reconciliation. */
export interface LoadedSession {
/** Logical source metadata. */
record: SessionRecord
/** Detached canonical events. */
events: SessionEvent[]
}
interface StoredEventExtractor {
version: string
extract(event: SessionEvent): readonly string[]
}
interface StoredContentExtractor {
version: string
extract(block: ContentBlock): readonly string[]
}
/** Owns core/custom semantic extractors and builds versioned index snapshots. */
export class SessionTextExtractors {
private readonly _eventExtractors = new Map<SessionEventType, StoredEventExtractor>()
private readonly _contentExtractors = new Map<ContentBlockType, StoredContentExtractor>()
constructor() {
this._installCoreExtractors()
}
/**
* Register one effect-scoped event extractor.
* @param ctx - contributing caller context.
* @param type - event discriminant.
* @param extractor - versioned semantic extractor.
* @returns disposer for the registration.
*/
registerEvent<K extends SessionEventType>(
ctx: Context,
type: K,
extractor: SessionEventTextExtractor<K>,
): () => void {
this._validateVersion(type, extractor.version)
if (this._eventExtractors.has(type)) {
throw new SessionQueryError(`session event text extractor "${type}" is already registered`, 'SESSION_QUERY_DUPLICATE_EXTRACTOR')
}
const stored: StoredEventExtractor = {
version: extractor.version,
extract: event => extractor.extract(event as SessionEvent<K>),
}
const dispose = ctx.effect(function* (this: SessionTextExtractors) {
this._eventExtractors.set(type, stored)
yield () => {
this._eventExtractors.delete(type)
}
}.bind(this), `sessionQuery.eventExtractor(${type})`)
return () => void dispose()
}
/**
* Register one effect-scoped content-block extractor.
* @param ctx - contributing caller context.
* @param type - content-block discriminant.
* @param extractor - versioned semantic extractor.
* @returns disposer for the registration.
*/
registerContent<K extends ContentBlockType>(
ctx: Context,
type: K,
extractor: SessionContentTextExtractor<K>,
): () => void {
this._validateVersion(type, extractor.version)
if (this._contentExtractors.has(type)) {
throw new SessionQueryError(`session content text extractor "${type}" is already registered`, 'SESSION_QUERY_DUPLICATE_EXTRACTOR')
}
const stored: StoredContentExtractor = {
version: extractor.version,
extract: block => extractor.extract(block as ContentBlockMap[K]),
}
const dispose = ctx.effect(function* (this: SessionTextExtractors) {
this._contentExtractors.set(type, stored)
yield () => {
this._contentExtractors.delete(type)
}
}.bind(this), `sessionQuery.contentExtractor(${type})`)
return () => void dispose()
}
/**
* Build one provider-neutral snapshot and SHA-256 source/version fingerprint.
* @param loaded - detached canonical source.
* @returns lightweight documents and stable fingerprint.
*/
buildSnapshot(loaded: LoadedSession): SessionIndexSnapshot {
const records = eventRecords(loaded.record.header.id, loaded.events)
const documents: SessionIndexDocument[] = []
const eventVersions = new Set<string>()
const blockVersions = new Set<string>()
for (const event of loaded.events) {
const extractor = this._eventExtractors.get(event.type)
if (extractor === undefined) continue
eventVersions.add(`${event.type}@${extractor.version}`)
collectBlockVersions(event.data, this._contentExtractors, blockVersions)
const text = normalizeText(extractor.extract(event))
if (text.length === 0) continue
// The event record array parallels the contiguous log.
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
documents.push({ ...records[event.seq]!, text })
}
const fingerprint = createHash('sha256').update(canonicalJson({
header: loaded.record.header,
events: loaded.events,
eventExtractors: [...eventVersions].sort(),
contentExtractors: [...blockVersions].sort(),
})).digest('hex')
return {
session: cloneRecord(loaded.record),
fingerprint,
documents,
}
}
private _installCoreExtractors(): void {
this._contentExtractors.set('text', { version: '1', extract: block => [(block as ContentBlockMap['text']).text] })
this._contentExtractors.set('reasoning', { version: '1', extract: block => [(block as ContentBlockMap['reasoning']).text] })
this._contentExtractors.set('tool-call', {
version: '1',
extract: (block) => {
const call = block as ContentBlockMap['tool-call']
return [call.name, call.arguments]
},
})
this._contentExtractors.set('tool-result', {
version: '1',
extract: block => this._extractBlocks((block as ContentBlockMap['tool-result']).content),
})
for (const type of ['user/message', 'assistant/message', 'context/message', 'steering/message'] as const) {
this._eventExtractors.set(type, {
version: '1',
extract: event => this._extractBlocks((event as SessionEvent<typeof type>).data.content),
})
}
this._eventExtractors.set('prompt/blocked', {
version: '1',
extract: (event) => {
const data = (event as SessionEvent<'prompt/blocked'>).data
return [...this._extractBlocks(data.content), data.reason]
},
})
this._eventExtractors.set('tool/call', {
version: '1',
extract: (event) => {
const data = (event as SessionEvent<'tool/call'>).data
return [data.name, data.arguments]
},
})
this._eventExtractors.set('tool/result', {
version: '1',
extract: (event) => {
const data = (event as SessionEvent<'tool/result'>).data
return [...this._extractBlocks(data.content), data.error?.name ?? '', data.error?.code ?? '']
},
})
this._eventExtractors.set('todo/write', {
version: '1',
extract: event => (event as SessionEvent<'todo/write'>).data.todos.map(todo => `${todo.status} ${todo.content}`),
})
this._eventExtractors.set('turn/end', {
version: '1',
extract: (event) => {
const reason = (event as SessionEvent<'turn/end'>).data.reason
switch (reason.kind) {
case 'error': return ['error', reason.message, reason.code ?? '']
case 'aborted': return ['aborted', reason.reason ?? '']
case 'rejected': return ['rejected', reason.reason]
case 'disposed': return ['disposed']
case 'max-tokens': return ['max-tokens']
case 'interrupted': return ['interrupted']
case 'completed': return []
// TurnEndReasonMap is merge-extensible; unknown variants contribute no text.
/* v8 ignore next -- only an external declaration-merged reason can reach this fallback */
default: return []
}
},
})
}
private _extractBlocks(blocks: readonly ContentBlock[]): string[] {
const fragments: string[] = []
for (const block of blocks) {
const extractor = this._contentExtractors.get(block.type)
if (extractor !== undefined) fragments.push(...extractor.extract(block))
}
return fragments
}
private _validateVersion(type: string, version: string): void {
if (version.trim().length === 0) {
throw new SessionQueryError(`session-query extractor "${type}" requires a non-blank version`, 'SESSION_QUERY_INVALID_EXTRACTOR')
}
}
}
/**
* Encode canonical JSON with recursively sorted object keys.
* @param value - JSON-compatible source value.
* @returns deterministic JSON text.
*/
export function canonicalJson(value: unknown): string {
if (value === null || typeof value !== 'object') return JSON.stringify(value)
if (Array.isArray(value)) return `[${value.map(canonicalJson).join(',')}]`
const object = value as Record<string, unknown>
return `{${Object.keys(object).sort().map(key => `${JSON.stringify(key)}:${canonicalJson(object[key])}`).join(',')}}`
}
function normalizeText(fragments: readonly string[]): string {
return fragments.map(fragment => fragment.trim()).filter(Boolean).join('\n')
}
function collectBlockVersions(
value: unknown,
extractors: ReadonlyMap<ContentBlockType, StoredContentExtractor>,
versions: Set<string>,
): void {
if (Array.isArray(value)) {
for (const item of value) collectBlockVersions(item, extractors, versions)
return
}
if (value === null || typeof value !== 'object') return
const object = value as Record<string, unknown>
if (typeof object.type === 'string') {
const type = object.type as ContentBlockType
const extractor = extractors.get(type)
if (extractor !== undefined) versions.add(`${type}@${extractor.version}`)
}
for (const nested of Object.values(object)) collectBlockVersions(nested, extractors, versions)
}
function cloneRecord(record: SessionRecord): SessionRecord {
return { ...record, header: structuredClone(record.header) }
}

View File

@@ -1,123 +0,0 @@
/** Pure serializable session-query result filters. */
import { assertNever } from '@deepseek-ai/dsh-llm'
import type {
SessionEventRecord,
SessionEventResultFilter,
SessionQueryRange,
SessionRecord,
SessionResultFilter,
} from './types.ts'
import { SessionQueryError } from './config.ts'
const AVAILABILITIES = ['live', 'persisted'] as const
const SURFACE_STATES = ['current', 'shadowed', 'log-only'] as const
/**
* Apply an ordered AND-chain of session filters while preserving item order
* and the concrete generic item type.
* @param results - session records or richer session search hits.
* @param filters - serializable filters applied in order.
* @returns a fresh filtered array.
*/
export function filterSessionResults<T extends SessionRecord>(
results: readonly T[],
filters: readonly SessionResultFilter[],
): T[] {
for (const filter of filters) validateSessionFilter(filter)
return results.filter(result => filters.every(filter => matchesSessionFilter(result, filter)))
}
/**
* Apply an ordered AND-chain of event filters while preserving item order and
* the concrete generic item type.
* @param results - event records or richer event search hits.
* @param filters - serializable filters applied in order.
* @returns a fresh filtered array.
*/
export function filterEventResults<T extends SessionEventRecord>(
results: readonly T[],
filters: readonly SessionEventResultFilter[],
): T[] {
for (const filter of filters) validateEventFilter(filter)
return results.filter(result => filters.every(filter => matchesEventFilter(result, filter)))
}
function matchesSessionFilter(record: SessionRecord, filter: SessionResultFilter): boolean {
switch (filter.kind) {
case 'id': return filter.values.includes(record.header.id)
case 'cwd': return filter.values.includes(record.header.cwd ?? null)
case 'created-at': return inRange(record.header.createdAt, filter.range)
case 'parent': return filter.values.includes(record.header.parentSession ?? null)
case 'availability': return filter.values.some(value => value === 'live' ? record.live : record.persisted)
/* v8 ignore next -- closed discriminated union exhaustiveness guard */
default: return assertNever(filter)
}
}
function matchesEventFilter(record: SessionEventRecord, filter: SessionEventResultFilter): boolean {
switch (filter.kind) {
case 'seq': return inRange(record.seq, filter.range)
case 'time': return inRange(record.time, filter.range)
case 'type': return filter.values.includes(record.type)
case 'surface': return filter.values.includes(record.surface)
/* v8 ignore next -- closed discriminated union exhaustiveness guard */
default: return assertNever(filter)
}
}
function validateSessionFilter(filter: SessionResultFilter): void {
switch (filter.kind) {
case 'id':
case 'cwd':
case 'parent':
return
case 'created-at':
validateRange('created-at', filter.range)
return
case 'availability':
for (const value of filter.values) {
if (!(AVAILABILITIES as readonly string[]).includes(value)) invalidFilter(`unknown availability "${value}"`)
}
return
/* v8 ignore next -- closed discriminated union exhaustiveness guard */
default:
assertNever(filter)
}
}
function validateEventFilter(filter: SessionEventResultFilter): void {
switch (filter.kind) {
case 'seq':
case 'time':
validateRange(filter.kind, filter.range)
return
case 'type':
return
case 'surface':
for (const value of filter.values) {
if (!(SURFACE_STATES as readonly string[]).includes(value)) invalidFilter(`unknown surface status "${value}"`)
}
return
/* v8 ignore next -- closed discriminated union exhaustiveness guard */
default:
assertNever(filter)
}
}
function validateRange(name: string, range: SessionQueryRange): void {
if (range.from !== undefined && !Number.isFinite(range.from)) invalidFilter(`${name}.from must be finite`)
if (range.to !== undefined && !Number.isFinite(range.to)) invalidFilter(`${name}.to must be finite`)
if (range.from !== undefined && range.to !== undefined && range.from > range.to) {
invalidFilter(`${name}.from must be <= ${name}.to`)
}
}
function invalidFilter(message: string): never {
throw new SessionQueryError(`session-query filter: ${message}`, 'SESSION_QUERY_INVALID_FILTER')
}
function inRange(value: number, range: SessionQueryRange): boolean {
return (range.from === undefined || value >= range.from)
&& (range.to === undefined || value <= range.to)
}

View File

@@ -1,53 +1,29 @@
/**
* Provider-neutral session-history retrieval over live and optionally
* persisted session logs. The public service composes logical-corpus reads,
* pure filters and tracing, semantic extraction, and provider coordination.
* Exact session-history reads over live and optionally persisted logs.
*
* @module @deepseek-ai/dsh-session-query
*/
import { Context, Service } from 'cordis'
import z from 'schemastery'
import type { ContentBlockType } from '@deepseek-ai/dsh-llm'
import type { SessionEventType, SessionId } from '@deepseek-ai/dsh-session'
import { foldSurface } from '@deepseek-ai/dsh-session'
import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session'
import type {
SessionContentTextExtractor,
SessionEventReadRequest,
SessionEventRecord,
SessionEventSearchHit,
SessionEventSearchRequest,
SessionEventTextExtractor,
SessionEventTrace,
SessionEventWindow,
SessionLineageTrace,
SessionRecord,
SessionSearchHit,
SessionSearchPage,
SessionSearchProvider,
SessionSearchRequest,
SessionQueryExecContext,
} from './types.ts'
import {
SESSION_QUERY_DEFAULT_LIMIT,
SESSION_QUERY_MAX_LIMIT,
SESSION_QUERY_READ_WINDOW_MAX,
SessionQueryError,
type Config,
} from './config.ts'
import { SessionTextExtractors } from './extraction.ts'
import { SessionCorpus } from './corpus.ts'
import { SessionProviderCoordinator } from './provider.ts'
import { eventRecords, traceEventLog, traceLineage } from './tracing.ts'
export type * from './types.ts'
export type { Config, SessionQueryErrorCode } from './config.ts'
export {
SESSION_QUERY_DEFAULT_LIMIT,
SESSION_QUERY_MAX_LIMIT,
SESSION_QUERY_READ_WINDOW_MAX,
SessionQueryError,
} from './config.ts'
export { filterEventResults, filterSessionResults } from './filters.ts'
export { SESSION_QUERY_READ_WINDOW_MAX, SessionQueryError } from './config.ts'
declare module 'cordis' {
interface Context {
@@ -55,35 +31,25 @@ declare module 'cordis' {
}
}
/** Session-history retrieval and provider coordination service. */
/** Live-preferred logical-corpus and exact-event read service. */
export class SessionQueryService extends Service {
static inject = ['sessions']
static Config: z<Config> = z.object({
searchProvider: z.string(),
defaultLimit: z.number().step(1).min(1).default(SESSION_QUERY_DEFAULT_LIMIT),
maxLimit: z.number().step(1).min(1).default(SESSION_QUERY_MAX_LIMIT),
readWindowMax: z.number().step(1).min(0).default(SESSION_QUERY_READ_WINDOW_MAX),
})
private readonly _readWindowMax: number
private readonly _extractors: SessionTextExtractors
private readonly _providers: SessionProviderCoordinator
private readonly _corpus: SessionCorpus
constructor(ctx: Context, config: Config = {}) {
super(ctx, 'sessionQuery')
const defaultLimit = config.defaultLimit ?? SESSION_QUERY_DEFAULT_LIMIT
const maxLimit = config.maxLimit ?? SESSION_QUERY_MAX_LIMIT
this._readWindowMax = config.readWindowMax ?? SESSION_QUERY_READ_WINDOW_MAX
if (defaultLimit > maxLimit) {
throw new SessionQueryError('session-query: defaultLimit must be <= maxLimit', 'SESSION_QUERY_INVALID_CONFIG')
if (!Number.isInteger(this._readWindowMax) || this._readWindowMax < 0) {
throw new SessionQueryError(
'session-query: readWindowMax must be a non-negative integer',
'SESSION_QUERY_INVALID_CONFIG',
)
}
this._extractors = new SessionTextExtractors()
this._providers = new SessionProviderCoordinator({
...config.searchProvider !== undefined ? { searchProvider: config.searchProvider } : {},
defaultLimit,
maxLimit,
}, () => this._corpus, this._extractors)
this._corpus = new SessionCorpus(ctx)
}
@@ -101,7 +67,7 @@ export class SessionQueryService extends Service {
* @returns event records in ascending seq order.
*/
async listEvents(sessionId: SessionId): Promise<SessionEventRecord[]> {
const loaded = await this._corpus.loadLogical(sessionId)
const loaded = await this._corpus.load(sessionId)
return eventRecords(sessionId, loaded.events)
}
@@ -113,113 +79,58 @@ export class SessionQueryService extends Service {
async readEvent(request: SessionEventReadRequest): Promise<SessionEventWindow> {
const before = this._readWindow('before', request.before)
const after = this._readWindow('after', request.after)
const loaded = await this._corpus.loadLogical(request.sessionId)
const loaded = await this._corpus.load(request.sessionId)
const target = loaded.events[request.seq]
if (target === undefined || target.seq !== request.seq) {
throw new SessionQueryError(`session "${request.sessionId}" has no event at seq ${request.seq}`, 'SESSION_QUERY_EVENT_NOT_FOUND')
throw new SessionQueryError(
`session "${request.sessionId}" has no event at seq ${request.seq}`,
'SESSION_QUERY_EVENT_NOT_FOUND',
)
}
const startSeq = Math.max(0, request.seq - before)
const endSeq = Math.min(loaded.events.length - 1, request.seq + after)
return {
session: cloneRecord(loaded.record),
target: structuredClone(target),
events: loaded.events.slice(startSeq, endSeq + 1).map(event => structuredClone(event)),
session: loaded.header,
target,
events: loaded.events.slice(startSeq, endSeq + 1),
startSeq,
endSeq,
}
}
/**
* Trace parent ancestry and the complete known descendant tree of a session.
* @param sessionId - logical session id to trace.
* @returns complete or explicitly partial lineage.
*/
async traceSession(sessionId: SessionId): Promise<SessionLineageTrace> {
return traceLineage(await this._corpus.listSessions(), sessionId)
}
/**
* Trace direct provenance and surface replacement relationships for any event.
* @param sessionId - logical session containing the target.
* @param seq - target event seq.
* @returns lightweight trace with related seq links.
*/
async traceEvent(sessionId: SessionId, seq: number): Promise<SessionEventTrace> {
return traceEventLog(sessionId, (await this._corpus.loadLogical(sessionId)).events, seq)
}
/**
* Register one full-text provider with effect-scoped disposal.
* @param provider - provider and synchronization implementation.
* @returns async disposer that immediately unregisters selection and awaits accepted provider work.
*/
registerSearchProvider(provider: SessionSearchProvider): () => Promise<void> {
return this._providers.register(this.ctx, provider)
}
/**
* Register semantic text extraction for one event type.
* @param type - declaration-merged event discriminant.
* @param extractor - stable version and typed extraction callback.
* @returns disposer that removes the extractor.
*/
registerEventTextExtractor<K extends SessionEventType>(
type: K,
extractor: SessionEventTextExtractor<K>,
): () => void {
return this._extractors.registerEvent(this.ctx, type, extractor)
}
/**
* Register semantic text extraction for one content block type.
* @param type - declaration-merged content-block discriminant.
* @param extractor - stable version and typed extraction callback.
* @returns disposer that removes the extractor.
*/
registerContentTextExtractor<K extends ContentBlockType>(
type: K,
extractor: SessionContentTextExtractor<K>,
): () => void {
return this._extractors.registerContent(this.ctx, type, extractor)
}
/**
* Search the complete logical corpus and rank one result per session.
* @param request - query, pre-ranking filters, and pagination.
* @param exec - optional cancellation context.
* @returns ranked provider page.
*/
searchSessions(
request: SessionSearchRequest,
exec?: SessionQueryExecContext,
): Promise<SessionSearchPage<SessionSearchHit>> {
return this._providers.searchSessions(request, exec)
}
/**
* Search events within one logical session.
* @param request - target session, query, filters, and pagination.
* @param exec - optional cancellation context.
* @returns ranked provider page.
*/
searchEvents(
request: SessionEventSearchRequest,
exec?: SessionQueryExecContext,
): Promise<SessionSearchPage<SessionEventSearchHit>> {
return this._providers.searchEvents(request, exec)
}
private _readWindow(name: 'before' | 'after', value: number | undefined): number {
if (value === undefined) return 0
if (!Number.isInteger(value) || value < 0 || value > this._readWindowMax) {
throw new SessionQueryError(`${name} must be an integer between 0 and ${this._readWindowMax}`, 'SESSION_QUERY_INVALID_WINDOW')
throw new SessionQueryError(
`${name} must be an integer between 0 and ${this._readWindowMax}`,
'SESSION_QUERY_INVALID_WINDOW',
)
}
return value
}
}
function cloneRecord(record: SessionRecord): SessionRecord {
return { ...record, header: structuredClone(record.header) }
function eventRecords(sessionId: SessionId, events: readonly SessionEvent[]): SessionEventRecord[] {
let folded: ReturnType<typeof foldSurface>
try {
folded = foldSurface(events)
} catch (error: unknown) {
throw new SessionQueryError(
/* v8 ignore next -- foldSurface throws Error instances */
`invalid session surface: ${error instanceof Error ? error.message : 'unknown error'}`,
'SESSION_QUERY_INVALID_SURFACE',
{ cause: error },
)
}
const current = new Set(folded.nodes.map(node => node.seq))
const shadowed = new Set(folded.replacements.flatMap(replacement => replacement.shadowedSeqs))
return events.map(event => ({
sessionId,
seq: event.seq,
type: event.type,
time: event.time,
surface: current.has(event.seq) ? 'current' : shadowed.has(event.seq) ? 'shadowed' : 'log-only',
}))
}
export default SessionQueryService

View File

@@ -1,310 +0,0 @@
/** Search-provider selection, synchronization, pagination, and cancellation. */
import type { Context } from 'cordis'
import type { Session, SessionId } from '@deepseek-ai/dsh-session'
import type { SessionTextExtractors } from './extraction.ts'
import type { PersistenceView, SessionCorpus } from './corpus.ts'
import type {
SessionEventRecord,
SessionEventSearchHit,
SessionEventSearchRequest,
SessionEventSearchSpec,
SessionQueryExecContext,
SessionRecord,
SessionSearchHit,
SessionSearchPage,
SessionSearchProvider,
SessionSearchRequest,
SessionSearchSpec,
} from './types.ts'
import type { Config } from './config.ts'
import { SessionQueryError } from './config.ts'
import { filterEventResults, filterSessionResults } from './filters.ts'
interface ProviderState {
provider: SessionSearchProvider
chain: Promise<void>
liveIds: Set<SessionId>
}
/** Coordinates one selected provider against live and persisted corpus layers. */
export class SessionProviderCoordinator {
private readonly _configuredProviderId: string | undefined
private readonly _defaultLimit: number
private readonly _maxLimit: number
private readonly _providers = new Map<string, ProviderState>()
constructor(
config: Required<Pick<Config, 'defaultLimit' | 'maxLimit'>> & Pick<Config, 'searchProvider'>,
private readonly _corpus: () => SessionCorpus,
private readonly _extractors: SessionTextExtractors,
) {
this._configuredProviderId = config.searchProvider
this._defaultLimit = config.defaultLimit
this._maxLimit = config.maxLimit
}
/**
* Register one effect-scoped provider.
* @param ctx - contributing caller context.
* @param provider - provider implementation.
* @returns async disposer that deselects immediately and drains accepted work.
*/
register(ctx: Context, provider: SessionSearchProvider): () => Promise<void> {
if (this._providers.has(provider.id)) {
throw new SessionQueryError(`a session-query provider with id "${provider.id}" is already registered`, 'SESSION_QUERY_DUPLICATE_PROVIDER')
}
const state: ProviderState = {
provider,
chain: Promise.resolve(),
liveIds: new Set(),
}
const dispose = ctx.effect(function* (this: SessionProviderCoordinator) {
this._providers.set(provider.id, state)
yield async () => {
this._providers.delete(provider.id)
await state.chain
}
}.bind(this), 'sessionQuery.registerSearchProvider()')
return async () => { await dispose() }
}
/**
* Search and group the complete logical corpus.
* @param request - normalized provider-neutral request input.
* @param exec - optional cancellation controls.
* @returns ranked session page.
*/
async searchSessions(
request: SessionSearchRequest,
exec?: SessionQueryExecContext,
): Promise<SessionSearchPage<SessionSearchHit>> {
const state = this._resolveProvider()
const normalized = this._normalizeSessionSearch(request)
const work = this._runFullSearch(state, undefined, async () => {
if (exec?.signal?.aborted) throw aborted()
const result = await state.provider.searchSessions(normalized, exec)
return this._validateSearchPage(state, result, normalized.limit)
})
return waitFor(work, exec?.signal)
}
/**
* Search events within one logical session.
* @param request - target and provider-neutral request input.
* @param exec - optional cancellation controls.
* @returns ranked event page.
*/
async searchEvents(
request: SessionEventSearchRequest,
exec?: SessionQueryExecContext,
): Promise<SessionSearchPage<SessionEventSearchHit>> {
const state = this._resolveProvider()
const normalized = this._normalizeEventSearch(request)
const query = async (): Promise<SessionSearchPage<SessionEventSearchHit>> => {
if (exec?.signal?.aborted) throw aborted()
const result = await state.provider.searchEvents(normalized, exec)
return this._validateSearchPage(state, result, normalized.limit)
}
const live = this._corpus().getLive(request.sessionId)
let work: Promise<SessionSearchPage<SessionEventSearchHit>>
if (live !== undefined) {
work = this._runLiveSearch(state, live, query)
} else {
work = this._runFullSearch(state, request.sessionId, query)
}
return waitFor(work, exec?.signal)
}
private _runFullSearch<T>(
state: ProviderState,
requiredSessionId: SessionId | undefined,
query: () => Promise<T>,
): Promise<T> {
const liveSessions = this._corpus().listLive()
return this._serialize(state, async () => {
await this._synchronize(state, async () => {
const persistence = await this._corpus().persistenceView()
const missingRequired = requiredSessionId !== undefined
&& (persistence === undefined || !persistence.headers.some(header => header.id === requiredSessionId))
if (missingRequired) {
throw new SessionQueryError(`session "${requiredSessionId}" not found`, 'SESSION_QUERY_SESSION_NOT_FOUND')
}
if (persistence === undefined) {
await state.provider.setPersistedActive(false)
} else {
await this._syncPersisted(state, persistence)
}
await this._replaceLiveCorpus(state, liveSessions)
})
return query()
})
}
private async _syncPersisted(state: ProviderState, persistence: PersistenceView): Promise<void> {
await state.provider.setPersistedActive(false)
const inventory = new Map((await state.provider.persistedInventory()).map(entry => [entry.sessionId, entry.fingerprint]))
for (const header of persistence.headers) {
const snapshot = this._extractors.buildSnapshot(await persistence.load(header.id))
if (inventory.get(header.id) !== snapshot.fingerprint) await state.provider.replacePersisted(snapshot)
inventory.delete(header.id)
}
for (const staleId of inventory.keys()) await state.provider.removePersisted(staleId)
await state.provider.setPersistedActive(true)
}
private async _replaceLiveCorpus(state: ProviderState, sessions: readonly Session[]): Promise<void> {
const liveIds = new Set(sessions.map(session => session.id))
for (const staleId of state.liveIds) {
if (!liveIds.has(staleId)) await state.provider.removeLive(staleId)
}
for (const session of sessions) {
await state.provider.replaceLive(this._snapshotLive(session))
}
state.liveIds = liveIds
}
private _runLiveSearch<T>(state: ProviderState, session: Session, query: () => Promise<T>): Promise<T> {
let snapshot: ReturnType<SessionTextExtractors['buildSnapshot']>
try {
snapshot = this._snapshotLive(session)
} catch (error: unknown) {
return Promise.reject(this._synchronizationError(state, error))
}
return this._serialize(state, async () => {
await this._synchronize(state, async () => {
await state.provider.replaceLive(snapshot)
state.liveIds.add(session.id)
})
return query()
})
}
private _snapshotLive(session: Session): ReturnType<SessionTextExtractors['buildSnapshot']> {
return this._extractors.buildSnapshot(this._corpus().snapshotLive(session))
}
/** Serialize reconciliation and its provider query as one stable transaction. */
private _serialize<T>(state: ProviderState, operation: () => Promise<T>): Promise<T> {
const next = state.chain.then(operation, operation)
state.chain = next.then(() => undefined, () => undefined)
return next
}
/** Translate only derived-index update failures, never provider query failures. */
private async _synchronize(state: ProviderState, operation: () => Promise<void>): Promise<void> {
try {
await operation()
} catch (error: unknown) {
throw this._synchronizationError(state, error)
}
}
private _synchronizationError(state: ProviderState, error: unknown): SessionQueryError {
/* v8 ignore next -- service-created typed synchronization errors pass through unchanged */
if (error instanceof SessionQueryError) return error
return new SessionQueryError(`session-query provider "${state.provider.id}" synchronization failed: ${errorMessage(error)}`, 'SESSION_QUERY_INDEX_FAILED', { cause: error })
}
private _resolveProvider(): ProviderState {
if (this._configuredProviderId !== undefined) {
const state = this._providers.get(this._configuredProviderId)
if (state === undefined) {
throw new SessionQueryError(`configured session-query provider "${this._configuredProviderId}" is not registered`, 'SESSION_QUERY_PROVIDER_CONFIGURED_MISSING')
}
if (!state.provider.status().available) {
throw new SessionQueryError(`configured session-query provider "${this._configuredProviderId}" is unavailable`, 'SESSION_QUERY_PROVIDER_CONFIGURED_UNAVAILABLE')
}
return state
}
const usable = [...this._providers.values()].filter(state => state.provider.status().available)
const [single] = usable
if (single === undefined) {
throw new SessionQueryError('no usable session-query provider is registered', 'SESSION_QUERY_PROVIDER_UNAVAILABLE')
}
if (usable.length > 1) {
throw new SessionQueryError(`multiple usable session-query providers are registered (${usable.map(state => state.provider.id).join(', ')}); configure one explicitly`, 'SESSION_QUERY_PROVIDER_AMBIGUOUS')
}
return single
}
private _normalizeSessionSearch(request: SessionSearchRequest): SessionSearchSpec {
const query = this._queryText(request.query)
const limit = this._limitValue(request.limit)
filterSessionResults<SessionRecord>([], request.sessionFilters ?? [])
filterEventResults<SessionEventRecord>([], request.eventFilters ?? [])
return { ...request, query, limit }
}
private _normalizeEventSearch(request: SessionEventSearchRequest): SessionEventSearchSpec {
const query = this._queryText(request.query)
const limit = this._limitValue(request.limit)
filterEventResults<SessionEventRecord>([], request.filters ?? [])
return { ...request, query, limit }
}
private _queryText(query: string): string {
const normalized = query.trim()
if (normalized.length === 0) {
throw new SessionQueryError('session-query search text must not be blank', 'SESSION_QUERY_INVALID_QUERY')
}
return normalized
}
private _limitValue(limit: number | undefined): number {
const value = limit ?? this._defaultLimit
if (!Number.isInteger(value) || value < 1 || value > this._maxLimit) {
throw new SessionQueryError(`session-query limit must be an integer between 1 and ${this._maxLimit}`, 'SESSION_QUERY_INVALID_LIMIT')
}
return value
}
private _validateSearchPage<T>(state: ProviderState, page: SessionSearchPage<T>, limit: number): SessionSearchPage<T> {
if (page.providerId !== state.provider.id) {
throw new SessionQueryError(`session-query provider "${state.provider.id}" returned providerId "${page.providerId}"`, 'SESSION_QUERY_PROVIDER_ERROR')
}
if (page.items.length > limit) {
throw new SessionQueryError(`session-query provider "${state.provider.id}" returned ${page.items.length} items for limit ${limit}`, 'SESSION_QUERY_PROVIDER_ERROR')
}
return page
}
}
function waitFor<T>(work: Promise<T>, signal: AbortSignal | undefined): Promise<T> {
const observed = work.catch((error: unknown) => { throw operationError(error) })
if (signal === undefined) return observed
if (signal.aborted) {
// Cancellation supersedes the caller's result, but shared work must still
// have a rejection observer when it has already failed synchronously.
void observed.catch((_supersededError: unknown) => undefined)
return Promise.reject(aborted())
}
return new Promise<T>((resolve, reject) => {
const onAbort = () => { reject(aborted()) }
signal.addEventListener('abort', onAbort, { once: true })
observed.then(
(value) => {
signal.removeEventListener('abort', onAbort)
resolve(value)
},
(error: unknown) => {
signal.removeEventListener('abort', onAbort)
reject(operationError(error))
},
)
})
}
function operationError(error: unknown): Error {
if (error instanceof Error) return error
return new SessionQueryError('session-query operation failed with a non-Error rejection', 'SESSION_QUERY_PROVIDER_ERROR', { cause: error })
}
function aborted(): SessionQueryError {
return new SessionQueryError('session-query operation aborted', 'SESSION_QUERY_ABORTED')
}
function errorMessage(error: unknown): string {
/* v8 ignore next -- provider update contracts reject Error instances */
return error instanceof Error ? error.message : 'unknown error'
}

View File

@@ -1,158 +0,0 @@
/** Session lineage and event surface/provenance tracing. */
import { foldSurface, isSurfaceEvent } from '@deepseek-ai/dsh-session'
import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session'
import type {
SessionEventRecord,
SessionEventTrace,
SessionLineageNode,
SessionLineageTrace,
SessionRecord,
} from './types.ts'
import { SessionQueryError } from './config.ts'
/**
* Classify raw events against the canonical surface fold.
* @param sessionId - owner of the event log.
* @param events - detached raw log.
* @returns lightweight records in seq order.
*/
export function eventRecords(sessionId: SessionId, events: readonly SessionEvent[]): SessionEventRecord[] {
const fold = safeFold(events)
const current = new Set(fold.nodes.map(node => node.seq))
const shadowed = new Set(fold.replacements.flatMap(replacement => replacement.shadowedSeqs))
return events.map(event => ({
sessionId,
seq: event.seq,
type: event.type,
time: event.time,
surface: current.has(event.seq) ? 'current' : shadowed.has(event.seq) ? 'shadowed' : 'log-only',
}))
}
/**
* Build one event trace from a validated logical event log.
* @param sessionId - owner of the event log.
* @param events - detached raw log.
* @param seq - target event seq.
* @returns direct provenance and replacement relationships.
*/
export function traceEventLog(sessionId: SessionId, events: readonly SessionEvent[], seq: number): SessionEventTrace {
const target = events[seq]
if (target === undefined || target.seq !== seq) {
throw new SessionQueryError(`session "${sessionId}" has no event at seq ${seq}`, 'SESSION_QUERY_EVENT_NOT_FOUND')
}
const records = eventRecords(sessionId, events)
const fold = safeFold(events)
const shadowedBy = new Map<number, number>()
const shadows = new Map<number, number[]>()
for (const replacement of fold.replacements) {
shadows.set(replacement.seq, [...replacement.shadowedSeqs])
for (const shadowed of replacement.shadowedSeqs) shadowedBy.set(shadowed, replacement.seq)
}
const references: number[] = []
const referencedBy: number[] = []
for (const event of events) {
if (!isSurfaceEvent(event)) continue
for (const source of event.sourceEventSeqs ?? []) {
if (event.seq === seq) references.push(source)
if (source === seq) referencedBy.push(event.seq)
}
}
const replacementChain: number[] = []
let replacement = shadowedBy.get(seq)
while (replacement !== undefined) {
replacementChain.push(replacement)
replacement = shadowedBy.get(replacement)
}
const immediate = shadowedBy.get(seq)
// seq was checked against the contiguous event log, so its parallel record exists.
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
const targetRecord = records[seq]!
return {
target: { ...targetRecord },
...immediate !== undefined ? { shadowedBy: immediate } : {},
replacementChain,
shadows: shadows.get(seq) ?? [],
references,
referencedBy,
}
}
/**
* Trace ancestry and descendants within one materialized logical corpus.
* @param records - complete visible logical corpus.
* @param sessionId - target session id.
* @returns complete known lineage or explicit unresolved parent.
*/
export function traceLineage(records: readonly SessionRecord[], sessionId: SessionId): SessionLineageTrace {
const byId = new Map(records.map(record => [record.header.id, record]))
const target = byId.get(sessionId)
if (target === undefined) {
throw new SessionQueryError(`session "${sessionId}" not found`, 'SESSION_QUERY_SESSION_NOT_FOUND')
}
const parents: SessionRecord[] = []
const ancestrySeen = new Set<SessionId>([sessionId])
let unresolvedParentId: SessionId | undefined
let parentId = target.header.parentSession
while (parentId !== undefined) {
if (ancestrySeen.has(parentId)) lineageCycle(parentId)
ancestrySeen.add(parentId)
const parent = byId.get(parentId)
if (parent === undefined) {
unresolvedParentId = parentId
break
}
parents.push(parent)
parentId = parent.header.parentSession
}
const childrenByParent = new Map<SessionId, SessionRecord[]>()
for (const record of records) {
const parent = record.header.parentSession
if (parent === undefined) continue
const children = childrenByParent.get(parent) ?? []
children.push(record)
childrenByParent.set(parent, children)
}
for (const children of childrenByParent.values()) children.sort(compareSessionsAscending)
const buildChildren = (id: SessionId): SessionLineageNode[] => (childrenByParent.get(id) ?? []).map(child => ({
session: cloneRecord(child),
children: buildChildren(child.header.id),
}))
return {
target: cloneRecord(target),
parents: parents.map(cloneRecord),
...unresolvedParentId !== undefined
? { unresolvedParentId }
: { root: cloneRecord(parents.at(-1) ?? target) },
children: buildChildren(sessionId),
}
}
function safeFold(events: readonly SessionEvent[]): ReturnType<typeof foldSurface> {
try {
return foldSurface(events)
} catch (error: unknown) {
throw new SessionQueryError(`invalid session surface: ${errorMessage(error)}`, 'SESSION_QUERY_INVALID_SURFACE', { cause: error })
}
}
function cloneRecord(record: SessionRecord): SessionRecord {
return { ...record, header: structuredClone(record.header) }
}
function compareSessionsAscending(a: SessionRecord, b: SessionRecord): number {
return a.header.createdAt - b.header.createdAt || a.header.id.localeCompare(b.header.id)
}
function lineageCycle(id: SessionId): never {
throw new SessionQueryError(`session lineage contains a cycle at "${id}"`, 'SESSION_QUERY_INVALID_LINEAGE')
}
function errorMessage(error: unknown): string {
/* v8 ignore next -- foldSurface throws Error instances */
return error instanceof Error ? error.message : 'unknown error'
}

View File

@@ -1,31 +1,24 @@
/**
* Public vocabulary for the session-query retrieval service: lightweight
* records, composable filters, traces, search requests/results, extractor
* registrations, and the provider synchronization contract.
* Public records for exact reads over the live-preferred logical session corpus.
*
* @module @deepseek-ai/dsh-session-query/types
*/
import type { ContentBlockMap, ContentBlockType } from '@deepseek-ai/dsh-llm'
import type {
SessionEvent,
SessionEventType,
SessionHeader,
SessionId,
} from '@deepseek-ai/dsh-session'
import type { SessionEvent, SessionEventType, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
/** Whether an event is on the current surface, was replaced, or is log-only. */
/** Whether an event is current model context, replaced context, or raw-log-only. */
export type SessionEventSurface = 'current' | 'shadowed' | 'log-only'
/** Lightweight identity and availability for one logical session. */
/** Lightweight identity and source availability for one logical session. */
export interface SessionRecord {
/** Cloned immutable session header selected from the live-preferred corpus. */
/** Cloned session header selected from the live-preferred corpus. */
header: SessionHeader
/** Whether the id currently exists in `ctx.sessions`. */
live: boolean
/** Whether the active persistence backend currently materializes the id. */
persisted: boolean
}
/** Lightweight metadata for one event within a logical session. */
export interface SessionEventRecord {
/** Session that owns the event. */
@@ -40,102 +33,6 @@ export interface SessionEventRecord {
surface: SessionEventSurface
}
/** Inclusive numeric range used by result and search filters. */
export interface SessionQueryRange {
/** Inclusive lower bound. */
from?: number
/** Inclusive upper bound. */
to?: number
}
/** Serializable filter applied to session records. */
export type SessionResultFilter =
| { kind: 'id'; values: readonly SessionId[] }
| { kind: 'cwd'; values: readonly (string | null)[] }
| { kind: 'created-at'; range: SessionQueryRange }
| { kind: 'parent'; values: readonly (SessionId | null)[] }
| { kind: 'availability'; values: readonly ('live' | 'persisted')[] }
/** Serializable filter applied to event records. */
export type SessionEventResultFilter =
| { kind: 'seq'; range: SessionQueryRange }
| { kind: 'time'; range: SessionQueryRange }
| { kind: 'type'; values: readonly SessionEventType[] }
| { kind: 'surface'; values: readonly SessionEventSurface[] }
/** Caller cancellation threaded through synchronization and provider search. */
export interface SessionQueryExecContext {
/** Abort signal for waiting and provider-owned query work. */
readonly signal?: AbortSignal
}
/** Cheap local usability status returned by a search provider. */
export type SessionSearchProviderStatus =
| { readonly available: true }
| { readonly available: false; readonly reason: 'misconfigured' | 'unavailable' }
/** Common pagination fields accepted by both search scopes. */
export interface SessionSearchPageRequest {
/** Maximum number of hits on this page. */
limit?: number
/** Opaque cursor returned by the same provider/request. */
cursor?: string
}
/** Cross-session full-text request. */
export interface SessionSearchRequest extends SessionSearchPageRequest {
/** Plain text query interpreted by the selected provider. */
query: string
/** Session metadata filters applied before event ranking/grouping. */
sessionFilters?: readonly SessionResultFilter[]
/** Event metadata filters applied before best-event grouping. */
eventFilters?: readonly SessionEventResultFilter[]
}
/** Full-text request scoped to one session's events. */
export interface SessionEventSearchRequest extends SessionSearchPageRequest {
/** Session whose events form the search corpus. */
sessionId: SessionId
/** Plain text query interpreted by the selected provider. */
query: string
/** Event metadata filters applied before ranking. */
filters?: readonly SessionEventResultFilter[]
}
/** Provider-facing cross-session search spec after service normalization. */
export interface SessionSearchSpec extends SessionSearchRequest {
/** Required page size validated and defaulted by the query service. */
limit: number
}
/** Provider-facing event search spec after service normalization. */
export interface SessionEventSearchSpec extends SessionEventSearchRequest {
/** Required page size validated and defaulted by the query service. */
limit: number
}
/** One lightweight event search hit with provider-produced evidence text. */
export interface SessionEventSearchHit extends SessionEventRecord {
/** Plain-text excerpt explaining the match. */
snippet: string
}
/** One session-ranked search hit and its strongest matching event. */
export interface SessionSearchHit extends SessionRecord {
/** Strongest matching event used as the session's ranking evidence. */
bestMatch: SessionEventSearchHit
}
/** One provider-owned page of search results. */
export interface SessionSearchPage<T> {
/** Stable id of the provider that produced this page. */
providerId: string
/** Ranked hits in deterministic provider order, no longer than the requested limit. */
items: readonly T[]
/** Opaque next-page cursor, absent when the result is exhausted. */
nextCursor?: string
}
/** Request for one event plus raw neighboring log context. */
export interface SessionEventReadRequest {
/** Session that owns the target event. */
@@ -150,8 +47,8 @@ export interface SessionEventReadRequest {
/** Full target event and a bounded raw-log window. */
export interface SessionEventWindow {
/** Logical session metadata at read time. */
session: SessionRecord
/** Cloned header for the live-preferred source read. */
session: SessionHeader
/** Full cloned target event. */
target: SessionEvent
/** Full cloned events from `startSeq` through `endSeq`. */
@@ -161,144 +58,3 @@ export interface SessionEventWindow {
/** Last seq included in `events`. */
endSeq: number
}
/** Recursive child node in a session lineage trace. */
export interface SessionLineageNode {
/** Session represented by this lineage node. */
session: SessionRecord
/** Direct children in deterministic creation order. */
children: SessionLineageNode[]
}
/** Complete known lineage around one session. */
export interface SessionLineageTrace {
/** Session that was traced. */
target: SessionRecord
/** Known parents from immediate parent outward. */
parents: SessionRecord[]
/** Root when the complete parent chain is available. */
root?: SessionRecord
/** First parent id outside the visible corpus, when the trace is partial. */
unresolvedParentId?: SessionId
/** Complete known descendant forest rooted at the target's direct children. */
children: SessionLineageNode[]
}
/** Surface and provenance relationships for one event. */
export interface SessionEventTrace {
/** Lightweight target record. */
target: SessionEventRecord
/** Immediate replacement event that shadowed the target. */
shadowedBy?: number
/** Replacement seqs from the target toward the current descendant. */
replacementChain: number[]
/** Surface nodes directly shadowed by the target replacement event. */
shadows: number[]
/** Direct provenance sources from `sourceEventSeqs`. */
references: number[]
/** Events that directly name the target in `sourceEventSeqs`. */
referencedBy: number[]
}
/** Typed extractor for one declaration-merged session event type. */
export interface SessionEventTextExtractor<K extends SessionEventType = SessionEventType> {
/** Stable cache-invalidation version chosen by the extractor owner. */
version: string
/**
* Extract semantic searchable fragments from one event.
* @param event - event narrowed to the registered type.
* @returns plain-text fragments; blanks are discarded by the service.
*/
extract(event: SessionEvent<K>): readonly string[]
}
/** Typed extractor for one declaration-merged content block type. */
export interface SessionContentTextExtractor<K extends ContentBlockType = ContentBlockType> {
/** Stable cache-invalidation version chosen by the extractor owner. */
version: string
/**
* Extract semantic searchable fragments from one content block.
* @param block - block narrowed to the registered type.
* @returns plain-text fragments; blanks are discarded by the service.
*/
extract(block: ContentBlockMap[K]): readonly string[]
}
/** One provider-neutral event document produced by registered extractors. */
export interface SessionIndexDocument extends SessionEventRecord {
/** Normalized newline-joined text indexed by a search provider. */
text: string
}
/** One complete index layer for a live session or persisted checkpoint. */
export interface SessionIndexSnapshot {
/** Layer metadata and live/persisted availability exposed in results. */
session: SessionRecord
/** Stable SHA-256 identity of canonical source data and extractor versions. */
fingerprint: string
/** Searchable event documents in seq order. */
documents: readonly SessionIndexDocument[]
}
/** Durable provider inventory entry used to reuse unchanged persisted rows. */
export interface SessionPersistedIndexEntry {
/** Persisted session id. */
sessionId: SessionId
/** Last indexed source/extractor fingerprint. */
fingerprint: string
}
/** Search and synchronization backend registered into `ctx.sessionQuery`. */
export interface SessionSearchProvider {
/** Stable provider id, unique within the query service. */
readonly id: string
/**
* Return cheap local usability without performing index or search I/O.
* @returns whether the provider can be selected.
*/
status(): SessionSearchProviderStatus
/**
* Read reusable persisted-layer fingerprints from derived storage.
* @returns durable inventory entries.
*/
persistedInventory(): Promise<readonly SessionPersistedIndexEntry[]>
/**
* Hide or expose reconciled persisted rows without deleting their cache.
* @param active - whether canonical persistence is mounted and reconciled.
*/
setPersistedActive(active: boolean): Promise<void>
/**
* Atomically replace one persisted session's derived documents.
* @param snapshot - canonical persisted checkpoint and fingerprint.
*/
replacePersisted(snapshot: SessionIndexSnapshot): Promise<void>
/**
* Delete one durable derived entry after canonical reconciliation proves it absent.
* @param sessionId - persisted id to remove.
*/
removePersisted(sessionId: SessionId): Promise<void>
/**
* Replace one connection-local live override.
* @param snapshot - current live snapshot and availability.
*/
replaceLive(snapshot: SessionIndexSnapshot): Promise<void>
/**
* Drop one live override, revealing its active persisted base when present.
* @param sessionId - live id to remove.
*/
removeLive(sessionId: SessionId): Promise<void>
/**
* Search and group the complete logical corpus by session.
* @param request - query, pre-ranking filters, and pagination.
* @param exec - optional cancellation context.
* @returns one ranked session page.
*/
searchSessions(request: SessionSearchSpec, exec?: SessionQueryExecContext): Promise<SessionSearchPage<SessionSearchHit>>
/**
* Search events within one logical session.
* @param request - target session, query, filters, and pagination.
* @param exec - optional cancellation context.
* @returns one ranked event page.
*/
searchEvents(request: SessionEventSearchSpec, exec?: SessionQueryExecContext): Promise<SessionSearchPage<SessionEventSearchHit>>
}