Dynamic workflow run'larını izler: yan pane'de aşamalar, başlayan/aktif/biten agent sayısı ve her agent'ın son adımı; /wf ile metin özeti (cloud session'da da…

hooks/register.js 321 lines1// wf-monitor: dynamic workflow run'larını izler.
2// Veri kaynakları:
3// - Workflow tool sonucu: taskId, workflowName, transcriptDir, scriptPath
4// - transcriptDir/agent-<id>.jsonl: başlayan agent'lar (mtime → aktif mi)
5// - transcriptDir/journal.jsonl: dönen agent sonuçları (satır sayısı → biten)
6// - workflow agent'larının tool.call event'leri (agentId): son adım
7// - Stop / SubagentStop hook'undaki background_tasks: run bitti mi
8// Çıktılar:
9// - Terminal / Desktop: yan pane
10// - Her yerde (cloud dahil): /wf metin özeti
11// - WF_MONITOR_URL + WF_MONITOR_TOKEN tanımlıysa: durumu dashboard'a POST eder
12
13const PANE_ID = 'wf-monitor'
14const POLL_MS = 2000
15const ACTIVE_MS = 20000
16const HEARTBEAT_MS = 15000
17
18const runs = new Map() // taskId -> run
19const lastAction = new Map() // agentId -> metin
20let host = null // session.start'ta kurulan, mods API çağrılarını saran fonksiyonlar
21let drawing = false
22let paneOpen = false
23let stopPolling = null
24let nowMs = 0
25let lastKey = ''
26let lastPushAt = 0
27
28function short(s, n) {
29 const t = String(s ?? '').replace(/\s+/g, ' ').trim()
30 return t.length > n ? t.slice(0, n - 1) + '…' : t
31}
32
33function describe(e) {
34 switch (e.tool) {
35 case 'Bash':
36 return 'Bash: ' + short(e.command, 48)
37 case 'Read':
38 case 'Edit':
39 case 'Write':
40 case 'MultiEdit':
41 return e.tool + ' ' + short(String(e.file_path ?? '').split('/').pop(), 40)
42 case 'Grep':
43 case 'Glob':
44 return e.tool + ' ' + short(e.pattern, 36)
45 default:
46 return String(e.tool)
47 }
48}
49
50function since(ms) {
51 const s = Math.max(0, Math.round(ms / 1000))
52 const m = Math.floor(s / 60)
53 return m ? `${m}dk ${s % 60}sn` : `${s}sn`
54}
55
56function phasesOf(source) {
57 const meta = /export\s+const\s+meta\s*=\s*\{([\s\S]*?)\n\}/.exec(source)
58 if (!meta) return []
59 return [...meta[1].matchAll(/title:\s*['"`]([^'"`]+)['"`]/g)].map(m => m[1])
60}
61
62function hasRunning() {
63 for (const r of runs.values()) if (r.status === 'running') return true
64 return false
65}
66
67function redraw() {
68 if (host && drawing && paneOpen) host.invalidate()
69}
70
71function view(r) {
72 const end = r.endedAt ?? nowMs
73 return {
74 taskId: r.taskId,
75 name: r.name,
76 status: r.status,
77 elapsedMs: Math.max(0, end - r.startedAt),
78 phases: r.phases,
79 done: r.done,
80 lastLabel: r.lastLabel,
81 agents: r.agents.map(a => ({
82 id: a.id.slice(0, 8),
83 active: end - a.mtime < ACTIVE_MS,
84 last: lastAction.get(a.id) ?? null,
85 })),
86 }
87}
88
89async function push() {
90 if (!host || !host.post || runs.size === 0) return
91 const list = [...runs.values()].slice(-10).map(view)
92 const key = JSON.stringify(list.map(({ elapsedMs, ...rest }) => rest))
93 if (key === lastKey && nowMs - lastPushAt < HEARTBEAT_MS) return
94 lastKey = key
95 lastPushAt = nowMs
96 await host.post({ session: host.session, runs: list }).catch(() => {})
97}
98
99function summaryText() {
100 if (runs.size === 0) return 'Bu session\'da workflow çalışmadı.'
101 const lines = []
102 for (const r of runs.values()) {
103 const v = view(r)
104 const active = v.agents.filter(a => a.active)
105 lines.push(
106 `${v.name} · ${v.status} · ${since(v.elapsedMs)} · ` +
107 `başlayan ${v.agents.length} / aktif ${active.length} / biten ${v.done}`,
108 )
109 if (v.phases.length) lines.push(' aşamalar: ' + v.phases.join(' → '))
110 if (v.lastLabel) lines.push(' son dönen: ' + v.lastLabel)
111 if (v.status === 'running') for (const a of active) lines.push(` ${a.id} ${a.last ?? '…'}`)
112 }
113 return lines.join('\n')
114}
115
116async function refresh(run) {
117 if (!host || !run.transcriptDir) return
118 let entries
119 try {
120 entries = await host.list(run.transcriptDir)
121 } catch {
122 return
123 }
124 const agents = []
125 for (const ent of entries) {
126 const m = /^agent-(.+)\.jsonl$/.exec(ent.name)
127 if (!m || ent.kind !== 'file') continue
128 let mtime = 0
129 try {
130 mtime = (await host.stat(`${run.transcriptDir}/${ent.name}`)).mtimeMs
131 } catch {}
132 agents.push({ id: m[1], mtime })
133 }
134 run.agents = agents
135 if (entries.some(ent => ent.name === 'journal.jsonl')) {
136 try {
137 const rows = (await host.read(`${run.transcriptDir}/journal.jsonl`)).split('\n').filter(Boolean)
138 run.done = rows.length
139 try {
140 const last = JSON.parse(rows[rows.length - 1])
141 run.lastLabel = last.label ?? last.opts?.label ?? run.lastLabel
142 } catch {}
143 } catch {}
144 }
145}
146
147async function finish(run, status) {
148 if (!host || run.status !== 'running') return
149 run.status = status
150 run.endedAt = await host.now()
151 nowMs = run.endedAt
152 await refresh(run) // son durumu bir kez daha oku
153 const msg = `workflow ${run.name}: ${status} (${since(run.endedAt - run.startedAt)}, ${run.done} agent sonucu)`
154 if (drawing) host.toast(msg)
155 host.log(msg)
156 if (!hasRunning() && stopPolling) {
157 stopPolling()
158 stopPolling = null
159 }
160 redraw()
161 await push()
162}
163
164async function poll() {
165 if (!host) return
166 nowMs = await host.now()
167 for (const run of runs.values()) if (run.status === 'running') await refresh(run)
168 redraw()
169 await push()
170}
171
172function startPolling() {
173 if (host && !stopPolling) stopPolling = host.every(POLL_MS, () => void poll().catch(() => {}))
174}
175
176async function openPane() {
177 if (!host || !drawing) return
178 paneOpen = true
179 await host.open().catch(() => {})
180}
181
182async function reconcile(tasks) {
183 if (!Array.isArray(tasks)) return
184 for (const run of runs.values()) {
185 if (run.status !== 'running') continue
186 const t = tasks.find(x => x.id === run.taskId || (x.type === 'workflow' && x.name === run.name))
187 // background_tasks yalnızca hâlâ süren işleri listeler: listede yoksa bitmiştir
188 if (!t) await finish(run, 'bitti')
189 else if (!['running', 'pending'].includes(t.status)) await finish(run, t.status)
190 }
191}
192
193export function register(on) {
194 on('session.start', async ($, e, next) => {
195 const url = await $.env.get('WF_MONITOR_URL')
196 const token = await $.env.get('WF_MONITOR_TOKEN')
197 const repo = await $.session.repo().catch(() => null)
198 const id = await $.session.id()
199 host = {
200 now: () => $.clock.now(),
201 every: (ms, fn) => $.clock.every(ms, fn),
202 list: path => $.fs.list(path),
203 stat: path => $.fs.stat(path),
204 read: path => $.fs.read(path),
205 invalidate: () => $.ui.invalidate('ui.render'),
206 toast: text => $.ui.toast(text),
207 log: text => $.ui.log(text),
208 open: () => $.ui.open({ id: PANE_ID, title: 'Workflow' }),
209 session: { id, repo: (repo?.remote ?? repo?.root ?? e.cwd ?? '').replace(/\.git$/, '').split('/').slice(-2).join('/') },
210 post:
211 url && token
212 ? body =>
213 $.http.fetch(url, {
214 method: 'POST',
215 headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${token}` },
216 body: JSON.stringify(body),
217 })
218 : null,
219 }
220 nowMs = await $.clock.now()
221 const surfaces = await $.session.surfaces()
222 drawing = surfaces.some(s => s === 'terminal' || s === 'desktop')
223 await $.command
224 .register({ name: 'wf', description: 'Workflow durumunu göster', immediate: true })
225 .catch(() => {})
226 return next(e)
227 })
228
229 // Workflow başlatıldığında run'ı kaydet
230 on('tool.call', { tool: 'Workflow' }, async ($, e, next) => {
231 const res = await next(e)
232 const r = res?.result
233 if (r && r.taskId && r.status === 'async_launched') {
234 const startedAt = await $.clock.now()
235 let phases = []
236 if (r.scriptPath) {
237 try {
238 phases = phasesOf(await $.fs.read(r.scriptPath))
239 } catch {}
240 }
241 runs.set(r.taskId, {
242 taskId: r.taskId,
243 name: r.workflowName ?? e.name ?? 'workflow',
244 transcriptDir: r.transcriptDir,
245 phases,
246 status: 'running',
247 startedAt,
248 endedAt: null,
249 agents: [],
250 done: 0,
251 lastLabel: null,
252 })
253 nowMs = startedAt
254 startPolling()
255 await openPane()
256 await push()
257 }
258 return res
259 })
260
261 // Workflow agent'larının son adımı
262 on('tool.call', async ($, e, next) => {
263 if (e.agentId && hasRunning()) lastAction.set(e.agentId, describe(e))
264 return next(e)
265 })
266
267 on('classic.Stop', async ($, e, next) => {
268 await reconcile(e.background_tasks)
269 return next(e)
270 })
271
272 on('classic.SubagentStop', async ($, e, next) => {
273 await reconcile(e.background_tasks)
274 return next(e)
275 })
276
277 // Çizen yüzeylerde bitiş bildirimi satırı da durumu söyler
278 on('ui.render', { component: 'UserMessage' }, async ($, e, next) => {
279 const task = e.props?.task
280 const run = task?.id ? runs.get(task.id) : undefined
281 if (run && run.status === 'running' && task.status) await finish(run, task.status)
282 return next(e)
283 })
284
285 on('command.run', { command: 'wf' }, async ($, e, next) => {
286 nowMs = await $.clock.now()
287 for (const run of runs.values()) if (run.status === 'running') await refresh(run)
288 if (runs.size > 0) await openPane()
289 return { text: summaryText() }
290 })
291
292 on('ui.close', { id: PANE_ID }, async ($, e, next) => {
293 paneOpen = false
294 return next(e)
295 })
296
297 on('ui.render', { component: 'Pane' }, async ($, e, next) => {
298 if (e.requestId !== PANE_ID) return next(e)
299 const { Box, Text } = $.ui.resolve(e)
300 const width = Math.max(20, (e.props.bodyColumns ?? 60) - 2)
301 const blocks = []
302 if (runs.size === 0) blocks.push(h(Text, { dimColor: true }, 'Henüz workflow yok.'))
303 for (const r of [...runs.values()].reverse()) {
304 const v = view(r)
305 const active = v.agents.filter(a => a.active)
306 const color = v.status === 'running' ? 'yellow' : v.status === 'failed' ? 'red' : 'green'
307 const rows = [
308 h(Text, { bold: true, color }, short(`${v.name} · ${v.status} · ${since(v.elapsedMs)}`, width)),
309 h(Text, null, `başlayan ${v.agents.length} aktif ${active.length} biten ${v.done}`),
310 ]
311 if (v.phases.length) rows.push(h(Text, { dimColor: true }, short(v.phases.join(' → '), width)))
312 if (v.lastLabel) rows.push(h(Text, { dimColor: true }, short('son dönen: ' + v.lastLabel, width)))
313 if (v.status === 'running') {
314 for (const a of active) rows.push(h(Text, null, short(`${a.id} ${a.last ?? '…'}`, width)))
315 }
316 blocks.push(h(Box, { key: r.taskId, flexDirection: 'column' }, ...rows))
317 }
318 return h(Box, { flexDirection: 'column', gap: 1 }, ...blocks)
319 })
320}
321