diff options
| author | Adam Malczewski <[email protected]> | 2026-06-26 20:24:18 +0900 |
|---|---|---|
| committer | Adam Malczewski <[email protected]> | 2026-06-26 20:24:18 +0900 |
| commit | 12955cdf1d822ff395fd62d30916fbdb02d10e12 (patch) | |
| tree | c23a4b0700a08926eaf3636af7e5ca92a39fbe2f /packages/heartbeat/src | |
| parent | c5c34ed70e0f04b7b936fa7a1d88ef807472fb96 (diff) | |
| download | dispatch-12955cdf1d822ff395fd62d30916fbdb02d10e12.tar.gz dispatch-12955cdf1d822ff395fd62d30916fbdb02d10e12.zip | |
feat(heartbeat): resolve [type:name] variables in heartbeat prompts (CR-HB-1)
Diffstat (limited to 'packages/heartbeat/src')
| -rw-r--r-- | packages/heartbeat/src/config-store.test.ts | 12 | ||||
| -rw-r--r-- | packages/heartbeat/src/config-store.ts | 9 | ||||
| -rw-r--r-- | packages/heartbeat/src/extension.ts | 49 | ||||
| -rw-r--r-- | packages/heartbeat/src/heartbeat.test.ts | 87 | ||||
| -rw-r--r-- | packages/heartbeat/src/heartbeat.ts | 56 | ||||
| -rw-r--r-- | packages/heartbeat/src/index.ts | 18 | ||||
| -rw-r--r-- | packages/heartbeat/src/run-store.ts | 14 | ||||
| -rw-r--r-- | packages/heartbeat/src/scheduler.ts | 18 |
8 files changed, 204 insertions, 59 deletions
diff --git a/packages/heartbeat/src/config-store.test.ts b/packages/heartbeat/src/config-store.test.ts index 0cb6afa..17d3ee9 100644 --- a/packages/heartbeat/src/config-store.test.ts +++ b/packages/heartbeat/src/config-store.test.ts @@ -48,15 +48,15 @@ describe("applyConfigUpdate (pure)", () => { }); it("clamps intervalMinutes to a minimum of 1", () => { - expect(applyConfigUpdate(DEFAULT_HEARTBEAT_CONFIG, { intervalMinutes: 0 }).intervalMinutes).toBe( - 1, - ); + expect( + applyConfigUpdate(DEFAULT_HEARTBEAT_CONFIG, { intervalMinutes: 0 }).intervalMinutes, + ).toBe(1); expect( applyConfigUpdate(DEFAULT_HEARTBEAT_CONFIG, { intervalMinutes: -5 }).intervalMinutes, ).toBe(1); - expect(applyConfigUpdate(DEFAULT_HEARTBEAT_CONFIG, { intervalMinutes: 7 }).intervalMinutes).toBe( - 7, - ); + expect( + applyConfigUpdate(DEFAULT_HEARTBEAT_CONFIG, { intervalMinutes: 7 }).intervalMinutes, + ).toBe(7); }); it("truncates a non-integer interval to an integer", () => { diff --git a/packages/heartbeat/src/config-store.ts b/packages/heartbeat/src/config-store.ts index 74c01e2..5791601 100644 --- a/packages/heartbeat/src/config-store.ts +++ b/packages/heartbeat/src/config-store.ts @@ -63,9 +63,7 @@ export interface HeartbeatConfigStore { readonly listWorkspaceIds: () => Promise<readonly string[]>; } -export function createHeartbeatConfigStore( - storage: StorageNamespace, -): HeartbeatConfigStore { +export function createHeartbeatConfigStore(storage: StorageNamespace): HeartbeatConfigStore { return { async get(workspaceId: string): Promise<HeartbeatConfig> { const raw = await storage.get(configKey(workspaceId)); @@ -88,10 +86,7 @@ export function createHeartbeatConfigStore( } }, - async update( - workspaceId: string, - update: UpdateHeartbeatRequest, - ): Promise<HeartbeatConfig> { + async update(workspaceId: string, update: UpdateHeartbeatRequest): Promise<HeartbeatConfig> { const current = await this.get(workspaceId); const next = applyConfigUpdate(current, update); await storage.set(configKey(workspaceId), JSON.stringify(next)); diff --git a/packages/heartbeat/src/extension.ts b/packages/heartbeat/src/extension.ts index 66580e6..e79b01f 100644 --- a/packages/heartbeat/src/extension.ts +++ b/packages/heartbeat/src/extension.ts @@ -3,14 +3,21 @@ * * Wires the heartbeat service against the session-orchestrator + a storage * namespace, registers the typed service handle, and arms every enabled - * workspace's scheduler on boot. + * workspace's scheduler on boot. Prompt templates (`systemPrompt` / + * `taskPrompt`) are resolved against the SAME variable catalog the global + * system-prompt template uses (via the system-prompt service's `resolveText`), + * so `[type:name]` placeholders reach the model substituted — not raw. */ +import type { ConversationStore } from "@dispatch/conversation-store"; +import { conversationStoreHandle } from "@dispatch/conversation-store"; import type { Extension, HostAPI, Manifest } from "@dispatch/kernel"; import { type SessionOrchestrator, sessionOrchestratorHandle, } from "@dispatch/session-orchestrator"; +import type { SystemPromptService } from "@dispatch/system-prompt"; +import { systemPromptHandle } from "@dispatch/system-prompt"; import { createHeartbeatService, heartbeatServiceHandle } from "./heartbeat.js"; export const manifest: Manifest = { @@ -19,7 +26,12 @@ export const manifest: Manifest = { version: "0.0.0", apiVersion: "^0.1.0", trust: "bundled", - dependsOn: ["session-orchestrator"], + // system-prompt provides `resolveText` (the variable resolver used to + // substitute [type:name] placeholders in heartbeat prompts); conversation- + // store resolves the workspace's default cwd (the resolver runs git / reads + // files against it, mirroring the global template). Both lookups are lazy + // (at fire time, not activation), but declaring them keeps the DAG honest. + dependsOn: ["session-orchestrator", "system-prompt", "conversation-store"], activation: "eager", contributes: { services: ["heartbeat"] }, }; @@ -34,7 +46,38 @@ export const extension: Extension = { const storage = host.storage("heartbeat"); const logger = host.logger; - const service = createHeartbeatService({ storage, orchestrator, logger }); + // Resolve [type:name] placeholders in heartbeat prompts via the + // system-prompt service (same resolver + variables as the global + // template). The cwd is the workspace's defaultCwd (resolved the same + // way the orchestrator resolves a new conversation's effective cwd); + // falling back to process.cwd() when the workspace has none. Both + // services are declared `dependsOn` (always activated before heartbeat). + const systemPromptService = host.getService<SystemPromptService>(systemPromptHandle); + const conversationStore = host.getService<ConversationStore>(conversationStoreHandle); + + const resolvePrompt = async ( + template: string, + ctx: { + readonly workspaceId: string; + readonly conversationId: string; + readonly model: string; + }, + ): Promise<string> => { + const workspace = await conversationStore.getWorkspace(ctx.workspaceId); + const cwd = workspace?.defaultCwd ?? process.cwd(); + return systemPromptService.resolveText(template, cwd, { + ...(ctx.conversationId !== "" ? { conversationId: ctx.conversationId } : {}), + ...(ctx.model !== "" ? { model: ctx.model } : {}), + workspaceId: ctx.workspaceId, + }); + }; + + const service = createHeartbeatService({ + storage, + orchestrator, + logger, + resolvePrompt, + }); // Reconcile stale runs + arm enabled workspaces on boot. await service.startAll(); diff --git a/packages/heartbeat/src/heartbeat.test.ts b/packages/heartbeat/src/heartbeat.test.ts index c185c92..4c01904 100644 --- a/packages/heartbeat/src/heartbeat.test.ts +++ b/packages/heartbeat/src/heartbeat.test.ts @@ -145,6 +145,14 @@ const flush = async (): Promise<void> => { function createService(opts: { readonly orch: ReturnType<typeof createFakeOrchestrator>; readonly storage?: StorageNamespace; + readonly resolvePrompt?: ( + template: string, + ctx: { + readonly workspaceId: string; + readonly conversationId: string; + readonly model: string; + }, + ) => Promise<string>; }) { const fake = createFakeTimers(); let id = 0; @@ -154,6 +162,7 @@ function createService(opts: { orchestrator: opts.orch, timers: fake.timers, generateId: () => `id-${++id}`, + ...(opts.resolvePrompt !== undefined ? { resolvePrompt: opts.resolvePrompt } : {}), }); return { svc, advance: fake.advance, storage }; } @@ -225,14 +234,75 @@ describe("createHeartbeatService", () => { await flush(); }); + it("resolves [type:name] variables in both systemPrompt and taskPrompt before sending", async () => { + const orch = createFakeOrchestrator(); + // A fake resolver that mirrors the real resolver's contract: substitute + // known [type:name] placeholders, leave unknown text verbatim. + const resolvePrompt = async ( + template: string, + ctx: { + readonly workspaceId: string; + readonly conversationId: string; + readonly model: string; + }, + ): Promise<string> => { + return template + .replaceAll("[system:os]", "Linux (WSL)") + .replaceAll("[prompt:cwd]", "/repo") + .replaceAll("[prompt:workspace_id]", ctx.workspaceId) + .replaceAll("[prompt:conversation_id]", ctx.conversationId) + .replaceAll("[prompt:model]", ctx.model); + }; + const { svc, advance } = createService({ orch, resolvePrompt }); + await svc.updateConfig("ws-1", { + enabled: true, + systemPrompt: "You run on [system:os] in [prompt:cwd] (ws [prompt:workspace_id]).", + taskPrompt: "Check chats for [prompt:conversation_id] on [system:os].", + model: "opencode/gpt-4o", + intervalMinutes: 1, + }); + + advance(60_000); + await flush(); + expect(orch.pending).toHaveLength(1); + const turn = orch.pending[0]!; + // Variables substituted — NOT left as literal [type:name] text. + expect(turn.systemPrompt).toBe("You run on Linux (WSL) in /repo (ws ws-1)."); + expect(turn.text).toBe(`Check chats for ${turn.conversationId} on Linux (WSL).`); + expect(turn.modelName).toBe("opencode/gpt-4o"); + turn.resolve(); + await flush(); + }); + + it("passes prompts through raw when no resolver is wired (resolution is optional)", async () => { + const orch = createFakeOrchestrator(); + // No resolvePrompt → raw pass-through (the default). + const { svc, advance } = createService({ orch }); + await svc.updateConfig("ws-1", { + enabled: true, + systemPrompt: "raw [system:os] prompt", + taskPrompt: "raw [system:date] task", + intervalMinutes: 1, + }); + + advance(60_000); + await flush(); + const turn = orch.pending[0]!; + // Unresolved — literals reach the orchestrator verbatim. + expect(turn.systemPrompt).toBe("raw [system:os] prompt"); + expect(turn.text).toBe("raw [system:date] task"); + turn.resolve(); + await flush(); + }); + it("stopRun aborts the turn and marks the run stopped (not overwritten on completion)", async () => { const orch = createFakeOrchestrator(); const { svc, advance } = createService({ orch }); await svc.updateConfig("ws-1", { enabled: true, taskPrompt: "go", intervalMinutes: 1 }); advance(60_000); await flush(); - const conversationId = orch.pending[0]!.conversationId; - const runId = (await svc.listRuns("ws-1"))[0]!.id; + const conversationId = orch.pending[0]?.conversationId; + const runId = (await svc.listRuns("ws-1"))[0]?.id; const res = await svc.stopRun("ws-1", runId); expect(res).toEqual({ ok: true }); @@ -248,10 +318,10 @@ describe("createHeartbeatService", () => { await svc.updateConfig("ws-1", { enabled: true, taskPrompt: "go", intervalMinutes: 1 }); advance(60_000); await flush(); - orch.pending[0]!.resolve(); + orch.pending[0]?.resolve(); await flush(); - const runId = (await svc.listRuns("ws-1"))[0]!.id; + const runId = (await svc.listRuns("ws-1"))[0]?.id; const res = await svc.stopRun("ws-1", runId); expect(res).toEqual({ ok: true }); expect(orch.stopped).toEqual([]); // no abort on a completed run @@ -284,7 +354,14 @@ describe("createHeartbeatService", () => { ); await storage.set( "config:ws-1", - JSON.stringify({ enabled: false, systemPrompt: "", taskPrompt: "", intervalMinutes: 30, model: "", reasoningEffort: null }), + JSON.stringify({ + enabled: false, + systemPrompt: "", + taskPrompt: "", + intervalMinutes: 30, + model: "", + reasoningEffort: null, + }), ); const { svc } = createService({ orch: createFakeOrchestrator(), storage }); diff --git a/packages/heartbeat/src/heartbeat.ts b/packages/heartbeat/src/heartbeat.ts index d1efe47..7311479 100644 --- a/packages/heartbeat/src/heartbeat.ts +++ b/packages/heartbeat/src/heartbeat.ts @@ -7,12 +7,9 @@ import type { StopHeartbeatRunResponse, UpdateHeartbeatRequest, } from "@dispatch/transport-contract"; -import { - type HeartbeatConfigStore, - createHeartbeatConfigStore, -} from "./config-store.js"; -import { type HeartbeatRunStore, createHeartbeatRunStore } from "./run-store.js"; -import { HeartbeatScheduler, type Timers, realTimers } from "./scheduler.js"; +import { createHeartbeatConfigStore, type HeartbeatConfigStore } from "./config-store.js"; +import { createHeartbeatRunStore, type HeartbeatRunStore } from "./run-store.js"; +import { HeartbeatScheduler, realTimers, type Timers } from "./scheduler.js"; /** * The heartbeat service surface — what transport-http consumes and what the @@ -36,10 +33,7 @@ export interface HeartbeatService { * Stop an in-flight run (abort its turn). Idempotent for an already-finished * run. Throws when the run id is unknown (→ HTTP 404). */ - readonly stopRun: ( - workspaceId: string, - runId: string, - ) => Promise<StopHeartbeatRunResponse>; + readonly stopRun: (workspaceId: string, runId: string) => Promise<StopHeartbeatRunResponse>; /** Boot: sweep stale runs + arm every enabled workspace's scheduler. */ readonly startAll: () => Promise<void>; /** Shutdown: stop every scheduler. */ @@ -60,6 +54,23 @@ export interface HeartbeatServiceDeps { readonly timers?: Timers; /** Injectable id generator (default: crypto.randomUUID). */ readonly generateId?: () => string; + /** + * Resolve `[type:name]` variable placeholders in a prompt template against + * the current environment — the SAME resolver + variable catalog the global + * system-prompt template uses (system/file/prompt/git groups). Applied once + * per run, when the turn is constructed (mirrors the global template's + * construct-once-per-conversation resolution). When omitted, templates pass + * through UNRESOLVED (raw) — the extension wires the real resolver; tests + * inject a fake. + */ + readonly resolvePrompt?: ( + template: string, + ctx: { + readonly workspaceId: string; + readonly conversationId: string; + readonly model: string; + }, + ) => Promise<string>; } interface ActiveRun { @@ -75,6 +86,10 @@ export function createHeartbeatService(deps: HeartbeatServiceDeps): HeartbeatSer const configStore: HeartbeatConfigStore = createHeartbeatConfigStore(deps.storage); const runStore: HeartbeatRunStore = createHeartbeatRunStore(deps.storage); const orchestrator = deps.orchestrator; + // Default: pass templates through UNRESOLVED (raw). The extension wires the + // real resolver so [type:name] placeholders are substituted like the global + // system-prompt template; tests inject a fake. + const resolvePrompt = deps.resolvePrompt ?? ((template: string) => Promise.resolve(template)); // runId → active-run tracking (in-memory; the durable record lives in the // run store). Used to (a) map a stop request to its conversation, and @@ -106,17 +121,32 @@ export function createHeartbeatService(deps: HeartbeatServiceDeps): HeartbeatSer logger?.info("heartbeat: run started", { workspaceId, runId, conversationId }); + // Resolve [type:name] variables in the system/task prompts ONCE, when + // the turn is constructed — the same resolver + variable catalog the + // global system-prompt template uses. The orchestrator sends an + // explicit systemPrompt override AS-IS (bypassing its own templated + // prompt), so resolution must happen HERE, before handleMessage. For + // prompts using stable variables (os/cwd/git) this yields a stable, + // cache-warm prompt across runs; time-bearing variables refresh per + // run (mirroring the global template's per-conversation resolution). + const resolveCtx = { workspaceId, conversationId, model: config.model }; + const [systemPrompt, taskPrompt] = await Promise.all([ + resolvePrompt(config.systemPrompt, resolveCtx), + resolvePrompt(config.taskPrompt, resolveCtx), + ]); + try { await orchestrator.handleMessage({ conversationId, - text: config.taskPrompt, + text: taskPrompt, // Fire-and-forget: the heartbeat loop does not consume the // streamed events (it only awaits turn completion to mark the // run done). A no-op onEvent satisfies the required callback. onEvent: () => {}, // Always an explicit override (incl. the empty string = no system - // prompt) — bypasses the templated workspace prompt. - systemPrompt: config.systemPrompt, + // prompt) — bypasses the templated workspace prompt. Resolved + // above so [type:name] placeholders are substituted, not raw. + systemPrompt, ...(config.model !== "" ? { modelName: config.model } : {}), ...(config.reasoningEffort !== null ? { reasoningEffort: config.reasoningEffort } : {}), workspaceId, diff --git a/packages/heartbeat/src/index.ts b/packages/heartbeat/src/index.ts index 3886be5..985ad68 100644 --- a/packages/heartbeat/src/index.ts +++ b/packages/heartbeat/src/index.ts @@ -1,16 +1,16 @@ +export { + applyConfigUpdate, + createHeartbeatConfigStore, + DEFAULT_HEARTBEAT_CONFIG, + type HeartbeatConfigStore, + MIN_INTERVAL_MINUTES, +} from "./config-store.js"; export { extension, manifest } from "./extension.js"; export { - heartbeatServiceHandle, createHeartbeatService, type HeartbeatService, type HeartbeatServiceDeps, + heartbeatServiceHandle, } from "./heartbeat.js"; -export { - createHeartbeatConfigStore, - applyConfigUpdate, - DEFAULT_HEARTBEAT_CONFIG, - MIN_INTERVAL_MINUTES, - type HeartbeatConfigStore, -} from "./config-store.js"; export { createHeartbeatRunStore, type HeartbeatRunStore } from "./run-store.js"; -export { HeartbeatScheduler, realTimers, type Timers, type TimerHandle } from "./scheduler.js"; +export { HeartbeatScheduler, realTimers, type TimerHandle, type Timers } from "./scheduler.js"; diff --git a/packages/heartbeat/src/run-store.ts b/packages/heartbeat/src/run-store.ts index 7522aab..819ae00 100644 --- a/packages/heartbeat/src/run-store.ts +++ b/packages/heartbeat/src/run-store.ts @@ -19,10 +19,7 @@ function parseRunId(key: string, workspaceId: string): string { export interface HeartbeatRunStore { /** Create a new run record (status `"running"`) and persist it. */ - readonly create: ( - workspaceId: string, - run: HeartbeatRun, - ) => Promise<HeartbeatRun>; + readonly create: (workspaceId: string, run: HeartbeatRun) => Promise<HeartbeatRun>; /** Update the status of an existing run. No-op if the run is unknown. */ readonly setStatus: ( workspaceId: string, @@ -35,13 +32,8 @@ export interface HeartbeatRunStore { readonly list: (workspaceId: string) => Promise<readonly HeartbeatRun[]>; } -export function createHeartbeatRunStore( - storage: StorageNamespace, -): HeartbeatRunStore { - async function readRun( - workspaceId: string, - runId: string, - ): Promise<HeartbeatRun | null> { +export function createHeartbeatRunStore(storage: StorageNamespace): HeartbeatRunStore { + async function readRun(workspaceId: string, runId: string): Promise<HeartbeatRun | null> { const raw = await storage.get(runKey(workspaceId, runId)); if (raw === null) return null; try { diff --git a/packages/heartbeat/src/scheduler.ts b/packages/heartbeat/src/scheduler.ts index 53ffb96..e366bbc 100644 --- a/packages/heartbeat/src/scheduler.ts +++ b/packages/heartbeat/src/scheduler.ts @@ -72,17 +72,25 @@ export class HeartbeatScheduler { * run is in progress, the new `intervalMinutes` takes effect on the next * re-arm (the in-progress run is never cancelled here). */ - arm(workspaceId: string, config: { - readonly enabled: boolean; - readonly intervalMinutes: number; - }): void { + arm( + workspaceId: string, + config: { + readonly enabled: boolean; + readonly intervalMinutes: number; + }, + ): void { if (!config.enabled) { this.disarm(workspaceId); return; } let schedule = this.schedules.get(workspaceId); if (schedule === undefined) { - schedule = { intervalMinutes: config.intervalMinutes, timer: undefined, running: false, armed: true }; + schedule = { + intervalMinutes: config.intervalMinutes, + timer: undefined, + running: false, + armed: true, + }; this.schedules.set(workspaceId, schedule); } else { schedule.intervalMinutes = config.intervalMinutes; |
