/** * Workspace entity registry (`ctx.workspace`): durable workspace records, * stable registry order, and header-validated session membership over the * domain data form. * @module @deepseek-ai/dsh-workspace */ import { randomUUID } from 'node:crypto' import { stat } from 'node:fs/promises' import { basename } from 'node:path' import { Context, Service } from 'cordis' import type { SessionHeader, SessionId } from '@deepseek-ai/dsh-session' import type {} from '@deepseek-ai/dsh-session-persistence' import type { DomainGlobal, KvTable } from '@deepseek-ai/dsh-storage-domain' import { WorkspaceEntity } from './entity.ts' import type { WorkspaceEntityHost } from './entity.ts' export { WorkspaceMoveInvalidError } from './entity.ts' import { realpathNormalize } from './paths.ts' import { workspaceDomainSpec } from './spec.ts' import type { WorkspaceDomainState, WorkspaceRecord } from './spec.ts' import type { Workspace, WorkspaceId as WorkspaceIdBrand } from './types.ts' export type { Workspace } from './types.ts' export { workspaceDomainState, workspaceRecord, workspaceDomainSpec } from './spec.ts' export type { WorkspaceDomainState, WorkspaceRecord } from './spec.ts' export { realpathNormalize } from './paths.ts' /** Identifies one workspace record (see `src/types.ts` for the brand rationale). */ export type WorkspaceId = WorkspaceIdBrand /** * Brand a string as a {@link WorkspaceId}. * @param id - Raw workspace id string. * @returns the same string, branded at compile time. */ export function WorkspaceId(id: string): WorkspaceId { return id as WorkspaceId } /** A create request would give two Workspaces the same display name. */ export class WorkspaceNameConflictError extends Error { /** * @param workspaceName - Conflicting display name. */ constructor(readonly workspaceName: string) { super(`workspace name '${workspaceName}' is already in use`) this.name = 'WorkspaceNameConflictError' } } declare module 'cordis' { interface Context { workspace: WorkspaceRegistry } } interface BootstrapGroup { readonly path: string readonly headers: SessionHeader[] readonly newestAt: number } const sameIds = (left: readonly WorkspaceId[], right: readonly WorkspaceId[]): boolean => left.length === right.length && left.every((id, index) => id === right[index]) const compareHeaders = (left: SessionHeader, right: SessionHeader): number => right.createdAt - left.createdAt || String(left.id).localeCompare(String(right.id)) /** * Durable workspace registry. Startup waits for `sessionPersistence`, builds * one canonical-cwd header index, and completes the one-time history * bootstrap before the service becomes active. The persistence dependency is * mandatory so an unavailable peer can never be mistaken for an empty * history and commit the initialized marker. */ export class WorkspaceRegistry extends Service { static inject = ['storageDomain', 'sessionPersistence'] private table?: KvTable private global?: DomainGlobal private state?: WorkspaceDomainState private readonly entities = new Map() private readonly headers = new Map() private readonly sessionPaths = new Map() private readonly invalidSessionPaths = new Map() private operationTail: Promise = Promise.resolve() private readonly host: WorkspaceEntityHost = { table: () => this.requireTable(), sessionPath: id => this.sessionPaths.get(id), readSessionHeader: id => this.readSessionHeader(id), rememberSessionPath: (id, path) => { this.sessionPaths.set(id, path) this.invalidSessionPaths.delete(id) }, } constructor(ctx: Context) { super(ctx, 'workspace') } /** Open the domain, finish bootstrap when required, and rebuild the ordered cache. */ protected async [Service.init](): Promise { const domain = await this.ctx.storageDomain.open(workspaceDomainSpec) this.ctx.effect(() => () => domain.close(), 'workspace.domainClose') this.table = domain.table('workspaces') this.global = domain.global this.state = domain.global.get() await this.recoverPendingMutation() this.validateStoredState(this.state) if (!this.state.initialized) { const headers = await this.ctx.sessionPersistence.list() await this.replaceHeaderIndex(headers) await this.bootstrap(headers) } else if (this.table.size > 0) { await this.replaceHeaderIndex(await this.ctx.sessionPersistence.list()) } await this.indexLiveSessions() this.validateStoredState(this.requireState()) this.rebuildEntities() this.reportFilteredCandidates() } /** * Create or reuse a workspace for an existing directory. The path is * canonicalized through `fs.realpath`; a nonexistent path rejects with the * original error and a non-directory rejects. Repeated calls for the same * canonical path return the existing entity without changing its title. * A newly created workspace is prepended to the durable registry order. * A different canonical path cannot create a duplicate display title. * @param path - Existing directory to own, in any path spelling. * @param title - Display title used only when a new record is created. * @returns the existing or newly durable workspace. */ async create(path: string, title?: string): Promise { const canonical = await realpathNormalize(path) if (!(await stat(canonical)).isDirectory()) { throw new Error(`cannot create a workspace at '${canonical}': path is not a directory`) } return await this.enqueueOperation(() => this.createCanonical(canonical, title)) } /** * Look up a workspace by id. * @param id - Workspace id. * @returns the workspace, or `undefined` when unknown. */ get(id: WorkspaceId): Workspace | undefined { return this.entities.get(id) } /** * Synchronous workspace projection in durable registry order. Every * entity's `sessionIds` getter is already filtered by the startup/live * canonical-cwd header index; this method performs no persistence reads. * @returns a fresh ordered array of workspace entities. */ list(): Workspace[] { return this.requireState().workspaceIds.map((id) => { const entity = this.entities.get(id) if (entity === undefined) { throw new Error(`workspace registry order references missing workspace '${id}'`) } return entity }) } /** * Delete one workspace registration while retaining its directory and every * session log. The durable order is updated before the table deletion; a * failed table write restores the prior order and keeps the entity * published. Unknown ids are an idempotent no-op for domain callers. * @param id - Workspace registration to remove. * @returns `true` when a record was deleted, `false` when it was unknown. */ delete(id: WorkspaceId): Promise { return this.enqueueOperation(() => this.deleteKnown(id)) } /** * Resolve by canonical directory path without creating or mutating a * workspace. A missing path rejects during `realpath`; an existing unowned * directory returns `undefined`. * @param path - Existing directory path in any spelling. * @returns the workspace owning the canonical path, when one exists. */ async resolveByPath(path: string): Promise { const canonical = await realpathNormalize(path) for (const entity of this.entities.values()) { if (entity.path === canonical) return entity } return undefined } private async createCanonical(canonical: string, title?: string): Promise { for (const entity of this.entities.values()) { if (entity.path === canonical) return entity } const workspaceName = title ?? basename(canonical) if ([...this.entities.values()].some(entity => entity.title === workspaceName)) { throw new WorkspaceNameConflictError(workspaceName) } const table = this.requireTable() const state = this.requireState() const id = WorkspaceId(randomUUID()) const now = new Date().toISOString() const record: WorkspaceRecord = { path: canonical, title: workspaceName, sessionIds: [], createdAt: now, updatedAt: now, } const entity = new WorkspaceEntity(this.host, id, record) this.entities.set(id, entity) const pendingState: WorkspaceDomainState = { ...state, pendingMutation: { operation: 'create', workspaceId: id }, } try { await this.setState(pendingState) } catch (error) { this.entities.delete(id) throw error } try { await table.put(id, record) } catch (error) { this.entities.delete(id) try { await this.setState(state) } catch (rollbackError) { throw new AggregateError( [error, rollbackError], `workspace '${id}' record write and pending-marker rollback both failed`, ) } throw error } try { await this.setState({ initialized: true, workspaceIds: [id, ...state.workspaceIds] }) } catch (error) { this.entities.delete(id) try { await table.delete(id) } catch (rollbackError) { throw new AggregateError( [error, rollbackError], `workspace '${id}' order write and record rollback both failed; the pending marker remains recoverable`, ) } try { await this.setState(state) } catch (rollbackError) { throw new AggregateError( [error, rollbackError], `workspace '${id}' order write and pending-marker rollback both failed`, ) } throw error } return entity } private async deleteKnown(id: WorkspaceId): Promise { const entity = this.entities.get(id) if (entity === undefined) return false const state = this.requireState() const nextState = { initialized: true, workspaceIds: state.workspaceIds.filter(workspaceId => workspaceId !== id), } await this.setState({ ...nextState, pendingMutation: { operation: 'delete', workspaceId: id }, }) this.entities.delete(id) try { await this.requireTable().delete(id) } catch (error) { this.entities.set(id, entity) try { await this.setState(state) } catch (rollbackError) { // The durable marker still says to finish deletion, so the cache must // agree with that recoverable direction rather than republish a row // absent from the persisted order. this.entities.delete(id) throw new AggregateError( [error, rollbackError], `workspace '${id}' record deletion and registry-order rollback both failed`, ) } throw error } try { await this.setState(nextState) } catch (error) { // The deletion committed at the table write and was already published // to Host streams. Keep the durable marker for startup recovery rather // than reporting failure after the requested state became true. this.ctx.logger.warn( `workspace '${id}' was deleted but its pending marker could not be cleared: ${String(error)}`, ) } return true } /** * Complete the one mutation explicitly named by durable state. Unexplained * order/table divergence still reaches {@link validateStoredState} and * fails loud; this path never infers provenance from shape alone. */ private async recoverPendingMutation(): Promise { const state = this.requireState() const pending = state.pendingMutation if (pending === undefined) return if (state.workspaceIds.includes(pending.workspaceId)) { throw new Error( `workspace domain is inconsistent: pending ${pending.operation} workspace ` + `'${pending.workspaceId}' is still present in registry order`, ) } await this.requireTable().delete(pending.workspaceId) await this.setState({ initialized: state.initialized, workspaceIds: state.workspaceIds }) } private async bootstrap(headers: readonly SessionHeader[]): Promise { const table = this.requireTable() const state = this.requireState() const groupsByPath = new Map() for (const header of headers) { const path = this.sessionPaths.get(header.id) if (path === undefined) continue const group = groupsByPath.get(path) if (group === undefined) groupsByPath.set(path, [header]) else group.push(header) } const groups: BootstrapGroup[] = [...groupsByPath].map(([path, groupHeaders]) => { groupHeaders.sort(compareHeaders) const newest = groupHeaders[0] as SessionHeader return { path, headers: groupHeaders, newestAt: newest.createdAt } }).sort((left, right) => right.newestAt - left.newestAt || left.path.localeCompare(right.path)) const byPath = new Map() const accounted = new Map() for (const [id, record] of table.entries()) { byPath.set(record.path, id) for (const sessionId of record.sessionIds) accounted.set(sessionId, id) } for (const group of groups) { let id = byPath.get(group.path) if (id === undefined) { const sessionIds = group.headers .map(header => header.id) .filter(sessionId => !accounted.has(sessionId)) if (sessionIds.length === 0) continue id = WorkspaceId(randomUUID()) const createdAt = new Date(group.newestAt).toISOString() const record: WorkspaceRecord = { path: group.path, title: basename(group.path), sessionIds, createdAt, updatedAt: createdAt, } await table.put(id, record) byPath.set(group.path, id) for (const sessionId of sessionIds) accounted.set(sessionId, id) continue } const current = table.get(id) as WorkspaceRecord const historical = group.headers .map(header => header.id) .filter(sessionId => accounted.get(sessionId) === undefined || accounted.get(sessionId) === id) const historicalSet = new Set(historical) const sessionIds = [ ...historical, ...current.sessionIds.filter(sessionId => !historicalSet.has(sessionId)), ] if (sameSessionIds(current.sessionIds, sessionIds)) continue await table.update(id, record => ({ ...record, sessionIds, updatedAt: new Date().toISOString(), })) for (const sessionId of historical) accounted.set(sessionId, id) } const groupRank = new Map(groups.map(group => [group.path, group.newestAt])) const priorRank = new Map(state.workspaceIds.map((id, index) => [id, index])) const workspaceIds = [...table.entries()] .sort(([leftId, left], [rightId, right]) => { const leftTime = groupRank.get(left.path) ?? Date.parse(left.createdAt) const rightTime = groupRank.get(right.path) ?? Date.parse(right.createdAt) return rightTime - leftTime || (priorRank.get(leftId) ?? Number.MAX_SAFE_INTEGER) - (priorRank.get(rightId) ?? Number.MAX_SAFE_INTEGER) || String(leftId).localeCompare(String(rightId)) }) .map(([id]) => id) if (!sameIds(state.workspaceIds, workspaceIds)) { await this.setState({ initialized: false, workspaceIds }) } await this.setState({ initialized: true, workspaceIds }) } private validateStoredState(state: WorkspaceDomainState): void { const table = this.requireTable() const order = new Set() for (const id of state.workspaceIds) { if (order.has(id)) { throw new Error(`workspace domain is inconsistent: registry order repeats workspace '${id}'`) } if (table.get(id) === undefined) { throw new Error(`workspace domain is inconsistent: registry order references missing workspace '${id}'`) } order.add(id) } if (state.initialized && order.size !== table.size) { const orphan = [...table.keys()].find(id => !order.has(id)) throw new Error( `workspace domain is inconsistent: workspace '${orphan as WorkspaceId}' is absent from registry order`, ) } const paths = new Map() const accounted = new Map() for (const [id, record] of table.entries()) { const pathHolder = paths.get(record.path) if (pathHolder !== undefined) { throw new Error( `workspace domain is inconsistent: path '${record.path}' is claimed ` + `by both workspace '${pathHolder}' and workspace '${id}'`, ) } paths.set(record.path, id) for (const sessionId of record.sessionIds) { const holder = accounted.get(sessionId) if (holder !== undefined) { throw new Error( `workspace domain is inconsistent: session '${sessionId}' is accounted ` + `by both workspace '${holder}' and workspace '${id}'`, ) } accounted.set(sessionId, id) } } } private rebuildEntities(): void { this.entities.clear() for (const id of this.requireState().workspaceIds) { const record = this.requireTable().get(id) as WorkspaceRecord this.entities.set(id, new WorkspaceEntity(this.host, id, record)) } } private async replaceHeaderIndex(headers: readonly SessionHeader[]): Promise { this.headers.clear() this.sessionPaths.clear() this.invalidSessionPaths.clear() await this.indexHeaders(headers) } private async indexHeaders(headers: readonly SessionHeader[]): Promise { for (const header of headers) await this.indexHeader(header) } private async indexHeader(header: SessionHeader): Promise { this.headers.set(header.id, header) this.sessionPaths.delete(header.id) if (header.cwd === undefined) { this.invalidSessionPaths.set(header.id, 'header has no cwd') return } try { const path = await realpathNormalize(header.cwd) if (!(await stat(path)).isDirectory()) { this.invalidSessionPaths.set(header.id, `cwd '${header.cwd}' is not a directory`) return } this.sessionPaths.set(header.id, path) this.invalidSessionPaths.delete(header.id) } catch { this.invalidSessionPaths.set(header.id, `cwd '${header.cwd}' does not resolve`) } } private async indexLiveSessions(): Promise { const sessions = this.ctx.get('sessions') if (sessions === undefined) return await this.indexHeaders(sessions.list().map(session => session.header)) } private reportFilteredCandidates(): void { for (const entity of this.entities.values()) { const record = this.requireTable().get(entity.id) as WorkspaceRecord for (const sessionId of record.sessionIds) { const path = this.sessionPaths.get(sessionId) if (path === record.path) continue const reason = this.invalidSessionPaths.get(sessionId) ?? (this.headers.has(sessionId) ? `canonical cwd '${path}' differs from workspace path '${record.path}'` : 'session header is missing') this.ctx.logger.warn( `workspace '${entity.id}' filtered session '${sessionId}' from membership: ${reason}`, ) } } } private async readSessionHeader(id: SessionId): Promise { const live = this.ctx.get('sessions')?.get(id) if (live !== undefined) { this.headers.set(id, live.header) return live.header } const cached = this.headers.get(id) if (cached !== undefined) return cached const headers = await this.ctx.sessionPersistence.list() await this.indexHeaders(headers) const header = this.headers.get(id) if (header === undefined) { throw new Error(`cannot validate session '${id}': session persistence holds no such session`) } return header } private requireTable(): KvTable { if (this.table === undefined) throw new Error('workspace registry is not started yet') return this.table } private requireState(): WorkspaceDomainState { if (this.state === undefined) throw new Error('workspace registry is not started yet') return this.state } private async setState(state: WorkspaceDomainState): Promise { await (this.global as DomainGlobal).set(state) this.state = state } private enqueueOperation(operation: () => Promise): Promise { const result = this.operationTail.then(async () => { // A committed delete may leave only its marker cleanup pending. Retry // recovery before another create/delete can overwrite that provenance. await this.recoverPendingMutation() return await operation() }) this.operationTail = result.then(() => {}, () => {}) return result } } const sameSessionIds = (left: readonly SessionId[], right: readonly SessionId[]): boolean => left.length === right.length && left.every((id, index) => id === right[index]) export default WorkspaceRegistry