1use super::codec::Error;
8use super::fields::{walk_fields, FieldReader, FieldWriter, WIRE_BYTES, WIRE_U64};
9
10const F_HELLO_NAME: u16 = 1;
13const F_HELLO_VERSION: u16 = 2;
14const F_HELLO_CAPS: u16 = 3;
15const F_HELLO_PROTOCOL: u16 = 4;
16
17const F_ACK_PROTOCOL: u16 = 1;
18const F_ACK_PHI_VERSION: u16 = 2;
19const F_ACK_CWD: u16 = 3;
20const F_ACK_SESSION_ID: u16 = 4;
21const F_ACK_EXT_DIR: u16 = 5;
22
23const F_REG_CMD_NAME: u16 = 1;
24const F_REG_CMD_DESC: u16 = 2;
25
26const F_REG_TOOL_NAME: u16 = 1;
27const F_REG_TOOL_DESC: u16 = 2;
28const F_REG_TOOL_SCHEMA: u16 = 3;
29const F_REG_TOOL_TIMEOUT_SEC: u16 = 4;
30const F_REG_TOOL_HAS_DETAIL: u16 = 5;
31
32const F_SUB_EVENTS: u16 = 1;
33const F_SUB_INTERCEPT: u16 = 2;
34
35const F_TOOL_DETAIL_RESULT: u16 = 1;
36
37const F_CMD_INV_NAME: u16 = 1;
38const F_CMD_INV_ARGS: u16 = 2;
39
40const F_CMD_RES_OK: u16 = 1;
41const F_CMD_RES_ERROR: u16 = 2;
42const F_CMD_RES_NOTIFY: u16 = 3;
43const F_CMD_RES_SUBMIT: u16 = 4;
44
45const F_TOOL_INV_NAME: u16 = 1;
46const F_TOOL_INV_ARGS: u16 = 2;
47
48const F_TOOL_RES_CONTENT: u16 = 1;
49const F_TOOL_RES_DETAIL: u16 = 2;
50const F_TOOL_RES_OUTPUT: u16 = 3;
51const F_TOOL_RES_IS_ERROR: u16 = 4;
52const F_TOOL_RES_ERROR: u16 = 5;
53
54const F_IX_REQ_EVENT: u16 = 1;
55const F_IX_REQ_TOOL_NAME: u16 = 2;
56const F_IX_REQ_TOOL_CALL_ID: u16 = 3;
57const F_IX_REQ_INPUT: u16 = 4;
58const F_IX_REQ_CONTENT: u16 = 5;
59const F_IX_REQ_IS_ERROR: u16 = 6;
60const F_IX_REQ_ERR_TEXT: u16 = 7;
61const F_IX_REQ_PROMPT: u16 = 8;
62const F_IX_REQ_REASON: u16 = 9;
63const F_IX_REQ_TARGET_ID: u16 = 10;
64const F_IX_REQ_TURN_INDEX: u16 = 11;
65
66const F_IX_RES_BLOCK: u16 = 1;
67const F_IX_RES_STOP: u16 = 2;
68const F_IX_RES_CANCEL: u16 = 3;
69const F_IX_RES_REASON: u16 = 4;
70const F_IX_RES_INPUT: u16 = 5;
71const F_IX_RES_CONTENT: u16 = 6;
72const F_IX_RES_CONTEXT: u16 = 7;
73const F_IX_RES_SYS_APPEND: u16 = 8;
74const F_IX_RES_TOAST: u16 = 9;
75const F_IX_RES_HANDLED: u16 = 10;
76const F_IX_RES_PROMPT: u16 = 11;
77const F_IX_RES_CONTINUE: u16 = 12;
78
79const F_EV_EVENT: u16 = 1;
80const F_EV_TOOL_NAME: u16 = 2;
81const F_EV_TOOL_CALL_ID: u16 = 3;
82const F_EV_INPUT: u16 = 4;
83const F_EV_IS_ERROR: u16 = 5;
84const F_EV_PROMPT: u16 = 6;
85const F_EV_REASON: u16 = 7;
86const F_EV_TURN_INDEX: u16 = 8;
87const F_EV_SESSION_ID: u16 = 9;
88const F_EV_PREVIOUS_SESSION_ID: u16 = 10;
89const F_EV_TARGET_SESSION_ID: u16 = 11;
90
91const F_NOTIFY_LEVEL: u16 = 1;
92const F_NOTIFY_MESSAGE: u16 = 2;
93const F_NOTIFY_STATUS: u16 = 3;
94const F_NOTIFY_STATUS_SET: u16 = 4;
95
96const F_HOST_REQ_METHOD: u16 = 1;
97const F_HOST_REQ_ARG: u16 = 2;
98
99const F_HOST_RES_OK: u16 = 1;
100const F_HOST_RES_ERROR: u16 = 2;
101const F_HOST_RES_BODY: u16 = 3;
102
103const F_META_SESSION_ID: u16 = 1;
104const F_META_CWD: u16 = 2;
105
106fn take_u64(kind: u8, fr: &mut FieldReader<'_>) -> Result<u64, Error> {
107 if kind != WIRE_U64 {
108 fr.skip(kind)?;
109 return Err(Error::BadWire);
110 }
111 fr.u64()
112}
113
114fn take_bytes<'a>(kind: u8, fr: &mut FieldReader<'a>) -> Result<&'a [u8], Error> {
115 if kind != WIRE_BYTES {
116 fr.skip(kind)?;
117 return Err(Error::BadWire);
118 }
119 fr.bytes()
120}
121
122fn take_string(kind: u8, fr: &mut FieldReader<'_>) -> Result<String, Error> {
125 Ok(String::from_utf8_lossy(take_bytes(kind, fr)?).into_owned())
126}
127
128#[derive(Debug, Clone, PartialEq, Eq, Default)]
130pub struct Hello {
131 pub name: String,
132 pub version: String,
133 pub caps: u32,
134 pub protocol: u16,
135}
136
137pub fn encode_hello(h: &Hello) -> Vec<u8> {
138 let mut fw = FieldWriter::new();
139 fw.put_string(F_HELLO_NAME, &h.name);
140 fw.put_string(F_HELLO_VERSION, &h.version);
141 fw.put_u32(F_HELLO_CAPS, h.caps);
142 fw.put_u16(F_HELLO_PROTOCOL, h.protocol);
143 fw.into_vec()
144}
145
146pub fn decode_hello(b: &[u8]) -> Result<Hello, Error> {
147 let mut h = Hello::default();
148 walk_fields(b, |tag, kind, fr| {
149 match tag {
150 F_HELLO_NAME => h.name = take_string(kind, fr)?,
151 F_HELLO_VERSION => h.version = take_string(kind, fr)?,
152 F_HELLO_CAPS => h.caps = take_u64(kind, fr)? as u32,
153 F_HELLO_PROTOCOL => h.protocol = take_u64(kind, fr)? as u16,
154 _ => fr.skip(kind)?,
155 }
156 Ok(())
157 })?;
158 Ok(h)
159}
160
161#[derive(Debug, Clone, PartialEq, Eq, Default)]
163pub struct HelloAck {
164 pub protocol: u16,
165 pub phi_version: String,
166 pub cwd: String,
167 pub session_id: String,
168 pub extension_dir: String,
169}
170
171pub fn encode_hello_ack(h: &HelloAck) -> Vec<u8> {
172 let mut fw = FieldWriter::new();
173 fw.put_u16(F_ACK_PROTOCOL, h.protocol);
174 fw.put_string(F_ACK_PHI_VERSION, &h.phi_version);
175 fw.put_string(F_ACK_CWD, &h.cwd);
176 fw.put_string(F_ACK_SESSION_ID, &h.session_id);
177 fw.put_string(F_ACK_EXT_DIR, &h.extension_dir);
178 fw.into_vec()
179}
180
181pub fn decode_hello_ack(b: &[u8]) -> Result<HelloAck, Error> {
182 let mut h = HelloAck::default();
183 walk_fields(b, |tag, kind, fr| {
184 match tag {
185 F_ACK_PROTOCOL => h.protocol = take_u64(kind, fr)? as u16,
186 F_ACK_PHI_VERSION => h.phi_version = take_string(kind, fr)?,
187 F_ACK_CWD => h.cwd = take_string(kind, fr)?,
188 F_ACK_SESSION_ID => h.session_id = take_string(kind, fr)?,
189 F_ACK_EXT_DIR => h.extension_dir = take_string(kind, fr)?,
190 _ => fr.skip(kind)?,
191 }
192 Ok(())
193 })?;
194 Ok(h)
195}
196
197#[derive(Debug, Clone, PartialEq, Eq, Default)]
199pub struct RegisterCommand {
200 pub name: String,
201 pub description: String,
202}
203
204pub fn encode_register_command(r: &RegisterCommand) -> Vec<u8> {
205 let mut fw = FieldWriter::new();
206 fw.put_string(F_REG_CMD_NAME, &r.name);
207 fw.put_string(F_REG_CMD_DESC, &r.description);
208 fw.into_vec()
209}
210
211pub fn decode_register_command(b: &[u8]) -> Result<RegisterCommand, Error> {
212 let mut r = RegisterCommand::default();
213 walk_fields(b, |tag, kind, fr| {
214 match tag {
215 F_REG_CMD_NAME => r.name = take_string(kind, fr)?,
216 F_REG_CMD_DESC => r.description = take_string(kind, fr)?,
217 _ => fr.skip(kind)?,
218 }
219 Ok(())
220 })?;
221 Ok(r)
222}
223
224#[derive(Debug, Clone, PartialEq, Eq, Default)]
226pub struct RegisterTool {
227 pub name: String,
228 pub description: String,
229 pub schema_json: Vec<u8>,
230 pub timeout_sec: u32,
233 pub has_detail: bool,
236}
237
238pub fn encode_register_tool(r: &RegisterTool) -> Vec<u8> {
239 let mut fw = FieldWriter::new();
240 fw.put_string(F_REG_TOOL_NAME, &r.name);
241 fw.put_string(F_REG_TOOL_DESC, &r.description);
242 fw.put_bytes(F_REG_TOOL_SCHEMA, &r.schema_json);
243 if r.timeout_sec > 0 {
244 fw.put_u32(F_REG_TOOL_TIMEOUT_SEC, r.timeout_sec);
245 }
246 if r.has_detail {
247 fw.put_bool(F_REG_TOOL_HAS_DETAIL, true);
248 }
249 fw.into_vec()
250}
251
252pub fn decode_register_tool(b: &[u8]) -> Result<RegisterTool, Error> {
253 let mut r = RegisterTool::default();
254 walk_fields(b, |tag, kind, fr| {
255 match tag {
256 F_REG_TOOL_NAME => r.name = take_string(kind, fr)?,
257 F_REG_TOOL_DESC => r.description = take_string(kind, fr)?,
258 F_REG_TOOL_SCHEMA => r.schema_json = take_bytes(kind, fr)?.to_vec(),
259 F_REG_TOOL_TIMEOUT_SEC => r.timeout_sec = take_u64(kind, fr)? as u32,
260 F_REG_TOOL_HAS_DETAIL => r.has_detail = take_u64(kind, fr)? != 0,
261 _ => fr.skip(kind)?,
262 }
263 Ok(())
264 })?;
265 Ok(r)
266}
267
268#[derive(Debug, Clone, PartialEq, Eq, Default)]
270pub struct ToolDetailResult {
271 pub detail: String,
272}
273
274pub fn encode_tool_detail_result(r: &ToolDetailResult) -> Vec<u8> {
275 let mut fw = FieldWriter::new();
276 fw.put_string(F_TOOL_DETAIL_RESULT, &r.detail);
277 fw.into_vec()
278}
279
280pub fn decode_tool_detail_result(b: &[u8]) -> Result<ToolDetailResult, Error> {
281 let mut r = ToolDetailResult::default();
282 walk_fields(b, |tag, kind, fr| {
283 match tag {
284 F_TOOL_DETAIL_RESULT => r.detail = take_string(kind, fr)?,
285 _ => fr.skip(kind)?,
286 }
287 Ok(())
288 })?;
289 Ok(r)
290}
291
292#[derive(Debug, Clone, PartialEq, Eq, Default)]
294pub struct Subscribe {
295 pub events: Vec<u16>,
296 pub intercept: Vec<u16>,
297}
298
299pub fn encode_subscribe(s: &Subscribe) -> Vec<u8> {
300 let mut fw = FieldWriter::new();
301 fw.put_u16s(F_SUB_EVENTS, &s.events);
302 fw.put_u16s(F_SUB_INTERCEPT, &s.intercept);
303 fw.into_vec()
304}
305
306pub fn decode_subscribe(b: &[u8]) -> Result<Subscribe, Error> {
307 let mut s = Subscribe::default();
308 walk_fields(b, |tag, kind, fr| {
309 match tag {
310 F_SUB_EVENTS => s.events = decode_u16s(take_bytes(kind, fr)?)?,
311 F_SUB_INTERCEPT => s.intercept = decode_u16s(take_bytes(kind, fr)?)?,
312 _ => fr.skip(kind)?,
313 }
314 Ok(())
315 })?;
316 Ok(s)
317}
318
319fn decode_u16s(p: &[u8]) -> Result<Vec<u16>, Error> {
320 if p.len() < 2 {
322 return Err(Error::Truncated);
323 }
324 let n = u16::from_le_bytes([p[0], p[1]]) as usize;
325 if p.len() < 2 + n * 2 {
326 return Err(Error::Truncated);
327 }
328 let mut out = Vec::with_capacity(n);
329 for i in 0..n {
330 let off = 2 + i * 2;
331 out.push(u16::from_le_bytes([p[off], p[off + 1]]));
332 }
333 Ok(out)
334}
335
336#[derive(Debug, Clone, PartialEq, Eq, Default)]
338pub struct CommandInvoked {
339 pub name: String,
340 pub args: String,
341}
342
343pub fn encode_command_invoked(c: &CommandInvoked) -> Vec<u8> {
344 let mut fw = FieldWriter::new();
345 fw.put_string(F_CMD_INV_NAME, &c.name);
346 fw.put_string(F_CMD_INV_ARGS, &c.args);
347 fw.into_vec()
348}
349
350pub fn decode_command_invoked(b: &[u8]) -> Result<CommandInvoked, Error> {
351 let mut c = CommandInvoked::default();
352 walk_fields(b, |tag, kind, fr| {
353 match tag {
354 F_CMD_INV_NAME => c.name = take_string(kind, fr)?,
355 F_CMD_INV_ARGS => c.args = take_string(kind, fr)?,
356 _ => fr.skip(kind)?,
357 }
358 Ok(())
359 })?;
360 Ok(c)
361}
362
363#[derive(Debug, Clone, PartialEq, Eq, Default)]
365pub struct CommandResponse {
366 pub ok: bool,
367 pub error: String,
368 pub notify: String,
369 pub submit: String,
370}
371
372pub fn encode_command_response(c: &CommandResponse) -> Vec<u8> {
373 let mut fw = FieldWriter::new();
374 fw.put_bool(F_CMD_RES_OK, c.ok);
375 fw.put_string(F_CMD_RES_ERROR, &c.error);
376 fw.put_string(F_CMD_RES_NOTIFY, &c.notify);
377 fw.put_string(F_CMD_RES_SUBMIT, &c.submit);
378 fw.into_vec()
379}
380
381pub fn decode_command_response(b: &[u8]) -> Result<CommandResponse, Error> {
382 let mut c = CommandResponse::default();
383 walk_fields(b, |tag, kind, fr| {
384 match tag {
385 F_CMD_RES_OK => c.ok = take_u64(kind, fr)? != 0,
386 F_CMD_RES_ERROR => c.error = take_string(kind, fr)?,
387 F_CMD_RES_NOTIFY => c.notify = take_string(kind, fr)?,
388 F_CMD_RES_SUBMIT => c.submit = take_string(kind, fr)?,
389 _ => fr.skip(kind)?,
390 }
391 Ok(())
392 })?;
393 Ok(c)
394}
395
396#[derive(Debug, Clone, PartialEq, Eq, Default)]
398pub struct ToolInvoke {
399 pub name: String,
400 pub args: Vec<u8>,
401}
402
403pub fn encode_tool_invoke(t: &ToolInvoke) -> Vec<u8> {
404 let mut fw = FieldWriter::new();
405 fw.put_string(F_TOOL_INV_NAME, &t.name);
406 fw.put_bytes(F_TOOL_INV_ARGS, &t.args);
407 fw.into_vec()
408}
409
410pub fn decode_tool_invoke(b: &[u8]) -> Result<ToolInvoke, Error> {
411 let mut t = ToolInvoke::default();
412 walk_fields(b, |tag, kind, fr| {
413 match tag {
414 F_TOOL_INV_NAME => t.name = take_string(kind, fr)?,
415 F_TOOL_INV_ARGS => t.args = take_bytes(kind, fr)?.to_vec(),
416 _ => fr.skip(kind)?,
417 }
418 Ok(())
419 })?;
420 Ok(t)
421}
422
423#[derive(Debug, Clone, PartialEq, Eq, Default)]
425pub struct ToolResultMsg {
426 pub content: String,
427 pub detail: String,
428 pub output: String,
429 pub is_error: bool,
430 pub error: String,
431}
432
433pub fn encode_tool_result(t: &ToolResultMsg) -> Vec<u8> {
434 let mut fw = FieldWriter::new();
435 fw.put_string(F_TOOL_RES_CONTENT, &t.content);
436 fw.put_string(F_TOOL_RES_DETAIL, &t.detail);
437 fw.put_string(F_TOOL_RES_OUTPUT, &t.output);
438 fw.put_bool(F_TOOL_RES_IS_ERROR, t.is_error);
439 fw.put_string(F_TOOL_RES_ERROR, &t.error);
440 fw.into_vec()
441}
442
443pub fn decode_tool_result(b: &[u8]) -> Result<ToolResultMsg, Error> {
444 let mut t = ToolResultMsg::default();
445 walk_fields(b, |tag, kind, fr| {
446 match tag {
447 F_TOOL_RES_CONTENT => t.content = take_string(kind, fr)?,
448 F_TOOL_RES_DETAIL => t.detail = take_string(kind, fr)?,
449 F_TOOL_RES_OUTPUT => t.output = take_string(kind, fr)?,
450 F_TOOL_RES_IS_ERROR => t.is_error = take_u64(kind, fr)? != 0,
451 F_TOOL_RES_ERROR => t.error = take_string(kind, fr)?,
452 _ => fr.skip(kind)?,
453 }
454 Ok(())
455 })?;
456 Ok(t)
457}
458
459#[derive(Debug, Clone, PartialEq, Eq, Default)]
461pub struct InterceptReq {
462 pub event: u16,
463 pub tool_name: String,
464 pub tool_call_id: String,
465 pub input: Vec<u8>,
466 pub content: String,
467 pub is_error: bool,
468 pub err_text: String,
469 pub prompt: String,
470 pub reason: String,
471 pub target_id: String,
472 pub turn_index: u32,
473}
474
475pub fn encode_intercept_req(r: &InterceptReq) -> Vec<u8> {
476 let mut fw = FieldWriter::new();
477 fw.put_u16(F_IX_REQ_EVENT, r.event);
478 fw.put_string(F_IX_REQ_TOOL_NAME, &r.tool_name);
479 fw.put_string(F_IX_REQ_TOOL_CALL_ID, &r.tool_call_id);
480 fw.put_bytes(F_IX_REQ_INPUT, &r.input);
481 fw.put_string(F_IX_REQ_CONTENT, &r.content);
482 fw.put_bool(F_IX_REQ_IS_ERROR, r.is_error);
483 fw.put_string(F_IX_REQ_ERR_TEXT, &r.err_text);
484 fw.put_string(F_IX_REQ_PROMPT, &r.prompt);
485 fw.put_string(F_IX_REQ_REASON, &r.reason);
486 fw.put_string(F_IX_REQ_TARGET_ID, &r.target_id);
487 fw.put_u32(F_IX_REQ_TURN_INDEX, r.turn_index);
488 fw.into_vec()
489}
490
491pub fn decode_intercept_req(b: &[u8]) -> Result<InterceptReq, Error> {
492 let mut r = InterceptReq::default();
493 walk_fields(b, |tag, kind, fr| {
494 match tag {
495 F_IX_REQ_EVENT => r.event = take_u64(kind, fr)? as u16,
496 F_IX_REQ_TOOL_NAME => r.tool_name = take_string(kind, fr)?,
497 F_IX_REQ_TOOL_CALL_ID => r.tool_call_id = take_string(kind, fr)?,
498 F_IX_REQ_INPUT => r.input = take_bytes(kind, fr)?.to_vec(),
499 F_IX_REQ_CONTENT => r.content = take_string(kind, fr)?,
500 F_IX_REQ_IS_ERROR => r.is_error = take_u64(kind, fr)? != 0,
501 F_IX_REQ_ERR_TEXT => r.err_text = take_string(kind, fr)?,
502 F_IX_REQ_PROMPT => r.prompt = take_string(kind, fr)?,
503 F_IX_REQ_REASON => r.reason = take_string(kind, fr)?,
504 F_IX_REQ_TARGET_ID => r.target_id = take_string(kind, fr)?,
505 F_IX_REQ_TURN_INDEX => r.turn_index = take_u64(kind, fr)? as u32,
506 _ => fr.skip(kind)?,
507 }
508 Ok(())
509 })?;
510 Ok(r)
511}
512
513#[derive(Debug, Clone, PartialEq, Eq, Default)]
515pub struct InterceptResp {
516 pub block: bool,
517 pub stop: bool,
518 pub cancel: bool,
519 pub handled: bool,
520 pub continue_: bool,
523 pub reason: String,
524 pub input: Vec<u8>,
525 pub content: String,
526 pub context: String,
527 pub system_prompt_append: String,
528 pub toast: String,
529 pub prompt: String,
530}
531
532pub fn encode_intercept_resp(r: &InterceptResp) -> Vec<u8> {
533 let mut fw = FieldWriter::new();
534 fw.put_bool(F_IX_RES_BLOCK, r.block);
535 fw.put_bool(F_IX_RES_STOP, r.stop);
536 fw.put_bool(F_IX_RES_CANCEL, r.cancel);
537 fw.put_string(F_IX_RES_REASON, &r.reason);
538 fw.put_bytes(F_IX_RES_INPUT, &r.input);
539 fw.put_string(F_IX_RES_CONTENT, &r.content);
540 fw.put_string(F_IX_RES_CONTEXT, &r.context);
541 fw.put_string(F_IX_RES_SYS_APPEND, &r.system_prompt_append);
542 fw.put_string(F_IX_RES_TOAST, &r.toast);
543 fw.put_bool(F_IX_RES_HANDLED, r.handled);
544 fw.put_string(F_IX_RES_PROMPT, &r.prompt);
545 fw.put_bool(F_IX_RES_CONTINUE, r.continue_);
546 fw.into_vec()
547}
548
549pub fn decode_intercept_resp(b: &[u8]) -> Result<InterceptResp, Error> {
550 let mut r = InterceptResp::default();
551 walk_fields(b, |tag, kind, fr| {
552 match tag {
553 F_IX_RES_BLOCK => r.block = take_u64(kind, fr)? != 0,
554 F_IX_RES_STOP => r.stop = take_u64(kind, fr)? != 0,
555 F_IX_RES_CANCEL => r.cancel = take_u64(kind, fr)? != 0,
556 F_IX_RES_HANDLED => r.handled = take_u64(kind, fr)? != 0,
557 F_IX_RES_CONTINUE => r.continue_ = take_u64(kind, fr)? != 0,
558 F_IX_RES_REASON => r.reason = take_string(kind, fr)?,
559 F_IX_RES_INPUT => r.input = take_bytes(kind, fr)?.to_vec(),
560 F_IX_RES_CONTENT => r.content = take_string(kind, fr)?,
561 F_IX_RES_CONTEXT => r.context = take_string(kind, fr)?,
562 F_IX_RES_SYS_APPEND => r.system_prompt_append = take_string(kind, fr)?,
563 F_IX_RES_TOAST => r.toast = take_string(kind, fr)?,
564 F_IX_RES_PROMPT => r.prompt = take_string(kind, fr)?,
565 _ => fr.skip(kind)?,
566 }
567 Ok(())
568 })?;
569 Ok(r)
570}
571
572#[derive(Debug, Clone, PartialEq, Eq, Default)]
574pub struct EventNotify {
575 pub event: u16,
576 pub tool_name: String,
577 pub tool_call_id: String,
578 pub input: Vec<u8>,
579 pub is_error: bool,
580 pub prompt: String,
581 pub reason: String,
582 pub turn_index: u32,
583 pub session_id: String,
584 pub previous_session_id: String,
585 pub target_session_id: String,
586}
587
588pub fn encode_event_notify(e: &EventNotify) -> Vec<u8> {
589 let mut fw = FieldWriter::new();
590 fw.put_u16(F_EV_EVENT, e.event);
591 fw.put_string(F_EV_TOOL_NAME, &e.tool_name);
592 fw.put_string(F_EV_TOOL_CALL_ID, &e.tool_call_id);
593 fw.put_bytes(F_EV_INPUT, &e.input);
594 fw.put_bool(F_EV_IS_ERROR, e.is_error);
595 fw.put_string(F_EV_PROMPT, &e.prompt);
596 fw.put_string(F_EV_REASON, &e.reason);
597 fw.put_u32(F_EV_TURN_INDEX, e.turn_index);
598 fw.put_string(F_EV_SESSION_ID, &e.session_id);
599 fw.put_string(F_EV_PREVIOUS_SESSION_ID, &e.previous_session_id);
600 fw.put_string(F_EV_TARGET_SESSION_ID, &e.target_session_id);
601 fw.into_vec()
602}
603
604pub fn decode_event_notify(b: &[u8]) -> Result<EventNotify, Error> {
605 let mut e = EventNotify::default();
606 walk_fields(b, |tag, kind, fr| {
607 match tag {
608 F_EV_EVENT => e.event = take_u64(kind, fr)? as u16,
609 F_EV_TOOL_NAME => e.tool_name = take_string(kind, fr)?,
610 F_EV_TOOL_CALL_ID => e.tool_call_id = take_string(kind, fr)?,
611 F_EV_INPUT => e.input = take_bytes(kind, fr)?.to_vec(),
612 F_EV_IS_ERROR => e.is_error = take_u64(kind, fr)? != 0,
613 F_EV_PROMPT => e.prompt = take_string(kind, fr)?,
614 F_EV_REASON => e.reason = take_string(kind, fr)?,
615 F_EV_TURN_INDEX => e.turn_index = take_u64(kind, fr)? as u32,
616 F_EV_SESSION_ID => e.session_id = take_string(kind, fr)?,
617 F_EV_PREVIOUS_SESSION_ID => e.previous_session_id = take_string(kind, fr)?,
618 F_EV_TARGET_SESSION_ID => e.target_session_id = take_string(kind, fr)?,
619 _ => fr.skip(kind)?,
620 }
621 Ok(())
622 })?;
623 Ok(e)
624}
625
626#[derive(Debug, Clone, PartialEq, Eq, Default)]
628pub struct NotifyMsg {
629 pub level: String,
630 pub message: String,
631 pub status: String,
632 pub status_set: bool,
633}
634
635pub fn encode_notify(n: &NotifyMsg) -> Vec<u8> {
636 let mut fw = FieldWriter::new();
637 fw.put_string(F_NOTIFY_LEVEL, &n.level);
638 fw.put_string(F_NOTIFY_MESSAGE, &n.message);
639 fw.put_string(F_NOTIFY_STATUS, &n.status);
640 fw.put_bool(F_NOTIFY_STATUS_SET, n.status_set);
641 fw.into_vec()
642}
643
644pub fn decode_notify(b: &[u8]) -> Result<NotifyMsg, Error> {
645 let mut n = NotifyMsg::default();
646 walk_fields(b, |tag, kind, fr| {
647 match tag {
648 F_NOTIFY_LEVEL => n.level = take_string(kind, fr)?,
649 F_NOTIFY_MESSAGE => n.message = take_string(kind, fr)?,
650 F_NOTIFY_STATUS => n.status = take_string(kind, fr)?,
651 F_NOTIFY_STATUS_SET => n.status_set = take_u64(kind, fr)? != 0,
652 _ => fr.skip(kind)?,
653 }
654 Ok(())
655 })?;
656 Ok(n)
657}
658
659#[derive(Debug, Clone, PartialEq, Eq, Default)]
661pub struct HostRequest {
662 pub method: String, pub arg: String,
664}
665
666pub fn encode_host_request(r: &HostRequest) -> Vec<u8> {
667 let mut fw = FieldWriter::new();
668 fw.put_string(F_HOST_REQ_METHOD, &r.method);
669 fw.put_string(F_HOST_REQ_ARG, &r.arg);
670 fw.into_vec()
671}
672
673pub fn decode_host_request(b: &[u8]) -> Result<HostRequest, Error> {
674 let mut r = HostRequest::default();
675 walk_fields(b, |tag, kind, fr| {
676 match tag {
677 F_HOST_REQ_METHOD => r.method = take_string(kind, fr)?,
678 F_HOST_REQ_ARG => r.arg = take_string(kind, fr)?,
679 _ => fr.skip(kind)?,
680 }
681 Ok(())
682 })?;
683 Ok(r)
684}
685
686#[derive(Debug, Clone, PartialEq, Eq, Default)]
688pub struct HostResult {
689 pub ok: bool,
690 pub error: String,
691 pub body: String,
692}
693
694pub fn encode_host_result(r: &HostResult) -> Vec<u8> {
695 let mut fw = FieldWriter::new();
696 fw.put_bool(F_HOST_RES_OK, r.ok);
697 fw.put_string(F_HOST_RES_ERROR, &r.error);
698 fw.put_string(F_HOST_RES_BODY, &r.body);
699 fw.into_vec()
700}
701
702pub fn decode_host_result(b: &[u8]) -> Result<HostResult, Error> {
703 let mut r = HostResult::default();
704 walk_fields(b, |tag, kind, fr| {
705 match tag {
706 F_HOST_RES_OK => r.ok = take_u64(kind, fr)? != 0,
707 F_HOST_RES_ERROR => r.error = take_string(kind, fr)?,
708 F_HOST_RES_BODY => r.body = take_string(kind, fr)?,
709 _ => fr.skip(kind)?,
710 }
711 Ok(())
712 })?;
713 Ok(r)
714}
715
716#[derive(Debug, Clone, PartialEq, Eq, Default)]
718pub struct SessionMeta {
719 pub session_id: String,
720 pub cwd: String,
721}
722
723pub fn encode_session_meta(m: &SessionMeta) -> Vec<u8> {
724 let mut fw = FieldWriter::new();
725 fw.put_string(F_META_SESSION_ID, &m.session_id);
726 fw.put_string(F_META_CWD, &m.cwd);
727 fw.into_vec()
728}
729
730pub fn decode_session_meta(b: &[u8]) -> Result<SessionMeta, Error> {
731 let mut m = SessionMeta::default();
732 walk_fields(b, |tag, kind, fr| {
733 match tag {
734 F_META_SESSION_ID => m.session_id = take_string(kind, fr)?,
735 F_META_CWD => m.cwd = take_string(kind, fr)?,
736 _ => fr.skip(kind)?,
737 }
738 Ok(())
739 })?;
740 Ok(m)
741}
742
743#[cfg(test)]
744mod tests {
745 use super::*;
746
747 #[test]
750 fn all_messages_roundtrip() {
751 type Reencode = fn(&[u8]) -> Result<Vec<u8>, Error>;
752 let cases: Vec<(Vec<u8>, Reencode)> = vec![
753 (
754 encode_hello(&Hello {
755 name: "greet".into(),
756 version: "1.0.0".into(),
757 caps: 3,
758 protocol: 1,
759 }),
760 |b| Ok(encode_hello(&decode_hello(b)?)),
761 ),
762 (
763 encode_hello_ack(&HelloAck {
764 protocol: 1,
765 phi_version: "v0.19.0".into(),
766 cwd: "/tmp".into(),
767 session_id: "s1".into(),
768 extension_dir: "/ext".into(),
769 }),
770 |b| Ok(encode_hello_ack(&decode_hello_ack(b)?)),
771 ),
772 (
773 encode_register_command(&RegisterCommand {
774 name: "hi".into(),
775 description: "Say hi".into(),
776 }),
777 |b| Ok(encode_register_command(&decode_register_command(b)?)),
778 ),
779 (
780 encode_register_tool(&RegisterTool {
781 name: "t".into(),
782 description: "d".into(),
783 schema_json: br#"{"type":"object"}"#.to_vec(),
784 timeout_sec: 120,
785 has_detail: true,
786 }),
787 |b| Ok(encode_register_tool(&decode_register_tool(b)?)),
788 ),
789 (
790 encode_tool_detail_result(&ToolDetailResult {
791 detail: "path/to/file".into(),
792 }),
793 |b| Ok(encode_tool_detail_result(&decode_tool_detail_result(b)?)),
794 ),
795 (
796 encode_subscribe(&Subscribe {
797 events: vec![5, 10],
798 intercept: vec![1, 2],
799 }),
800 |b| Ok(encode_subscribe(&decode_subscribe(b)?)),
801 ),
802 (
803 encode_command_invoked(&CommandInvoked {
804 name: "hi".into(),
805 args: "a b".into(),
806 }),
807 |b| Ok(encode_command_invoked(&decode_command_invoked(b)?)),
808 ),
809 (
810 encode_command_response(&CommandResponse {
811 ok: false,
812 error: "boom".into(),
813 notify: String::new(),
814 submit: "next".into(),
815 }),
816 |b| Ok(encode_command_response(&decode_command_response(b)?)),
817 ),
818 (
819 encode_tool_invoke(&ToolInvoke {
820 name: "t".into(),
821 args: br#"{"k":1}"#.to_vec(),
822 }),
823 |b| Ok(encode_tool_invoke(&decode_tool_invoke(b)?)),
824 ),
825 (
826 encode_tool_result(&ToolResultMsg {
827 content: "c".into(),
828 detail: "d".into(),
829 output: "o".into(),
830 is_error: true,
831 error: "e".into(),
832 }),
833 |b| Ok(encode_tool_result(&decode_tool_result(b)?)),
834 ),
835 (
836 encode_intercept_req(&InterceptReq {
837 event: 1,
838 tool_name: "bash".into(),
839 tool_call_id: "c1".into(),
840 input: br#"{"command":"ls"}"#.to_vec(),
841 content: "out".into(),
842 is_error: false,
843 err_text: String::new(),
844 prompt: "p".into(),
845 reason: "r".into(),
846 target_id: "t2".into(),
847 turn_index: 3,
848 }),
849 |b| Ok(encode_intercept_req(&decode_intercept_req(b)?)),
850 ),
851 (
852 encode_intercept_resp(&InterceptResp {
853 block: true,
854 stop: false,
855 cancel: false,
856 handled: true,
857 continue_: true,
858 reason: "r".into(),
859 input: b"in".to_vec(),
860 content: "c".into(),
861 context: "ctx".into(),
862 system_prompt_append: "sys".into(),
863 toast: "t".into(),
864 prompt: "p".into(),
865 }),
866 |b| Ok(encode_intercept_resp(&decode_intercept_resp(b)?)),
867 ),
868 (
869 encode_event_notify(&EventNotify {
870 event: 5,
871 tool_name: "t".into(),
872 tool_call_id: "c".into(),
873 input: b"i".to_vec(),
874 is_error: true,
875 prompt: "p".into(),
876 reason: "r".into(),
877 turn_index: 2,
878 session_id: "s".into(),
879 previous_session_id: "ps".into(),
880 target_session_id: "ts".into(),
881 }),
882 |b| Ok(encode_event_notify(&decode_event_notify(b)?)),
883 ),
884 (
885 encode_notify(&NotifyMsg {
886 level: "info".into(),
887 message: "Hello".into(),
888 status: "st".into(),
889 status_set: true,
890 }),
891 |b| Ok(encode_notify(&decode_notify(b)?)),
892 ),
893 (
894 encode_host_request(&HostRequest {
895 method: "confirm".into(),
896 arg: r#"{"Title":"t"}"#.into(),
897 }),
898 |b| Ok(encode_host_request(&decode_host_request(b)?)),
899 ),
900 (
901 encode_host_result(&HostResult {
902 ok: true,
903 error: String::new(),
904 body: "b".into(),
905 }),
906 |b| Ok(encode_host_result(&decode_host_result(b)?)),
907 ),
908 (
909 encode_session_meta(&SessionMeta {
910 session_id: "s2".into(),
911 cwd: "/x".into(),
912 }),
913 |b| Ok(encode_session_meta(&decode_session_meta(b)?)),
914 ),
915 ];
916
917 for (bytes, reencode) in cases {
918 assert_eq!(reencode(&bytes).unwrap(), bytes);
919 }
920 }
921
922 #[test]
923 fn decode_skips_unknown_tags() {
924 let mut w = FieldWriter::new();
925 w.put_string(1, "name");
926 w.put_string(200, "future field"); w.put_string(2, "1.0.0");
928 let h = decode_hello(&w.into_vec()).unwrap();
929 assert_eq!(h.name, "name");
930 assert_eq!(h.version, "1.0.0");
931 }
932}