-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathClaudeEventAdapter.cs
More file actions
352 lines (332 loc) · 17.8 KB
/
Copy pathClaudeEventAdapter.cs
File metadata and controls
352 lines (332 loc) · 17.8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
using System.Text.Json;
using CodingAgentRunner.Events;
namespace CodingAgentRunner.Adapters;
/// <summary>
/// Maps Claude Code's <c>--output-format stream-json --verbose</c> NDJSON frames
/// onto the typed <see cref="CliRunEvent"/> contract. Stateless with respect to a
/// single frame; the caller holds whatever state it needs across frames (e.g. the
/// captured session id).
///
/// <para>
/// Frame-to-event mapping. Each line of stdout from claude is one JSON object
/// whose <c>type</c> we switch on:
/// </para>
/// <list type="bullet">
/// <item><c>system</c> with <c>subtype=init</c> → <see cref="CliRunEvent.SessionStarted"/> (carries the id).</item>
/// <item><c>system</c> with a <c>task_*</c> subtype (a delegated subtask reporting in) → <see cref="CliRunEvent.Heartbeat"/>.</item>
/// <item><c>system</c> with any other subtype → <see cref="CliRunEvent.SessionInitializing"/>.</item>
/// <item><c>rate_limit_event</c> → <see cref="CliRunEvent.RateLimitObserved"/>.</item>
/// <item><c>assistant</c> with <c>tool_use</c> part → <see cref="CliRunEvent.ToolStarted"/>.</item>
/// <item><c>assistant</c> with <c>text</c> part → <see cref="CliRunEvent.OutputDelta"/>.</item>
/// <item><c>user</c> with <c>tool_result</c> part → <see cref="CliRunEvent.ToolCompleted"/> (<c>is_error</c> flag).</item>
/// <item><c>result</c> with <c>is_error=true</c> → <see cref="CliRunEvent.TurnFailed"/>.</item>
/// <item><c>result</c> with <c>is_error=false</c> → <see cref="CliRunEvent.TurnCompleted"/> — the CLI's real completion signal.</item>
/// <item>everything else → <see cref="CliRunEvent.Unknown"/> with the raw <c>type</c> as sample.</item>
/// </list>
/// </summary>
public static class ClaudeEventAdapter
{
/// <summary>
/// Map one stream-json line to zero or more <see cref="CliRunEvent"/> instances.
/// A single line can produce multiple events (an assistant frame with a text
/// part AND a tool_use part yields <see cref="CliRunEvent.OutputDelta"/>
/// followed by <see cref="CliRunEvent.ToolStarted"/>).
///
/// Lines that are not JSON, not stdout, or empty produce no events. Lines whose
/// top-level <c>type</c> we cannot classify produce a single
/// <see cref="CliRunEvent.Unknown"/> with a 200-char prefix of the original
/// line as <see cref="CliRunEvent.Unknown.Sample"/>.
/// </summary>
public static IEnumerable<CliRunEvent> Map(string jsonLine, string runId)
{
if (string.IsNullOrWhiteSpace(jsonLine) || jsonLine[0] != '{') yield break;
JsonDocument? doc = null;
try { doc = JsonDocument.Parse(jsonLine); }
catch { yield break; }
using var _ = doc;
if (doc.RootElement.ValueKind != JsonValueKind.Object) yield break;
var root = doc.RootElement;
var type = root.TryGetProperty("type", out var t) && t.ValueKind == JsonValueKind.String ? t.GetString() : null;
switch (type)
{
case "system":
{
var subtype = root.TryGetProperty("subtype", out var st) ? st.GetString() : null;
var sessionId = root.TryGetProperty("session_id", out var sid) ? sid.GetString() : null;
if (string.Equals(subtype, "init", StringComparison.Ordinal))
{
// The init frame is Claude's session handshake AND its context
// report: it names the model it actually loaded, the effective
// permission mode, and the working directory. Carry those along
// instead of discarding them — a consumer rendering "what did
// the agent actually see" reads them from here.
yield return new CliRunEvent.SessionStarted(sessionId)
{
RunId = runId,
Model = TolerantString(root, "model"),
PermissionMode = TolerantString(root, "permissionMode", "permission_mode"),
Cwd = TolerantString(root, "cwd"),
ApiKeySource = TolerantString(root, "apiKeySource", "api_key_source"),
};
}
else if (IsDelegationLifecycle(subtype))
{
yield return new CliRunEvent.Heartbeat { RunId = runId };
}
else
{
yield return new CliRunEvent.SessionInitializing { RunId = runId };
}
yield break;
}
case "rate_limit_event":
{
// Tolerant on purpose: Claude has shipped camelCase fields, other
// stream surfaces commonly use snake_case, resets timestamps have
// appeared both as unix seconds and as ISO-8601 strings, and
// booleans as real booleans and as "true"/"false" strings. A frame
// this parser cannot read must degrade to nulls/zero, never fail
// the read loop or drop the event.
var info = TolerantObject(root, "rate_limit_info", "rateLimitInfo");
if (info.ValueKind == JsonValueKind.Object)
{
// Normalize the window label onto the built-in Claude probe's
// vocabulary (five_hour → 5-hour, seven_day → weekly):
// QuotaService.Observe merges by label, so a live event must
// land in the SAME QuotaWindow the probe wrote.
var window = Quota.ClaudeOAuthUsageProbe.WindowLabel(
TolerantString(info, "rateLimitType", "rate_limit_type", "window"));
var status = TolerantString(info, "status");
var resetsAt = TolerantUnixSeconds(info, "resetsAt", "resets_at");
var overage = TolerantString(info, "overageStatus", "overage_status");
var using_ = TolerantBoolean(info, "isUsingOverage", "is_using_overage");
yield return new CliRunEvent.RateLimitObserved(window, status, resetsAt, overage, using_) { RunId = runId };
}
yield break;
}
case "assistant":
{
if (!root.TryGetProperty("message", out var msg)) yield break;
if (!msg.TryGetProperty("content", out var content) || content.ValueKind != JsonValueKind.Array) yield break;
// We do not emit a separate TurnStarted: claude has no distinct
// "turn began" frame — OutputDelta on the first assistant text plays
// that role.
foreach (var part in content.EnumerateArray())
{
var partType = part.TryGetProperty("type", out var pt) ? pt.GetString() : null;
if (partType == "text")
{
var text = part.TryGetProperty("text", out var txt) ? txt.GetString() ?? "" : "";
if (!string.IsNullOrEmpty(text))
yield return new CliRunEvent.OutputDelta(text) { RunId = runId };
}
else if (partType == "tool_use")
{
var name = part.TryGetProperty("name", out var n) ? n.GetString() ?? "Tool" : "Tool";
var argument = ExtractToolArgument(part);
yield return new CliRunEvent.ToolStarted(name, argument) { RunId = runId };
// TodoWrite is Claude's native task-plan frame. Surface it as
// a typed PlanUpdated in addition to the tool call so a
// consumer can persist a plan snapshot.
if (string.Equals(name, "TodoWrite", StringComparison.Ordinal)
&& TryExtractTodoPlan(part, out var items))
{
yield return new CliRunEvent.PlanUpdated("claude/TodoWrite", items) { RunId = runId };
}
}
else if (partType == "thinking")
{
// Extended-thinking — dropped from the typed stream. The raw
// bytes still reach the on-disk output log.
}
// unknown part types: ignored here.
}
yield break;
}
case "user":
{
if (!root.TryGetProperty("message", out var msg)) yield break;
if (!msg.TryGetProperty("content", out var content) || content.ValueKind != JsonValueKind.Array) yield break;
foreach (var part in content.EnumerateArray())
{
var partType = part.TryGetProperty("type", out var pt) ? pt.GetString() : null;
if (partType != "tool_result") continue;
var isError = part.TryGetProperty("is_error", out var ie) && ie.ValueKind == JsonValueKind.True;
var name = "tool"; // tool_use_id is opaque; we report a generic name
var firstLine = ExtractFirstLine(part);
yield return new CliRunEvent.ToolCompleted(name, isError, firstLine) { RunId = runId };
}
yield break;
}
case "result":
{
var isError = root.TryGetProperty("is_error", out var ie) && ie.ValueKind == JsonValueKind.True;
if (isError)
{
var subtype = root.TryGetProperty("subtype", out var st) ? st.GetString() : "error";
yield return new CliRunEvent.TurnFailed(subtype ?? "error") { RunId = runId };
}
else
{
var usage = root.TryGetProperty("usage", out var u) && u.ValueKind == JsonValueKind.Object
? FormatUsage(u)
: null;
yield return new CliRunEvent.TurnCompleted(usage) { RunId = runId };
}
yield break;
}
default:
{
yield return new CliRunEvent.Unknown(Truncate(jsonLine, 200), jsonLine) { RunId = runId };
yield break;
}
}
}
/// <summary>
/// Pull the plan items out of a <c>TodoWrite</c> tool_use part. The input shape
/// is <c>{"todos":[{"content","status","activeForm"}]}</c>; we use the
/// imperative <c>content</c> as the stable title and normalize the status.
/// Returns false when the frame carries no usable todos so the caller emits no
/// PlanUpdated.
/// </summary>
private static bool TryExtractTodoPlan(JsonElement toolPart, out IReadOnlyList<PlanFrameItem> items)
{
items = System.Array.Empty<PlanFrameItem>();
if (!toolPart.TryGetProperty("input", out var input) || input.ValueKind != JsonValueKind.Object) return false;
if (!input.TryGetProperty("todos", out var todos) || todos.ValueKind != JsonValueKind.Array) return false;
var list = new List<PlanFrameItem>();
foreach (var todo in todos.EnumerateArray())
{
if (todo.ValueKind != JsonValueKind.Object) continue;
var content = todo.TryGetProperty("content", out var c) ? c.GetString() : null;
if (string.IsNullOrWhiteSpace(content)) continue;
var status = todo.TryGetProperty("status", out var s) ? s.GetString() : null;
list.Add(new PlanFrameItem(PlanItemId.From(content), content!.Trim(), PlanItemStatus.Normalize(status)));
}
if (list.Count == 0) return false;
items = list;
return true;
}
/// <summary>
/// Whether a <c>system</c> subtype is a delegated subtask reporting in, rather than
/// a session handshake. Claude Code streams <c>task_started</c> /
/// <c>task_progress</c> / <c>task_updated</c> / <c>task_notification</c> while a
/// subagent works, and the family is namespaced, so the prefix keeps a subtype
/// added later on the right side of the split.
/// <para>
/// These map to <see cref="CliRunEvent.Heartbeat"/>, not to
/// <see cref="CliRunEvent.SessionInitializing"/>, and both halves of that matter.
/// The phase must not fall back out of <c>ToolExecuting</c> while the delegation
/// tool is still running — that would swap the watchdog's tool budget for the
/// tighter session-initializing one at exactly the point a run is doing its
/// longest-running thing. And <c>task_progress</c> is literally the subtask
/// reporting progress, so it has to reset the silence clock; read as a phase-only
/// event it counts as silence, and a long delegated sweep that streams nothing else
/// gets classified hung while it is visibly working.
/// </para>
/// </summary>
private static bool IsDelegationLifecycle(string? subtype)
=> subtype is not null && subtype.StartsWith("task_", StringComparison.Ordinal);
private static string? ExtractToolArgument(JsonElement toolPart)
{
if (!toolPart.TryGetProperty("input", out var input) || input.ValueKind != JsonValueKind.Object) return null;
// Common tool argument keys we care about for the typed event. `subagent_type`
// and `description` come last and serve the delegation tool (`Agent` in Claude
// Code 2.1.220, `Task` in earlier versions): a delegated subtask otherwise
// reports no argument at all, which hides which cheap agent ran.
foreach (var key in new[] { "file_path", "path", "command", "pattern", "url", "query", "subagent_type", "description" })
{
if (input.TryGetProperty(key, out var v) && v.ValueKind == JsonValueKind.String)
{
var s = v.GetString();
return string.IsNullOrWhiteSpace(s) ? null : s;
}
}
return null;
}
private static string? ExtractFirstLine(JsonElement toolResultPart)
{
if (!toolResultPart.TryGetProperty("content", out var c)) return null;
string? text = null;
if (c.ValueKind == JsonValueKind.String) text = c.GetString();
else if (c.ValueKind == JsonValueKind.Array)
{
foreach (var p in c.EnumerateArray())
{
if (p.TryGetProperty("type", out var pt) && pt.GetString() == "text"
&& p.TryGetProperty("text", out var tx))
{
text = tx.GetString();
break;
}
}
}
if (string.IsNullOrEmpty(text)) return null;
var idx = text!.IndexOf('\n');
return idx >= 0 ? text[..idx] : text;
}
private static string FormatUsage(JsonElement usage)
{
var input = usage.TryGetProperty("input_tokens", out var i) && i.TryGetInt64(out var iv) ? iv : 0;
var output = usage.TryGetProperty("output_tokens", out var o) && o.TryGetInt64(out var ov) ? ov : 0;
var cacheRead = usage.TryGetProperty("cache_read_input_tokens", out var cr) && cr.TryGetInt64(out var crv) ? crv : 0;
return $"input={input} output={output} cache_read={cacheRead}";
}
private static string Truncate(string s, int max)
=> s.Length <= max ? s : s[..max];
// ── tolerant field readers (rate_limit_event / init-frame context) ──
// Field-name casing and value types have drifted across Claude Code
// releases; these readers accept every observed form and degrade to
// null/zero/false instead of dropping the frame.
private static JsonElement TolerantObject(JsonElement parent, params string[] names)
{
if (parent.ValueKind != JsonValueKind.Object) return default;
foreach (var name in names)
if (parent.TryGetProperty(name, out var value) && value.ValueKind == JsonValueKind.Object)
return value;
return default;
}
private static string? TolerantString(JsonElement parent, params string[] names)
{
if (parent.ValueKind != JsonValueKind.Object) return null;
foreach (var name in names)
if (parent.TryGetProperty(name, out var value)
&& value.ValueKind == JsonValueKind.String
&& !string.IsNullOrWhiteSpace(value.GetString()))
return value.GetString();
return null;
}
private static bool TolerantBoolean(JsonElement parent, params string[] names)
{
if (parent.ValueKind != JsonValueKind.Object) return false;
foreach (var name in names)
{
if (!parent.TryGetProperty(name, out var value)) continue;
if (value.ValueKind == JsonValueKind.True) return true;
if (value.ValueKind == JsonValueKind.False) return false;
if (value.ValueKind == JsonValueKind.String
&& bool.TryParse(value.GetString(), out var parsed))
return parsed;
}
return false;
}
private static long TolerantUnixSeconds(JsonElement parent, params string[] names)
{
if (parent.ValueKind != JsonValueKind.Object) return 0;
foreach (var name in names)
{
if (!parent.TryGetProperty(name, out var value)) continue;
if (value.ValueKind == JsonValueKind.Number && value.TryGetInt64(out var number))
return number;
if (value.ValueKind != JsonValueKind.String) continue;
var raw = value.GetString();
if (long.TryParse(raw, System.Globalization.NumberStyles.Integer,
System.Globalization.CultureInfo.InvariantCulture, out number))
return number;
if (DateTimeOffset.TryParse(raw, System.Globalization.CultureInfo.InvariantCulture,
System.Globalization.DateTimeStyles.AssumeUniversal, out var timestamp))
return timestamp.ToUnixTimeSeconds();
}
return 0;
}
}