6 Commits

Author SHA1 Message Date
Croissant Le Doux
6bee1a12aa test: point the live milestone assertion at P5 after the backlog reconcile
Closed the 15 genuinely-shipped issues on the dogfood repo (scheduler, Monte Carlo,
calibration, lifecycle, model router, query_project, capture_work, record_directive,
etc.) so the board reflects reality: 10 open / 24 done. P2 dropped off Runway once
all its work shipped (a milestone with no open scope isn't runway — correct), so the
live spec now drills into P5 (which still has open scope) instead of P2.

No code change — this only updates the live assertion to match the reconciled backlog.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 01:23:56 -04:00
Croissant Le Doux
1636d6bada feat: capacity-aware scheduling (#8) — real focus factors drive every forecast
Turns the single-serial-worker scheduler into a capacity-aware, multi-lane one.
Configured team members become lanes; an issue runs on its assignee's lane (or the
earliest-free lane), its duration scaled by that lane's throughput
(focusFactor × allocation). Every forecast — Focus cone, Runway, milestone
drill-in — is now capacity-aware.

core (@commitea/core):
- capacity/capacity-v0: CapacityMember + capacityPerWorkday + parseCapacityConfig
  (clamps, drops invalid; degrades to []).
- scheduler/scheduler-capacity-v0: scheduleWithCapacity reuses the v0 topo order +
  critical path, re-lays work across lanes (layoutOnLanes, resolveLanes, makespan).
  Empty workers → the single serial plan verbatim.
- forecast() gains options.workers: each MC trial lays sampled durations across the
  lanes and takes the makespan; serial path unchanged. SchedulableIssue gains
  assignee; ScheduledItem gains worker.
- 11 new tests (parse/clamp, parallelism halves makespan, speed scaling, assignee
  routing, cross-lane deps, forecast makespan shrinks with lanes).

app:
- pm-state capacity/members.json read (readCapacity + pmstate:capacity bridge);
  useCapacity hook → workers; forecastBacklog/runwayView/milestoneView pass workers.
- Runway Capacity card shows the real config (person · focus · alloc · pd/day).

Config lives in pm-state (D4); seeded christian(0.8)/stephen(0.6×0.5). Degrades to
the fixture/serial when absent.

Verified: 128 core tests green, desktop typecheck clean, 14 fixture e2e green. Live:
the capacity card is real, and the P2 forecast shifts 32d→37d — honest, since real
focus factors (<1) replace the v0 focus-1.0 assumption.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 00:59:03 -04:00
d80e1266ee Merge pull request 'feat: stream Reginald's replies token-by-token' (#49) from feat/streaming-chat into main
Reviewed-on: #49
2026-07-09 04:43:58 +00:00
cf03827cd5 Merge branch 'main' into feat/streaming-chat 2026-07-09 04:43:51 +00:00
69106b603d Merge pull request 'feat: Runway complete — per-milestone forecasts + real milestone drill-in' (#48) from feat/runway-real into main
Reviewed-on: #48
2026-07-09 04:43:46 +00:00
Croissant Le Doux
dbcdcda5e7 feat: stream Reginald's replies token-by-token
The 26b is slow (~30s/call); the chat now shows the answer forming instead of
freezing until it's done. The final prose streams over SSE; tool-calling turns
stay structured (no partial tokens), so streaming kicks in for the narration.

core (@commitea/core):
- chat-client.complete gains an optional onToken — when set, it requests
  stream:true and parses the OpenAI SSE stream, emitting content deltas and
  assembling streamed tool-call argument fragments into the final result.
- GiteaHttpResponse exposes the optional `body` stream (real fetch has it; stubs
  don't). agent-loop threads onToken to each completion.

app:
- model:chat forwards each delta to the renderer (event.sender.send); preload
  exposes model.onToken(cb) → unsubscribe. useChat accumulates the live stream
  into a growing bubble (with a cursor), replaced by the authoritative final
  content when the turn resolves. Unconfigured → scripted reply, unchanged.

Verified: 118 core tests green (2 streaming: SSE content deltas + tool-call
fragment assembly), desktop typecheck clean, 14 fixture e2e green. Live: a real
turn against gemma-4-26b assembles the correct answer via the streaming path
(live-reginald green) — the reply now renders token-by-token.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 00:22:12 -04:00
22 changed files with 573 additions and 32 deletions

View File

@@ -42,13 +42,17 @@ test.describe('live backlog', () => {
win.getByText(/cold-start priors · \d+\/20 closed issues estimated|calibrated on \d+ closed/),
).toBeVisible()
// Real per-milestone forecasts — these milestone names come from gitea, not the
// fixture (which lists Beta / Pilot-ready / v1.0).
await expect(win.getByText(/P2 — Scheduler/)).toBeVisible()
// fixture (which lists Beta / Pilot-ready / v1.0). Only milestones with open
// scope appear (P2 dropped off once all its work shipped — correct).
await expect(win.getByText(/P5 — Dogfood/)).toBeVisible()
// Real capacity config from pm-state (christian/stephen), not the fixture (Stephen/Ana K.)
await expect(win.getByText('christian', { exact: true })).toBeVisible()
await expect(win.getByText(/pd\/day/).first()).toBeVisible()
await win.screenshot({ path: join(here, '.artifacts', 'screens', 'live-runway.png'), fullPage: true, animations: 'disabled' })
// Milestone drill-in — clicking a real milestone opens its real detail
await win.getByText(/P2Scheduler/).click()
await expect(win.getByRole('heading', { name: 'P2Scheduler + Monte Carlo' })).toBeVisible()
await win.getByText(/P5Dogfood/).click()
await expect(win.getByRole('heading', { name: 'P5Dogfood + polish' })).toBeVisible()
await expect(win.getByText(/\d+ issues · est \d+d/)).toBeVisible()
await win.screenshot({ path: join(here, '.artifacts', 'screens', 'live-milestone.png'), fullPage: true, animations: 'disabled' })
await rail.getByRole('button', { name: 'Runway' }).click()

View File

@@ -20,6 +20,7 @@ import {
type IssueChange,
type LifecycleEvent,
makeDirectiveEntry,
parseCapacityConfig,
parseDirectiveLog,
planIssueChange,
type ProjectSnapshot,
@@ -109,6 +110,19 @@ export async function readDirectives(client: GiteaClient) {
return parseDirectiveLog(text)
}
const CAPACITY_PATH = 'capacity/members.json'
/** Read the capacity config from the pm-state repo (empty when absent). */
export async function readCapacity(client: GiteaClient) {
const file = await client.getFile(CAPACITY_PATH)
if (!file) return []
try {
return parseCapacityConfig(JSON.parse(Buffer.from(file.contentBase64, 'base64').toString('utf8')))
} catch {
return []
}
}
/** Full reconcile: issues + milestones + native deps + lifecycle timelines. */
export async function reconcileSnapshot(
client: GiteaClient,
@@ -266,4 +280,15 @@ export function registerGiteaIpc(): void {
return { ok: false as const, reason: 'error' as const, message: e instanceof Error ? e.message : String(e) }
}
})
// Read the capacity config from the pm-state repo (for capacity-aware forecasts).
ipcMain.handle('pmstate:capacity', async () => {
const pm = getPmStateClient()
if (!pm) return { ok: false as const, reason: 'unconfigured' as const, members: [] }
try {
return { ok: true as const, members: await readCapacity(pm) }
} catch {
return { ok: true as const, members: [] }
}
})
}

View File

@@ -86,11 +86,15 @@ export function registerModelIpc(): void {
return { configured: true, model }
})
ipcMain.handle('model:chat', async (_event, messages: ChatMessage[]) => {
ipcMain.handle('model:chat', async (event, messages: ChatMessage[]) => {
if (!router) return { ok: false as const, reason: 'unconfigured' as const }
const client = getGiteaClient()
const model = await resolveLoadedModel(router.small.baseUrl, router.small.model)
const chat = createChatClient({ ...router.small, model }, fetch)
// stream the model's prose to the renderer token-by-token
const onToken = (delta: string) => {
if (!event.sender.isDestroyed()) event.sender.send('model:chat:token', delta)
}
// Proposals the model formulates this turn; the renderer approves them (the
// write happens through gitea:applyChange, never inside the loop).
@@ -129,10 +133,11 @@ export function registerModelIpc(): void {
try {
const turn = await runAgentTurn({
complete: (m, t) => chat.complete(m, t),
complete: (m, t, ot) => chat.complete(m, t, ot),
messages: [{ role: 'system', content: REGINALD_SYSTEM }, ...messages],
tools: REGINALD_TOOLS,
execute,
onToken,
})
return { ok: true as const, content: turn.content, steps: turn.steps, proposals }
} catch (e) {

View File

@@ -19,12 +19,20 @@ const api = {
pmstate: {
/** Read the directive ledger from the pm-state repo. */
directives: () => ipcRenderer.invoke('pmstate:directives'),
/** Read the capacity config from the pm-state repo. */
capacity: () => ipcRenderer.invoke('pmstate:capacity'),
},
model: {
/** Whether a model endpoint is configured (else the UI keeps the scripted Reginald). */
status: () => ipcRenderer.invoke('model:status'),
/** One agent turn: messages in, Reginald's prose + the tools it consulted out. */
chat: (messages: unknown) => ipcRenderer.invoke('model:chat', messages),
/** Subscribe to streamed prose tokens for the in-flight turn; returns an unsubscribe. */
onToken: (cb: (delta: string) => void) => {
const listener = (_e: unknown, delta: string) => cb(delta)
ipcRenderer.on('model:chat:token', listener)
return () => ipcRenderer.removeListener('model:chat:token', listener)
},
/** Decompose a braindump into a proposed issue set (capture_work). */
capture: (braindump: string) => ipcRenderer.invoke('model:capture', braindump),
},

View File

@@ -1,23 +1,35 @@
import React from 'react'
import { type CapacityMember, capacityPerWorkday } from '@commitea/core'
import { RunwayBar } from '../charts/chart.js'
import { Card, Badge, Tag, Icon, IconButton } from '../ui/index.js'
import { RUNWAY, CAPACITY, type RunwayMilestone } from '../../data/fixtures.js'
// Runway — capacity vs milestone dates; ranges, never points.
// `milestones` (real per-milestone forecasts) overrides the demo when present.
// `milestones` (real per-milestone forecasts) + `capacity` (real config) override the demo.
export function RunwayScreen({
onOpenCalibration,
onOpenMilestone,
calibration,
milestones,
capacity,
}: {
onOpenCalibration: () => void
onOpenMilestone: (id?: number) => void
calibration?: { n: number; coldStart: boolean }
milestones?: RunwayMilestone[]
capacity?: CapacityMember[]
}) {
const rows = milestones && milestones.length ? milestones : RUNWAY
const capacityRows =
capacity && capacity.length
? capacity.map((m) => ({
who: m.person,
slices: `focus ${m.focusFactor} · alloc ${Math.round(m.allocation * 100)}%`,
hours: `${capacityPerWorkday(m).toFixed(2)} pd/day`,
}))
: CAPACITY
const calibNote = calibration
? calibration.coldStart
? `cold-start priors · ${calibration.n}/20 closed issues estimated`
@@ -59,7 +71,7 @@ export function RunwayScreen({
<div style={{ display: 'grid', gridTemplateColumns: '1fr 1fr', gap: 14, alignItems: 'start' }}>
<Card overline="Capacity" flush>
<div>
{CAPACITY.map((p, i) => (
{capacityRows.map((p, i) => (
<div key={p.who} style={{
display: 'flex', alignItems: 'center', gap: 12, padding: '12px 20px',
borderTop: i === 0 ? 'none' : '1px solid var(--line-1)',

View File

@@ -6,6 +6,7 @@ import type { IssueChange } from '@commitea/core'
import type { IssueRef } from '../../data/fixtures.js'
import {
backlogCalibration,
capacityWorkers,
forecastBacklog,
issuesToBoardColumns,
milestoneView,
@@ -13,6 +14,7 @@ import {
scheduleFocus,
} from '../../lib/backlog.js'
import { useBacklog } from '../../lib/use-backlog.js'
import { useCapacity } from '../../lib/use-capacity.js'
import { PrimitivesGallery } from '../gallery.js'
import { BoardScreen } from '../screens/board-screen.js'
import { CalibrationScreen } from '../screens/calibration-screen.js'
@@ -96,6 +98,8 @@ export function AppShell() {
const [readIds, setReadIds] = useState<number[]>([])
const [milestoneId, setMilestoneId] = useState<number | null>(null)
const [backlog, refetchBacklog] = useBacklog()
const capacityMembers = useCapacity()
const workers = capacityWorkers(capacityMembers)
const boardColumns =
backlog.status === 'ready' ? issuesToBoardColumns(backlog.issues, backlog.timelines) : undefined
const focus =
@@ -104,10 +108,12 @@ export function AppShell() {
backlog.status === 'ready' ? backlogCalibration(backlog.issues, backlog.timelines) : undefined
const forecast =
backlog.status === 'ready'
? (forecastBacklog(backlog.issues, backlog.deps, new Date(), calibration?.model) ?? undefined)
? (forecastBacklog(backlog.issues, backlog.deps, new Date(), calibration?.model, workers) ?? undefined)
: undefined
const runwayMilestones =
backlog.status === 'ready' ? runwayView(backlog.issues, backlog.milestones, backlog.deps) : undefined
backlog.status === 'ready'
? runwayView(backlog.issues, backlog.milestones, backlog.deps, new Date(), workers)
: undefined
const milestone =
backlog.status === 'ready' && milestoneId != null
? (milestoneView(
@@ -117,6 +123,8 @@ export function AppShell() {
backlog.deps,
backlog.timelines,
calibration?.model,
new Date(),
workers,
) ?? undefined)
: undefined
@@ -221,6 +229,7 @@ export function AppShell() {
}}
calibration={calibration ? { n: calibration.model.n, coldStart: calibration.model.coldStart } : undefined}
milestones={runwayMilestones}
capacity={capacityMembers}
/>
)
case 'calibration':

View File

@@ -19,7 +19,7 @@ export interface ChatPanelProps {
}
export function ChatPanel({ onOpenDirectives, offline, onApplyChange }: ChatPanelProps) {
const { msgs, thinking, live, model, steps, proposals, send: sendChat, approve, dismiss } = useChat(onApplyChange)
const { msgs, thinking, live, model, steps, proposals, streaming, send: sendChat, approve, dismiss } = useChat(onApplyChange)
// shorten "google/gemma-4-26b-a4b-qat" → "gemma-4-26b" for the header chip
const modelLabel = model ? (model.split('/').pop() ?? model).replace(/-(qat|instruct|it|gguf)$/i, '') : 'gemma-4'
const [text, setText] = useState('')
@@ -28,7 +28,7 @@ export function ChatPanel({ onOpenDirectives, offline, onApplyChange }: ChatPane
useEffect(() => {
const el = scrollRef.current
if (el) el.scrollTop = el.scrollHeight
}, [msgs, thinking])
}, [msgs, thinking, streaming])
const send = () => {
const t = text.trim()
@@ -98,7 +98,14 @@ export function ChatPanel({ onOpenDirectives, offline, onApplyChange }: ChatPane
The model is away from its desk. Reads still work; writes will wait their turn.
</div>
) : null}
{thinking ? <div style={{ font: 'var(--text-agent)', color: 'var(--ink-3)' }}>considering</div> : null}
{streaming ? (
<div style={{ font: 'var(--text-agent)', color: 'var(--ink-1)', lineHeight: 1.55 }}>
{streaming}
<span style={{ opacity: 0.5 }}></span>
</div>
) : thinking ? (
<div style={{ font: 'var(--text-agent)', color: 'var(--ink-3)' }}>considering</div>
) : null}
{!thinking && steps.length ? (
<div style={{ font: 'var(--text-caption)', color: 'var(--ink-3)', display: 'flex', alignItems: 'center', gap: 5 }}>
<Icon name="eye" size={11} /> consulted {Array.from(new Set(steps.map((s) => s.replace('query_project', 'the project').replace('propose_change', 'the labels').replace('record_directive', 'the directive ledger')))).join(', ')}

View File

@@ -1,5 +1,6 @@
import type {
AgentStep,
CapacityMember,
ChangeProposal,
ChatMessage,
CaptureProposal,
@@ -68,6 +69,8 @@ export interface ModelBridge {
status(): Promise<{ configured: boolean; model: string | null }>
chat(messages: ChatMessage[]): Promise<ChatResult>
capture(braindump: string): Promise<CaptureResult>
/** Subscribe to streamed prose tokens; returns an unsubscribe fn. */
onToken(cb: (delta: string) => void): () => void
}
/** The result of reading the directive ledger. */
@@ -75,9 +78,15 @@ export type DirectivesResult =
| { ok: false; reason: 'unconfigured' | 'error'; message?: string }
| { ok: true; directives: DirectiveRecord[] }
/** The result of reading the capacity config. */
export type CapacityResult =
| { ok: false; reason: 'unconfigured'; members: CapacityMember[] }
| { ok: true; members: CapacityMember[] }
/** The pm-state bridge (machine-derived state) exposed by the preload over IPC. */
export interface PmStateBridge {
directives(): Promise<DirectivesResult>
capacity(): Promise<CapacityResult>
}
declare global {

View File

@@ -2,10 +2,13 @@ import {
type CalibrationModel,
type CalibrationSample,
calibrationSamples,
type CapacityMember,
capacityPerWorkday,
COLD_START_THRESHOLD,
type DependencyEdge,
fitCalibration,
forecast,
type Worker,
type GiteaIssue,
type GiteaMilestone,
inferLifecycle,
@@ -107,9 +110,15 @@ function toSchedulable(issues: GiteaIssue[]) {
labels: i.labels,
estimateDays: i.facts.estimateDays,
priority: i.facts.priority,
assignee: i.assignee,
}))
}
/** gitea CapacityMembers → scheduler lanes (person + throughput). */
export function capacityWorkers(members: CapacityMember[]): Worker[] {
return members.map((m) => ({ person: m.person, speed: capacityPerWorkday(m) }))
}
export interface ForecastView {
scope: number
cone: BurnUpData
@@ -132,9 +141,10 @@ export function forecastBacklog(
deps: DependencyEdge[],
today: Date = new Date(),
calibration?: CalibrationModel,
workers: Worker[] = [],
): ForecastView | null {
const model = calibration ? toDurationModel(calibration) : undefined
const f = forecast(toSchedulable(issues), deps, model ? { model } : {})
const f = forecast(toSchedulable(issues), deps, { ...(model ? { model } : {}), workers })
const cone = buildBurnUpData(f, today)
if (!cone) return null
return {
@@ -252,6 +262,7 @@ export function runwayView(
milestones: GiteaMilestone[],
deps: DependencyEdge[],
today: Date = new Date(),
workers: Worker[] = [],
): RunwayMilestone[] {
const open = issues.filter((i) => i.state === 'open')
const rows = milestones
@@ -266,8 +277,10 @@ export function runwayView(
labels: i.labels,
estimateDays: i.facts.estimateDays,
priority: i.facts.priority,
assignee: i.assignee,
})),
deps,
{ workers },
)
const p90 = f.curve.length ? f.curve[f.curve.length - 1].p90Day : f.p95Day
const dueDay = m.dueOn ? workingDaysBetween(today, new Date(m.dueOn)) : null
@@ -324,6 +337,7 @@ export function milestoneView(
timelines: Timelines = {},
calibration?: CalibrationModel,
today: Date = new Date(),
workers: Worker[] = [],
): MilestoneView | null {
const m = milestones.find((x) => x.id === id)
if (!m) return null
@@ -340,9 +354,10 @@ export function milestoneView(
labels: i.labels,
estimateDays: i.facts.estimateDays,
priority: i.facts.priority,
assignee: i.assignee,
})),
deps,
model ? { model } : {},
{ ...(model ? { model } : {}), workers },
)
const cone = buildBurnUpData(f, today)

View File

@@ -0,0 +1,24 @@
import { useEffect, useState } from 'react'
import type { CapacityMember } from '@commitea/core'
/**
* The team's capacity config from the pm-state repo. Empty when unconfigured —
* forecasts then fall back to a single serial worker. Read once on mount.
*/
export function useCapacity(): CapacityMember[] {
const [members, setMembers] = useState<CapacityMember[]>([])
useEffect(() => {
let alive = true
window.commitea.pmstate
.capacity()
.then((r) => {
if (alive) setMembers(r.members)
})
.catch(() => {})
return () => {
alive = false
}
}, [])
return members
}

View File

@@ -20,6 +20,8 @@ export interface ChatState {
steps: string[]
/** Changes Reginald has proposed and is awaiting approval on. */
proposals: ChangeProposal[]
/** The in-flight streamed prose (grows token-by-token) before the turn finalizes. */
streaming: string
send: (text: string) => void
approve: (p: ChangeProposal) => void
dismiss: (p: ChangeProposal) => void
@@ -41,6 +43,7 @@ export function useChat(onApplyChange?: (change: IssueChange) => Promise<{ ok: b
const [model, setModel] = useState<string | null>(null)
const [steps, setSteps] = useState<string[]>([])
const [proposals, setProposals] = useState<ChangeProposal[]>([])
const [streaming, setStreaming] = useState('')
const convoRef = useRef(convo)
convoRef.current = convo
@@ -86,10 +89,17 @@ export function useChat(onApplyChange?: (change: IssueChange) => Promise<{ ok: b
role: m.from === 'user' ? 'user' : 'assistant',
content: m.text,
}))
setStreaming('')
const unsubscribe = window.commitea.model.onToken((delta) => setStreaming((s) => s + delta))
const finish = () => {
unsubscribe()
setThinking(false)
setStreaming('')
}
window.commitea.model
.chat(wire)
.then((res) => {
setThinking(false)
finish()
if (res.ok) {
setSteps(res.steps.map((s) => s.tool))
setProposals(res.proposals)
@@ -105,7 +115,7 @@ export function useChat(onApplyChange?: (change: IssueChange) => Promise<{ ok: b
}
})
.catch(() => {
setThinking(false)
finish()
setConvo((c) => [...c, { from: 'agent', text: 'I could not reach the model.' }])
})
},
@@ -137,5 +147,5 @@ export function useChat(onApplyChange?: (change: IssueChange) => Promise<{ ok: b
setConvo((c) => [...c, { from: 'agent', text: `Left #${p.change.issue} as it was.` }])
}, [])
return { msgs: [...seed, ...convo], thinking, live, model, steps, proposals, send, approve, dismiss }
return { msgs: [...seed, ...convo], thinking, live, model, steps, proposals, streaming, send, approve, dismiss }
}

View File

@@ -31,18 +31,20 @@ function stringify(result: unknown): string {
}
export async function runAgentTurn(opts: {
complete: (messages: ChatMessage[], tools?: ToolDecl[]) => Promise<CompletionResult>
complete: (messages: ChatMessage[], tools?: ToolDecl[], onToken?: (delta: string) => void) => Promise<CompletionResult>
messages: ChatMessage[]
tools: ToolDecl[]
execute: ToolExecutor
maxSteps?: number
/** Streams content deltas as the model produces prose (final-answer streaming). */
onToken?: (delta: string) => void
}): Promise<AgentTurn> {
const maxSteps = opts.maxSteps ?? DEFAULT_MAX_STEPS
const convo: ChatMessage[] = [...opts.messages]
const steps: AgentStep[] = []
for (let step = 0; step < maxSteps; step++) {
const { content, toolCalls } = await opts.complete(convo, opts.tools)
const { content, toolCalls } = await opts.complete(convo, opts.tools, opts.onToken)
if (toolCalls.length === 0) {
convo.push({ role: 'assistant', content })
return { content, steps, messages: convo }
@@ -63,7 +65,7 @@ export async function runAgentTurn(opts: {
}
// Out of tool budget — force a final prose answer with tools withheld.
const final = await opts.complete(convo, [])
const final = await opts.complete(convo, [], opts.onToken)
convo.push({ role: 'assistant', content: final.content })
return { content: final.content, steps, messages: convo }
}

View File

@@ -25,6 +25,48 @@ describe('createChatClient', () => {
return { fetch, calls }
}
function sseStub(chunks: string[]) {
const calls: { url: string; body: unknown }[] = []
const fetch: FetchLike = (url, init) => {
calls.push({ url, body: init?.body ? JSON.parse(init.body) : undefined })
const enc = new TextEncoder()
const body = new ReadableStream<Uint8Array>({
start(c) {
for (const ch of chunks) c.enqueue(enc.encode(ch))
c.close()
},
})
return Promise.resolve({ ok: true, status: 200, body, json: () => Promise.resolve({}), text: () => Promise.resolve('') })
}
return { fetch, calls }
}
it('streams content deltas via onToken and returns the assembled result', async () => {
const { fetch, calls } = sseStub([
'data: {"choices":[{"delta":{"content":"Right "}}]}\n\n',
'data: {"choices":[{"delta":{"content":"now: #2."}}]}\n\n',
'data: [DONE]\n\n',
])
const client = createChatClient({ baseUrl: 'http://x/v1', model: 'm' }, fetch)
const tokens: string[] = []
const res = await client.complete([{ role: 'user', content: 'now?' }], undefined, (d) => tokens.push(d))
expect(tokens).toEqual(['Right ', 'now: #2.'])
expect(res.content).toBe('Right now: #2.')
expect((calls[0].body as { stream?: boolean }).stream).toBe(true)
})
it('assembles a streamed tool call from argument deltas', async () => {
const { fetch } = sseStub([
'data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"c1","function":{"name":"query_project","arguments":"{\\"view\\""}}]}}]}\n\n',
'data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":":\\"focus\\"}"}}]}}]}\n\n',
'data: [DONE]\n\n',
])
const client = createChatClient({ baseUrl: 'http://x/v1', model: 'm' }, fetch)
const res = await client.complete([{ role: 'user', content: 'x' }], [{ name: 'query_project', description: '', parameters: {} }], () => {})
expect(res.toolCalls).toEqual([{ id: 'c1', name: 'query_project', arguments: '{"view":"focus"}' }])
})
it('POSTs to /chat/completions and parses content', async () => {
const { fetch, calls } = stub({ choices: [{ message: { content: 'the focus is #2' } }] })
const client = createChatClient({ baseUrl: 'http://localhost:1234/v1', model: 'gemma' }, fetch)

View File

@@ -48,8 +48,12 @@ export interface CompletionResult {
toolCalls: ToolCall[]
}
/** Called with each streamed content delta (final-prose streaming). */
export type OnToken = (delta: string) => void
export interface ChatClient {
complete(messages: ChatMessage[], tools?: ToolDecl[]): Promise<CompletionResult>
/** When `onToken` is given, the response streams (SSE) and each content delta is emitted. */
complete(messages: ChatMessage[], tools?: ToolDecl[], onToken?: OnToken): Promise<CompletionResult>
}
/** Map our message shape to the OpenAI wire shape. */
@@ -85,7 +89,7 @@ export function createChatClient(config: ModelConfig, fetchImpl: FetchLike): Cha
const url = `${config.baseUrl.replace(/\/+$/, '')}/chat/completions`
return {
async complete(messages, tools) {
async complete(messages, tools, onToken) {
const body: Record<string, unknown> = {
model: config.model,
messages: messages.map(toWireMessage),
@@ -95,11 +99,13 @@ export function createChatClient(config: ModelConfig, fetchImpl: FetchLike): Cha
body.tools = tools.map(toWireTool)
body.tool_choice = 'auto'
}
const stream = !!onToken
if (stream) body.stream = true
const res = await fetchImpl(url, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Accept: 'application/json',
Accept: stream ? 'text/event-stream' : 'application/json',
...(config.apiKey ? { Authorization: `Bearer ${config.apiKey}` } : {}),
},
body: JSON.stringify(body),
@@ -108,6 +114,8 @@ export function createChatClient(config: ModelConfig, fetchImpl: FetchLike): Cha
const text = await res.text().catch(() => '')
throw new Error(`model completion failed (${res.status}): ${text.slice(0, 200)}`)
}
if (stream && res.body) return readStream(res.body, onToken!)
const json = (await res.json()) as { choices?: { message: RawChoiceMessage }[] }
const msg = json.choices?.[0]?.message
return {
@@ -121,3 +129,56 @@ export function createChatClient(config: ModelConfig, fetchImpl: FetchLike): Cha
},
}
}
/** Streamed tool-call delta: name arrives first, arguments accumulate across chunks. */
interface RawToolDelta {
index: number
id?: string
function?: { name?: string; arguments?: string }
}
/** Parse an OpenAI SSE stream: emit content deltas via onToken, accumulate the final result. */
async function readStream(body: ReadableStream<Uint8Array>, onToken: OnToken): Promise<CompletionResult> {
const reader = body.getReader()
const decoder = new TextDecoder()
let buffer = ''
let content = ''
const toolAcc: { id: string; name: string; arguments: string }[] = []
const handle = (data: string) => {
if (data === '[DONE]') return
let chunk: { choices?: { delta?: { content?: string; tool_calls?: RawToolDelta[] } }[] }
try {
chunk = JSON.parse(data)
} catch {
return
}
const delta = chunk.choices?.[0]?.delta
if (!delta) return
if (delta.content) {
content += delta.content
onToken(delta.content)
}
for (const tc of delta.tool_calls ?? []) {
const slot = (toolAcc[tc.index] ??= { id: '', name: '', arguments: '' })
if (tc.id) slot.id = tc.id
if (tc.function?.name) slot.name = tc.function.name
if (tc.function?.arguments) slot.arguments += tc.function.arguments
}
}
for (;;) {
const { done, value } = await reader.read()
if (done) break
buffer += decoder.decode(value, { stream: true })
const lines = buffer.split('\n')
buffer = lines.pop() ?? ''
for (const line of lines) {
const t = line.trim()
if (t.startsWith('data:')) handle(t.slice(5).trim())
}
}
if (buffer.trim().startsWith('data:')) handle(buffer.trim().slice(5).trim())
return { content, toolCalls: toolAcc.filter((t) => t.name).map((t) => ({ id: t.id, name: t.name, arguments: t.arguments })) }
}

View File

@@ -0,0 +1,88 @@
import { describe, expect, it } from 'vitest'
import { type DependencyEdge, type SchedulableIssue } from '../scheduler/scheduler-v0.js'
import { makespan, scheduleWithCapacity, type Worker } from '../scheduler/scheduler-capacity-v0.js'
import { capacityPerWorkday, parseCapacityConfig } from './capacity-v0.js'
describe('capacity model', () => {
it('capacityPerWorkday = focusFactor × allocation', () => {
expect(capacityPerWorkday({ person: 'a', focusFactor: 0.8, allocation: 0.5 })).toBeCloseTo(0.4)
})
it('parses + clamps config, drops invalid members', () => {
const members = parseCapacityConfig({
members: [
{ person: 'christian', focusFactor: 0.8, allocation: 1 },
{ person: 'ak', focusFactor: 1.5, allocation: -1 }, // clamps to 1 / 0 → zero capacity → dropped
{ focusFactor: 0.8 }, // no person → dropped
],
})
expect(members).toEqual([{ person: 'christian', focusFactor: 0.8, allocation: 1 }])
})
it('returns [] for a non-array/absent members field', () => {
expect(parseCapacityConfig({})).toEqual([])
expect(parseCapacityConfig({ members: 'nope' })).toEqual([])
})
})
function issue(number: number, over: Partial<SchedulableIssue> = {}): SchedulableIssue {
return { number, title: `#${number}`, labels: [], estimateDays: 4, priority: 2, ...over }
}
describe('scheduleWithCapacity', () => {
it('with no workers, falls back to the single serial plan', () => {
const issues = [issue(1), issue(2)]
const plan = scheduleWithCapacity(issues, [], [])
expect(makespan(plan)).toBe(8) // 4 + 4 serial
})
it('parallelizes independent work across lanes (makespan shrinks)', () => {
const issues = [issue(1), issue(2), issue(3), issue(4)] // 4×4d = 16d serial
const one: Worker[] = [{ person: 'a', speed: 1 }]
const two: Worker[] = [
{ person: 'a', speed: 1 },
{ person: 'b', speed: 1 },
]
expect(makespan(scheduleWithCapacity(issues, [], one))).toBe(16)
expect(makespan(scheduleWithCapacity(issues, [], two))).toBe(8) // two lanes → half
})
it('scales duration by a lane speed (slower lane takes longer)', () => {
const plan = scheduleWithCapacity([issue(1, { estimateDays: 4 })], [], [{ person: 'a', speed: 0.5 }])
expect(plan.items[0].durationDays).toBe(8) // 4 / 0.5
expect(plan.items[0].worker).toBe('a')
})
it('routes an issue to its assignee lane', () => {
const issues = [issue(1, { assignee: 'ak' }), issue(2, { assignee: 'sm' })]
const workers: Worker[] = [
{ person: 'ak', speed: 1 },
{ person: 'sm', speed: 1 },
]
const plan = scheduleWithCapacity(issues, [], workers)
const byN = Object.fromEntries(plan.items.map((i) => [i.number, i]))
expect(byN[1].worker).toBe('ak')
expect(byN[2].worker).toBe('sm')
// both start at 0 (different lanes) → parallel
expect(byN[1].startDay).toBe(0)
expect(byN[2].startDay).toBe(0)
})
it('adds a lane for an assignee not in the config (mean speed)', () => {
const plan = scheduleWithCapacity([issue(1, { assignee: 'newbie' })], [], [{ person: 'a', speed: 0.5 }])
expect(plan.items[0].worker).toBe('newbie')
expect(plan.items[0].durationDays).toBe(8) // mean speed 0.5 → 4/0.5
})
it('respects dependencies across lanes (a blocker finishes before its dependent starts)', () => {
const edges: DependencyEdge[] = [{ issue: 2, dependsOn: 1 }]
const workers: Worker[] = [
{ person: 'a', speed: 1 },
{ person: 'b', speed: 1 },
]
const plan = scheduleWithCapacity([issue(1), issue(2)], edges, workers)
const byN = Object.fromEntries(plan.items.map((i) => [i.number, i]))
expect(byN[2].startDay).toBeGreaterThanOrEqual(byN[1].endDay) // #2 waits for #1 even on another lane
})
})

View File

@@ -0,0 +1,49 @@
/**
* Capacity model (#8). A member's throughput is `focusFactor × allocation` =
* ideal person-days of project work delivered per calendar working day. The
* scheduler treats each member as a lane whose task durations are scaled by that
* rate (a 0.5-capacity person takes twice as long on an est/2d task). Config
* lives in the pm-state repo (`capacity/members.yaml/json`, D4); estimates are
* in ideal person-days (pm-state.md).
*/
export interface CapacityMember {
/** gitea login. */
person: string
/** Productive fraction of a working day (0..1). */
focusFactor: number
/** Fraction of that allocated to this project (0..1). */
allocation: number
}
/** Ideal person-days delivered per calendar working day. */
export function capacityPerWorkday(m: CapacityMember): number {
return m.focusFactor * m.allocation
}
function clamp01(n: unknown, fallback: number): number {
return typeof n === 'number' && Number.isFinite(n) ? Math.min(1, Math.max(0, n)) : fallback
}
/**
* Validate a raw capacity config (`{ members: [...] }`) into members. Unknown
* shapes degrade to `[]` (→ the scheduler falls back to a single worker), never
* throw. focusFactor/allocation clamp to [0,1]; a member without a person is dropped.
*/
export function parseCapacityConfig(raw: unknown): CapacityMember[] {
const list = (raw as { members?: unknown })?.members
if (!Array.isArray(list)) return []
const out: CapacityMember[] = []
for (const m of list) {
const r = (m ?? {}) as Record<string, unknown>
const person = typeof r.person === 'string' ? r.person.trim() : ''
if (!person) continue
const member: CapacityMember = {
person,
focusFactor: clamp01(r.focusFactor, 0.8),
allocation: clamp01(r.allocation, 1),
}
if (capacityPerWorkday(member) > 0) out.push(member)
}
return out
}

View File

@@ -120,4 +120,20 @@ describe('forecast', () => {
expect(fitted.coldStart).toBe(false)
expect(fitted.p50Day).toBeLessThan(priors.p50Day)
})
it('capacity lanes parallelize the sim — the makespan shrinks with more workers', () => {
const independent = [issue(1), issue(2), issue(3), issue(4)] // no deps → fully parallelizable
const serial = forecast(independent, [], { trials: 2000, seed: 7 })
const twoLanes = forecast(independent, [], {
trials: 2000,
seed: 7,
workers: [
{ person: 'a', speed: 1 },
{ person: 'b', speed: 1 },
],
})
// two equal lanes ≈ half the serial landing (independent work)
expect(twoLanes.p80Day).toBeLessThan(serial.p80Day)
expect(twoLanes.p80Day).toBeLessThan(serial.p80Day * 0.7)
})
})

View File

@@ -18,6 +18,7 @@ import {
schedule,
type SchedulableIssue,
} from '../scheduler/scheduler-v0.js'
import { type LaneInputs, layoutOnLanes, resolveLanes, type Worker } from '../scheduler/scheduler-capacity-v0.js'
export interface LognormalPrior {
/** Median log-ratio: sampled median duration = estimate * e^mu. */
@@ -76,6 +77,12 @@ export interface ForecastOptions {
* the sim; otherwise the code-resident cold-start priors do.
*/
model?: DurationModel
/**
* Capacity lanes. When provided, each trial lays sampled durations across the
* lanes (parallel) instead of a single serial worker — the makespan shrinks
* toward the critical path. Empty/absent → single serial worker.
*/
workers?: Worker[]
}
/** Resolve the lognormal params for an estimate, preferring a fitted model. */
@@ -158,17 +165,34 @@ export function forecast(
const priors = order.map((it) => durationParams(it.durationDays, options.model))
const rng = mulberry32(seed)
// endByRank[k][t] = working day the (k+1)-th scheduled issue completes on trial t.
// Capacity lanes, when configured — the per-trial layout goes parallel.
const laneInputs: LaneInputs = {
order: order.map((it) => it.number),
blockedBy: new Map(order.map((it) => [it.number, it.blockedBy])),
assignee: new Map(issues.map((i) => [i.number, i.assignee ?? null])),
}
const lanes =
options.workers && options.workers.length ? resolveLanes(options.workers, [...laneInputs.assignee.values()]) : null
// endByRank[k][t] = working day the (k+1)-th issue *to finish* completes on trial t.
const endByRank: number[][] = Array.from({ length: n }, () => new Array<number>(trials))
for (let t = 0; t < trials; t++) {
const sampled = order.map((it, k) => it.durationDays * Math.exp(priors[k].mu + priors[k].sigma * standardNormal(rng)))
if (lanes) {
// parallel: lay the sampled durations across lanes, then sort finish days
const byNumber = new Map(order.map((it, k) => [it.number, sampled[k]]))
const { finishAt } = layoutOnLanes(laneInputs, lanes, (num) => byNumber.get(num)!)
const finishes = order.map((it) => finishAt.get(it.number)!).sort((a, b) => a - b)
for (let k = 0; k < n; k++) endByRank[k][t] = finishes[k]
} else {
// single serial worker: cumulative sum (already sorted ascending)
let cursor = 0
for (let k = 0; k < n; k++) {
const p = priors[k]
const sampled = order[k].durationDays * Math.exp(p.mu + p.sigma * standardNormal(rng))
cursor += sampled
cursor += sampled[k]
endByRank[k][t] = cursor
}
}
}
const curve: BurnUpPoint[] = endByRank.map((row, k) => {
const sorted = [...row].sort((a, b) => a - b)

View File

@@ -31,6 +31,8 @@ export interface GiteaHttpResponse {
status: number
json(): Promise<unknown>
text(): Promise<string>
/** Present on the real fetch Response; used for SSE streaming (chat). */
body?: ReadableStream<Uint8Array> | null
}
export type FetchLike = (url: string, init?: GiteaRequestInit) => Promise<GiteaHttpResponse>

View File

@@ -49,6 +49,12 @@ export type {
SchedulePlan,
} from './scheduler/scheduler-v0.js'
export { makespan, scheduleWithCapacity } from './scheduler/scheduler-capacity-v0.js'
export type { Worker } from './scheduler/scheduler-capacity-v0.js'
export { capacityPerWorkday, parseCapacityConfig } from './capacity/capacity-v0.js'
export type { CapacityMember } from './capacity/capacity-v0.js'
export {
COLD_START_PRIORS,
durationParams,
@@ -80,7 +86,7 @@ export type {
} from './calibration/calibration-v0.js'
export { createChatClient } from './agent/chat-client.js'
export type { ChatClient, ChatMessage, CompletionResult, ModelConfig, ToolCall, ToolDecl } from './agent/chat-client.js'
export type { ChatClient, ChatMessage, CompletionResult, ModelConfig, OnToken, ToolCall, ToolDecl } from './agent/chat-client.js'
export { pickModel } from './agent/model-router.js'
export type { ModelRouter, TaskKind } from './agent/model-router.js'
export { runAgentTurn } from './agent/agent-loop.js'

View File

@@ -0,0 +1,119 @@
/**
* Capacity-aware scheduler (#8). Reuses the single-worker scheduler's topological
* order + critical-path marking, then re-lays the work across lanes: an issue
* runs on its assignee's lane (or the earliest-free lane when unassigned), its
* duration scaled by that lane's speed (ideal person-days/workday). Makespan
* shrinks toward the critical path as lanes are added. Deterministic; falls back
* to the single serial worker when no capacity is configured. The lane layout is
* factored so the Monte Carlo forecast reuses it per trial with sampled durations.
*/
import {
type DependencyEdge,
schedule,
type SchedulableIssue,
type SchedulePlan,
} from './scheduler-v0.js'
/** A scheduling lane: a person and their throughput (ideal person-days / workday). */
export interface Worker {
person: string
speed: number
}
/** The topological order + per-issue relations the layout needs (from the base plan). */
export interface LaneInputs {
order: number[]
blockedBy: Map<number, number[]>
assignee: Map<number, string | null>
}
/** Ensure a lane exists for every assignee; unconfigured assignees get the mean speed. */
export function resolveLanes(workers: Worker[], assignees: (string | null)[]): Worker[] {
const lanes = [...workers]
const known = new Set(lanes.map((w) => w.person))
const meanSpeed = lanes.length ? lanes.reduce((s, w) => s + w.speed, 0) / lanes.length : 1
for (const a of assignees) {
if (a && !known.has(a)) {
lanes.push({ person: a, speed: meanSpeed })
known.add(a)
}
}
return lanes
}
function pickWorker(lanes: Worker[], freeAt: Map<string, number>, assignee: string | null | undefined): Worker {
if (assignee) {
const own = lanes.find((w) => w.person === assignee)
if (own) return own
}
let best = lanes[0]
for (const w of lanes) if (freeAt.get(w.person)! < freeAt.get(best.person)!) best = w
return best
}
/**
* Lay a topologically-ordered set out across lanes. `duration(n)` supplies each
* issue's duration for this layout (estimate, or a sampled value in a MC trial).
* Returns each issue's finish day + the lane it ran on. Order guarantees a
* dependency is always laid out before its dependents.
*/
export function layoutOnLanes(
inputs: LaneInputs,
lanes: Worker[],
duration: (n: number) => number,
): { finishAt: Map<number, number>; startAt: Map<number, number>; laneOf: Map<number, string> } {
const freeAt = new Map(lanes.map((w) => [w.person, 0]))
const finishAt = new Map<number, number>()
const startAt = new Map<number, number>()
const laneOf = new Map<number, string>()
for (const n of inputs.order) {
const worker = pickWorker(lanes, freeAt, inputs.assignee.get(n))
const blockers = inputs.blockedBy.get(n) ?? []
const depFinish = blockers.length ? Math.max(...blockers.map((d) => finishAt.get(d) ?? 0)) : 0
const start = Math.max(freeAt.get(worker.person)!, depFinish)
const end = start + duration(n) / worker.speed
startAt.set(n, start)
finishAt.set(n, end)
laneOf.set(n, worker.person)
freeAt.set(worker.person, end)
}
return { finishAt, startAt, laneOf }
}
/**
* Capacity-aware plan. An empty `workers` means no capacity is configured → the
* single-worker plan verbatim.
*/
export function scheduleWithCapacity(
issues: SchedulableIssue[],
edges: DependencyEdge[],
workers: Worker[],
): SchedulePlan {
if (workers.length === 0) return schedule(issues, edges)
const base = schedule(issues, edges)
if (base.cycle || base.items.length === 0) return base
const inputs: LaneInputs = {
order: base.items.map((it) => it.number),
blockedBy: new Map(base.items.map((it) => [it.number, it.blockedBy])),
assignee: new Map(issues.map((i) => [i.number, i.assignee ?? null])),
}
const lanes = resolveLanes(workers, [...inputs.assignee.values()])
const durationOf = new Map(base.items.map((it) => [it.number, it.durationDays]))
const { finishAt, startAt, laneOf } = layoutOnLanes(inputs, lanes, (n) => durationOf.get(n)!)
const items = base.items.map((it) => ({
...it,
startDay: startAt.get(it.number)!,
endDay: finishAt.get(it.number)!,
durationDays: finishAt.get(it.number)! - startAt.get(it.number)!,
worker: laneOf.get(it.number),
}))
return { items, cycle: null }
}
/** Makespan (last finish) of a plan — the project's landing day. */
export function makespan(plan: SchedulePlan): number {
return plan.items.reduce((m, it) => Math.max(m, it.endDay), 0)
}

View File

@@ -21,6 +21,8 @@ export interface SchedulableIssue {
estimateDays: number | null
/** 1 (most urgent) … 4; null when unset. */
priority: number | null
/** gitea assignee login, for capacity-aware lane routing; optional. */
assignee?: string | null
}
export interface DependencyEdge {
@@ -45,6 +47,8 @@ export interface ScheduledItem {
/** On a longest-duration dependency chain. */
critical: boolean
rationale: string
/** Lane it's scheduled on, when capacity-aware; absent for the single-worker plan. */
worker?: string
}
export interface SchedulePlan {