Move the 18 flat packages/<name> packages into role-grouped dirs: core/, llm/, bash/, session-persistence/, ui/, support/. Group dirs are pure containers; each package keeps its @deepseek-ai/dsh-* name. Collapse the per-package tsconfig paths maps (base + typecheck) into one @deepseek-ai/dsh-* wildcard with a candidate per group, and derive the publint list from the hierarchy. Update all depth-coupled globs/configs (workspace, tsdown, vitest, eslint, knip, tsconfig includes/refs, per-package tsconfigs, generators, doc-script scopes, type-equiv manifest) and the cross-package/script relative imports in tests. Fix doc-typecheck's workspacePaths() to parse tsconfig JSONC via the TypeScript API instead of a regex comment-strip, which corrupted the new wildcard `/*/` path candidates. WIP: doc cross-links and package/RFC docs still to update.
76 lines
2.2 KiB
TypeScript
76 lines
2.2 KiB
TypeScript
/**
|
|
* Per-agent message inbox: queued and steering FIFOs. Purely an in-memory
|
|
* mechanism of the loop driver — the public surface is `Agent.send()` and
|
|
* `Agent.steer()`.
|
|
*
|
|
* @module dsh-agent-loop/inbox
|
|
*/
|
|
|
|
import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
|
|
|
|
/** One message waiting in an agent's inbox. */
|
|
export interface InboxMessage {
|
|
content: ContentBlock[]
|
|
source: MessageSource
|
|
}
|
|
|
|
/**
|
|
* Per-agent inbox: a queued FIFO (drained at turn start) and a steering FIFO
|
|
* (drained between steps of a running turn). Purely an in-memory mechanism of
|
|
* the loop — the public surface is `Agent.send()` / `Agent.steer()`.
|
|
*/
|
|
export class Inbox {
|
|
private queuedMessages: InboxMessage[] = []
|
|
private steeringMessages: InboxMessage[] = []
|
|
private wakeup: (() => void) | undefined
|
|
|
|
/** Resolves when a queued message arrives (used by the idle loop). */
|
|
get hasQueued(): boolean {
|
|
return this.queuedMessages.length > 0
|
|
}
|
|
|
|
get hasSteering(): boolean {
|
|
return this.steeringMessages.length > 0
|
|
}
|
|
|
|
enqueue(message: InboxMessage): void {
|
|
this.queuedMessages.push(message)
|
|
this.wakeup?.()
|
|
}
|
|
|
|
steer(message: InboxMessage): void {
|
|
this.steeringMessages.push(message)
|
|
}
|
|
|
|
/** Drain all queued messages (turn start). */
|
|
drainQueued(): InboxMessage[] {
|
|
return this.queuedMessages.splice(0)
|
|
}
|
|
|
|
/** Drain all steering messages (between steps). */
|
|
drainSteering(): InboxMessage[] {
|
|
return this.steeringMessages.splice(0)
|
|
}
|
|
|
|
/**
|
|
* Discard all pending messages (queued + steering) without delivering them —
|
|
* used by `cancel()`, which drops un-started work rather than draining it into
|
|
* a turn. Unlike `drainQueued`/`drainSteering`, the messages are thrown away.
|
|
*/
|
|
clear(): void {
|
|
this.queuedMessages.length = 0
|
|
this.steeringMessages.length = 0
|
|
}
|
|
|
|
/** Wait until a queued message arrives or `cancel` resolves. */
|
|
waitForQueued(cancel: Promise<void>): Promise<void> {
|
|
if (this.hasQueued) return Promise.resolve()
|
|
const { promise, resolve } = Promise.withResolvers<void>()
|
|
this.wakeup = resolve
|
|
void cancel.then(resolve)
|
|
return promise.finally(() => {
|
|
if (this.wakeup === resolve) this.wakeup = undefined
|
|
})
|
|
}
|
|
}
|