qyle commited on
Commit
334d4f8
Β·
1 Parent(s): a4c6e98

Deploy from GitLab c3eac333

Browse files
agent/agent.py CHANGED
@@ -85,6 +85,7 @@ class Agent:
85
  conversation: ConversationHistory,
86
  *,
87
  documents: dict[str, str] | None = None,
 
88
  ) -> AgentResponse:
89
  conversation.record_user(query)
90
 
@@ -120,7 +121,7 @@ class Agent:
120
  # the whole skill pipeline (e.g. the wiki sub-agent).
121
  with timed_block(f"agent/tool/{tc.name}"):
122
  record, should_return = self._dispatch_tool_call(
123
- tc, conversation, documents
124
  )
125
  except Exception as exc:
126
  conversation.record_tool_exchange(
@@ -184,6 +185,7 @@ class Agent:
184
  tc: ToolCall,
185
  conversation: ConversationHistory,
186
  documents: dict[str, str] | None = None,
 
187
  ) -> tuple[ToolCallRecord, bool]:
188
  should_return = False
189
  try:
@@ -192,7 +194,10 @@ class Agent:
192
  result = f"Instructions: {instructions}"
193
  elif tc.name == "execute_function":
194
  result, should_return = self.skills.execute(
195
- **tc.arguments, conversation=conversation, documents=documents
 
 
 
196
  )
197
  else:
198
  result = f"Error: Unknown function: {tc.name}"
 
85
  conversation: ConversationHistory,
86
  *,
87
  documents: dict[str, str] | None = None,
88
+ request_id: str | None = None,
89
  ) -> AgentResponse:
90
  conversation.record_user(query)
91
 
 
121
  # the whole skill pipeline (e.g. the wiki sub-agent).
122
  with timed_block(f"agent/tool/{tc.name}"):
123
  record, should_return = self._dispatch_tool_call(
124
+ tc, conversation, documents, request_id
125
  )
126
  except Exception as exc:
127
  conversation.record_tool_exchange(
 
185
  tc: ToolCall,
186
  conversation: ConversationHistory,
187
  documents: dict[str, str] | None = None,
188
+ request_id: str | None = None,
189
  ) -> tuple[ToolCallRecord, bool]:
190
  should_return = False
191
  try:
 
194
  result = f"Instructions: {instructions}"
195
  elif tc.name == "execute_function":
196
  result, should_return = self.skills.execute(
197
+ **tc.arguments,
198
+ conversation=conversation,
199
+ documents=documents,
200
+ request_id=request_id,
201
  )
202
  else:
203
  result = f"Error: Unknown function: {tc.name}"
agent/agent_client.py CHANGED
@@ -101,7 +101,7 @@ class AgentClient:
101
  self._turn_start = len(self.conversation.ordered_transcript())
102
  self._documents = documents or {}
103
  try:
104
- response = self._invoke(query, lang=lang)
105
  except AgentException as e:
106
  outcome = ChatOutcome(
107
  reply="",
@@ -129,9 +129,11 @@ class AgentClient:
129
  self.logger_fn(outcome)
130
  return outcome
131
 
132
- def _invoke(self, query: str, lang: str | None = None) -> AgentResponse:
 
 
133
  assert self.agent is not None, "base _invoke requires self.agent"
134
- return self.agent.chat(query, self.conversation)
135
 
136
  def _inference_impacts(self, n_tokens: int) -> Any:
137
  """Return raw EcoLogits Impacts for this call, or None if unavailable."""
@@ -259,12 +261,16 @@ class AgentClient:
259
  class SkillsAgentClient(AgentClient):
260
  """gpt-oss-style client: post-call language-leak detection + one correction pass."""
261
 
262
- def _invoke(self, query: str, lang: str | None = None) -> AgentResponse:
 
 
263
  with timed_block("client/detect_language"):
264
  lang = lang or detect_language(query) or "en"
265
  docs = self._documents or None
266
  with timed_block("client/agent_chat"):
267
- response = self.agent.chat(query, self.conversation, documents=docs)
 
 
268
  with timed_block("client/find_leakage"):
269
  leaks = find_leakage(response.content, target_lang=lang)
270
  if not leaks:
@@ -273,9 +279,12 @@ class SkillsAgentClient(AgentClient):
273
  language=lang, mistranslated_words=leaks
274
  )
275
  # Full second pass through the agent (and potentially the whole wiki
276
- # pipeline again) β€” a prime latency suspect.
 
277
  with timed_block("client/translation_correction_chat"):
278
- return self.agent.chat(prompt, self.conversation, documents=docs)
 
 
279
 
280
  def _extract_context(self) -> list:
281
  return [
 
101
  self._turn_start = len(self.conversation.ordered_transcript())
102
  self._documents = documents or {}
103
  try:
104
+ response = self._invoke(query, lang=lang, reply_id=reply_id)
105
  except AgentException as e:
106
  outcome = ChatOutcome(
107
  reply="",
 
129
  self.logger_fn(outcome)
130
  return outcome
131
 
132
+ def _invoke(
133
+ self, query: str, lang: str | None = None, reply_id: str | None = None
134
+ ) -> AgentResponse:
135
  assert self.agent is not None, "base _invoke requires self.agent"
136
+ return self.agent.chat(query, self.conversation, request_id=reply_id)
137
 
138
  def _inference_impacts(self, n_tokens: int) -> Any:
139
  """Return raw EcoLogits Impacts for this call, or None if unavailable."""
 
261
  class SkillsAgentClient(AgentClient):
262
  """gpt-oss-style client: post-call language-leak detection + one correction pass."""
263
 
264
+ def _invoke(
265
+ self, query: str, lang: str | None = None, reply_id: str | None = None
266
+ ) -> AgentResponse:
267
  with timed_block("client/detect_language"):
268
  lang = lang or detect_language(query) or "en"
269
  docs = self._documents or None
270
  with timed_block("client/agent_chat"):
271
+ response = self.agent.chat(
272
+ query, self.conversation, documents=docs, request_id=reply_id
273
+ )
274
  with timed_block("client/find_leakage"):
275
  leaks = find_leakage(response.content, target_lang=lang)
276
  if not leaks:
 
279
  language=lang, mistranslated_words=leaks
280
  )
281
  # Full second pass through the agent (and potentially the whole wiki
282
+ # pipeline again) β€” a prime latency suspect. Same reply_id: it's still
283
+ # one turn from the caller's perspective.
284
  with timed_block("client/translation_correction_chat"):
285
+ return self.agent.chat(
286
+ prompt, self.conversation, documents=docs, request_id=reply_id
287
+ )
288
 
289
  def _extract_context(self) -> list:
290
  return [
agent/champ_client.py CHANGED
@@ -60,7 +60,12 @@ class ChampAgentClient(AgentClient):
60
  self._last_passages: list[str] = []
61
  self._last_triage: dict = {}
62
 
63
- def _invoke(self, query: str, lang: str | None = None) -> AgentResponse:
 
 
 
 
 
64
  if lang not in ("en", "fr"):
65
  lang = "en"
66
  self.conversation.record_user(query)
 
60
  self._last_passages: list[str] = []
61
  self._last_triage: dict = {}
62
 
63
+ def _invoke(
64
+ self, query: str, lang: str | None = None, reply_id: str | None = None
65
+ ) -> AgentResponse:
66
+ # reply_id unused: this client calls the provider directly, not
67
+ # through Agent/SkillsManager, so there's no draft/judge pipeline to
68
+ # tag with it. Accepted only to match AgentClient.call()'s signature.
69
  if lang not in ("en", "fr"):
70
  lang = "en"
71
  self.conversation.record_user(query)
agent/fake_client.py CHANGED
@@ -21,7 +21,10 @@ class FakeAgentClient(AgentClient):
21
  super().__init__(agent=None, logger_fn=logger_fn)
22
  self.provider = provider
23
 
24
- def _invoke(self, query: str, lang: str | None = None) -> AgentResponse:
 
 
 
25
  self.conversation.record_user(query)
26
  completion = self.provider.chat(
27
  messages=self.conversation.to_messages(),
 
21
  super().__init__(agent=None, logger_fn=logger_fn)
22
  self.provider = provider
23
 
24
+ def _invoke(
25
+ self, query: str, lang: str | None = None, reply_id: str | None = None
26
+ ) -> AgentResponse:
27
+ # reply_id unused β€” accepted only to match AgentClient.call()'s signature.
28
  self.conversation.record_user(query)
29
  completion = self.provider.chat(
30
  messages=self.conversation.to_messages(),
agent/frontier_client.py CHANGED
@@ -50,7 +50,12 @@ class _SingleShotClient(AgentClient):
50
  self.context_window = context_window_for(self.model_id)
51
  self._last_provider_impacts: Any = None
52
 
53
- def _invoke(self, query: str, lang: str | None = None) -> AgentResponse:
 
 
 
 
 
54
  if lang not in ("en", "fr"):
55
  lang = "en"
56
  self.conversation.record_user(query)
 
50
  self.context_window = context_window_for(self.model_id)
51
  self._last_provider_impacts: Any = None
52
 
53
+ def _invoke(
54
+ self, query: str, lang: str | None = None, reply_id: str | None = None
55
+ ) -> AgentResponse:
56
+ # reply_id unused: this client calls the provider directly, not
57
+ # through Agent/SkillsManager, so there's no draft/judge pipeline to
58
+ # tag with it. Accepted only to match AgentClient.call()'s signature.
59
  if lang not in ("en", "fr"):
60
  lang = "en"
61
  self.conversation.record_user(query)
agent/skill_decorators.py CHANGED
@@ -76,3 +76,31 @@ def needs_documents(func: F) -> F:
76
  """
77
  func._needs_documents = True # type: ignore[attr-defined]
78
  return func
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
76
  """
77
  func._needs_documents = True # type: ignore[attr-defined]
78
  return func
79
+
80
+
81
+ def needs_request_id(func: F) -> F:
82
+ """Mark a skill function as requiring the current turn's request id.
83
+
84
+ The dispatch path (`SkillsManager._import_and_call`) sees the
85
+ `_needs_request_id` attribute and injects a `request_id=` kwarg β€” the same
86
+ id `AgentClient.call()` generates once per turn as `reply_id`, threaded
87
+ down through `Agent.chat()`. Use it to tag log lines so a call can be told
88
+ apart from an unrelated concurrent or abandoned one sharing the same log
89
+ stream, instead of a skill minting its own id internally.
90
+ The function must declare a matching `request_id` parameter.
91
+
92
+ Usage:
93
+
94
+ from agent.skill_decorators import needs_conversation, needs_request_id
95
+
96
+ @needs_conversation
97
+ @needs_request_id
98
+ def my_skill_func(
99
+ user_query: str,
100
+ conversation: ConversationHistory,
101
+ request_id: str | None = None,
102
+ ) -> tuple[str, bool]:
103
+ ...
104
+ """
105
+ func._needs_request_id = True # type: ignore[attr-defined]
106
+ return func
agent/skill_manager.py CHANGED
@@ -222,6 +222,7 @@ class SkillsManager:
222
  conversation: ConversationHistory,
223
  params: dict = {},
224
  documents: dict | None = None,
 
225
  ) -> tuple[str, bool]:
226
  """Execute skill action by dynamically importing and calling Python functions."""
227
  if skill_name not in self.skills:
@@ -229,7 +230,7 @@ class SkillsManager:
229
 
230
  script_folder = self.skills[skill_name].path.parent / "scripts"
231
  result = self._import_and_call(
232
- script_folder, function_name, conversation, documents, **params
233
  )
234
  if result is not None:
235
  return result
@@ -245,6 +246,7 @@ class SkillsManager:
245
  action: str,
246
  conversation: ConversationHistory,
247
  documents: dict | None = None,
 
248
  **params,
249
  ) -> tuple[str, bool]:
250
  folder_str = str(folder)
@@ -283,6 +285,8 @@ class SkillsManager:
283
  params["conversation"] = conversation
284
  if getattr(func, "_needs_documents", False):
285
  params["documents"] = documents
 
 
286
 
287
  # Do NOT swallow exceptions raised inside the skill β€” surface them
288
  # so the caller sees what went wrong instead of the misleading
 
222
  conversation: ConversationHistory,
223
  params: dict = {},
224
  documents: dict | None = None,
225
+ request_id: str | None = None,
226
  ) -> tuple[str, bool]:
227
  """Execute skill action by dynamically importing and calling Python functions."""
228
  if skill_name not in self.skills:
 
230
 
231
  script_folder = self.skills[skill_name].path.parent / "scripts"
232
  result = self._import_and_call(
233
+ script_folder, function_name, conversation, documents, request_id, **params
234
  )
235
  if result is not None:
236
  return result
 
246
  action: str,
247
  conversation: ConversationHistory,
248
  documents: dict | None = None,
249
+ request_id: str | None = None,
250
  **params,
251
  ) -> tuple[str, bool]:
252
  folder_str = str(folder)
 
285
  params["conversation"] = conversation
286
  if getattr(func, "_needs_documents", False):
287
  params["documents"] = documents
288
+ if getattr(func, "_needs_request_id", False):
289
+ params["request_id"] = request_id
290
 
291
  # Do NOT swallow exceptions raised inside the skill β€” surface them
292
  # so the caller sees what went wrong instead of the misleading
agent/skills/pediatry_wiki/scripts/generate_wiki_response.py CHANGED
@@ -18,7 +18,7 @@ from agent.context_window import (
18
  )
19
  from agent.conversation_history import ConversationHistory
20
  from agent.documents import build_documents_block, render_documents
21
- from agent.skill_decorators import needs_conversation, needs_documents
22
  from agent.skills.pediatry_wiki.scripts.browse_category import (
23
  TOP_INDEX_PATH,
24
  browse_category,
@@ -320,6 +320,7 @@ def judge_answer(
320
  use_verbatim_check: bool = False,
321
  judge_prompt_override: str | None = None,
322
  echo_stream: bool = False,
 
323
  ) -> tuple[str, str, str | None, list | None]:
324
  """Return (verdict, reasoning, model_reasoning, claims). verdict is
325
  PASS / FAIL / ERROR.
@@ -349,9 +350,14 @@ def judge_answer(
349
  entirely (eval-only). Must take {question}/{reference}/{answer} and
350
  return the same verdict/reasoning JSON shape.
351
 
352
- `echo_stream=True` prints the judge's reasoning/content tokens to stdout
353
- live as they arrive (eval-only β€” off by default so production calls stay
354
- silent).
 
 
 
 
 
355
  """
356
  rendered_docs = render_documents(documents)
357
  documents_section = (
@@ -367,7 +373,7 @@ def judge_answer(
367
  documents_section=documents_section,
368
  )
369
 
370
- def _call():
371
  return _provider().chat(
372
  messages=[Message(role="user", content=prompt)],
373
  model_id=JUDGE_MODEL_ID,
@@ -378,13 +384,15 @@ def judge_answer(
378
  max_tokens=_JUDGE_MODEL_MAX_TOKENS or 65_536,
379
  stream=True, # avoids a gateway 504 on long reasoning traces
380
  echo_stream=echo_stream,
 
381
  )
382
 
383
  completion = None
384
  last_timeout = False
385
  for attempt in range(1, JUDGE_CALL_MAX_ATTEMPTS_ON_TIMEOUT + 1):
 
386
  try:
387
- completion = _judge_call_pool.submit(_call).result(
388
  timeout=JUDGE_CALL_TIMEOUT_SECONDS
389
  )
390
  last_timeout = False
@@ -392,8 +400,10 @@ def judge_answer(
392
  except FutureTimeoutError:
393
  last_timeout = True
394
  logger.warning(
395
- "judge_answer: call exceeded %ss, abandoning and retrying "
396
- "(attempt %d/%d)",
 
 
397
  JUDGE_CALL_TIMEOUT_SECONDS,
398
  attempt,
399
  JUDGE_CALL_MAX_ATTEMPTS_ON_TIMEOUT,
@@ -404,7 +414,7 @@ def judge_answer(
404
  # the 3900-token cap); if it still exhausts its own budget and raises,
405
  # don't let that crash the whole chat turn β€” degrade to the same
406
  # graceful ERROR verdict the unparseable-output path below returns.
407
- logger.warning("judge_answer: provider call failed: %s", e)
408
  return "ERROR", f"Judge provider call failed: {e}", None, None
409
  if last_timeout:
410
  # Deliberately not a graceful "ERROR" verdict: that would feed a
@@ -413,8 +423,9 @@ def judge_answer(
413
  # iteration for no signal, then likely timing out again on the same
414
  # slow request. Raise instead so the caller can fail once, locally.
415
  raise ProviderTimeoutExhaustedError(
416
- f"judge call timed out after {JUDGE_CALL_MAX_ATTEMPTS_ON_TIMEOUT} "
417
- f"attempts of {JUDGE_CALL_TIMEOUT_SECONDS}s each"
 
418
  )
419
  if usage_sink is not None:
420
  usage_sink["prompt_tokens"] = completion.usage.prompt_tokens
@@ -467,7 +478,8 @@ def judge_answer(
467
 
468
 
469
  def check_attribution(
470
- answer: str, *, usage_sink: dict | None = None, echo_stream: bool = False
 
471
  ) -> tuple[str, str, str | None]:
472
  """Return (verdict, reasoning, model_reasoning). FAIL means the reply
473
  attributes its answer to an external source instead of speaking directly.
@@ -478,8 +490,10 @@ def check_attribution(
478
  the thinking trace is worth persisting, not discarding.
479
 
480
  `usage_sink`, if given, is populated with this call's prompt_tokens /
481
- total_tokens / attempts (see judge_answer). `echo_stream` β€” see judge_answer."""
 
482
  prompt = ATTRIBUTION_PROMPT.format(answer=answer)
 
483
  try:
484
  completion = _provider().chat(
485
  messages=[Message(role="user", content=prompt)],
@@ -491,6 +505,7 @@ def check_attribution(
491
  max_tokens=_JUDGE_MODEL_MAX_TOKENS or 65_536,
492
  stream=True, # avoids a gateway 504 on long reasoning traces
493
  echo_stream=echo_stream,
 
494
  )
495
  except ProviderError as e:
496
  # Same reasoning as judge_answer's wrap: don't let a provider failure
@@ -498,7 +513,7 @@ def check_attribution(
498
  # caller treats any non-FAIL verdict as "ship the already-judge-passed
499
  # answer", so degrading to ERROR is a safe, bounded fallback instead of
500
  # crashing the whole generate_wiki_response call.
501
- logger.warning("check_attribution: provider call failed: %s", e)
502
  return "ERROR", f"Attribution provider call failed: {e}", None
503
  if usage_sink is not None:
504
  usage_sink["prompt_tokens"] = completion.usage.prompt_tokens
@@ -576,10 +591,12 @@ def _read_top_index() -> str:
576
 
577
  @needs_conversation
578
  @needs_documents
 
579
  def generate_wiki_response(
580
  user_query: str,
581
  conversation: ConversationHistory,
582
  documents: dict[str, str] | None = None,
 
583
  ) -> tuple[str, bool]:
584
  """Production entry point: short-answer draft prompt, reject ledger, and
585
  the split top-level category index β€” the configuration formerly split
@@ -597,6 +614,7 @@ def generate_wiki_response(
597
  # (or if) it ever finishes β€” this is prod, not eval, but nothing here
598
  # reaches the chat response itself.
599
  echo_stream=True,
 
600
  )
601
 
602
 
@@ -614,6 +632,7 @@ def generate_wiki_response_with_prompt(
614
  judge_prompt_override: str | None = None,
615
  run_attribution_check: bool = True,
616
  echo_stream: bool = False,
 
617
  ) -> tuple[str, bool]:
618
  """Restricted mini-agent: read wiki pages β†’ answer β†’ self-correct vs. the judge.
619
 
@@ -663,8 +682,17 @@ def generate_wiki_response_with_prompt(
663
  production wiki's category index).
664
 
665
  `echo_stream=True` streams every LLM call this makes (draft, judge,
666
- attribution) and prints tokens live to stdout as they arrive β€” for
667
- watching a long-running production call, not the chat response itself.
 
 
 
 
 
 
 
 
 
668
 
669
  Records every sub-step (tool call, tool result, draft, judge verdict,
670
  improve cycle) as `pipeline_step` events on the outer ConversationHistory.
@@ -685,6 +713,7 @@ def generate_wiki_response_with_prompt(
685
  "type": "pipeline_step",
686
  "step": "generate_wiki_response/start",
687
  "user_query": user_query,
 
688
  }
689
  )
690
 
@@ -726,8 +755,9 @@ def generate_wiki_response_with_prompt(
726
  # would let the model repeat a discarded mistake or answer
727
  # ungrounded instead of failing loudly, so bail out instead.
728
  logger.warning(
729
- "generate_wiki_response exceeded the context window for "
730
- "%s (iteration %d)",
 
731
  DRAFT_MODEL_ID,
732
  iteration,
733
  )
@@ -759,6 +789,7 @@ def generate_wiki_response_with_prompt(
759
  max_tokens=_DRAFT_MODEL_MAX_TOKENS or 65_536,
760
  stream=True, # avoids a gateway 504 on long generations
761
  echo_stream=echo_stream,
 
762
  )
763
  span["prompt_tokens"] = completion.usage.prompt_tokens
764
  span["total_tokens"] = completion.usage.total_tokens
@@ -774,7 +805,8 @@ def generate_wiki_response_with_prompt(
774
  # from scratch (wiping read_pages/rejected_claims and burning a
775
  # full redraft loop) instead of failing once, locally, here.
776
  logger.warning(
777
- "generate_wiki_response: subagent provider call failed: %s", e
 
778
  )
779
  conversation.record_event(
780
  {
@@ -812,7 +844,8 @@ def generate_wiki_response_with_prompt(
812
  wiki_root=wiki_root,
813
  )
814
  logger.info(
815
- "subagent tool call: %s(%s) -> %s", tc.name, tc.arguments, _clip(result)
 
816
  )
817
  # Only recall_topic results ground the judge β€” a category
818
  # browse is a navigation aid (which topics exist), not
@@ -854,7 +887,9 @@ def generate_wiki_response_with_prompt(
854
  }
855
  )
856
  logger.info(
857
- "subagent tried to answer before recalling anything (iteration %d): %s",
 
 
858
  iteration,
859
  _clip(answer),
860
  )
@@ -877,7 +912,8 @@ def generate_wiki_response_with_prompt(
877
  }
878
  )
879
  logger.info(
880
- "subagent draft answer (iteration %d, %d page(s) recalled): %s",
 
881
  iteration,
882
  len(read_pages),
883
  _clip(answer),
@@ -909,6 +945,7 @@ def generate_wiki_response_with_prompt(
909
  use_verbatim_check=use_verbatim_judge,
910
  judge_prompt_override=judge_prompt_override,
911
  echo_stream=echo_stream,
 
912
  )
913
  )
914
  except ProviderTimeoutExhaustedError as e:
@@ -916,7 +953,10 @@ def generate_wiki_response_with_prompt(
916
  # unlike a judge FAIL, there's nothing to feed the redraft loop.
917
  # Fail once, locally, same convention as the other no-safe-draft
918
  # exits above (context window, subagent provider_error).
919
- logger.warning("generate_wiki_response: judge call timed out: %s", e)
 
 
 
920
  conversation.record_event(
921
  {
922
  "type": "pipeline_step",
@@ -940,7 +980,8 @@ def generate_wiki_response_with_prompt(
940
  }
941
  )
942
  logger.info(
943
- "judge verdict (iteration %d): %s β€” %s", iteration, verdict, _clip(judge_reasoning)
 
944
  )
945
  if verdict == "PASS":
946
  if not run_attribution_check:
@@ -959,7 +1000,8 @@ def generate_wiki_response_with_prompt(
959
  with timed_block("subagent/attribution") as span:
960
  attr_verdict, attr_reasoning, attr_model_reasoning = (
961
  check_attribution(
962
- answer, usage_sink=span, echo_stream=echo_stream
 
963
  )
964
  )
965
  conversation.record_event(
@@ -1029,7 +1071,8 @@ def generate_wiki_response_with_prompt(
1029
  msgs.append(Message(role="user", content=improve_prompt))
1030
  except Exception as e:
1031
  logger.exception(
1032
- "generate_wiki_response: unexpected error, degrading to fallback: %s", e
 
1033
  )
1034
  conversation.record_event(
1035
  {
@@ -1042,7 +1085,8 @@ def generate_wiki_response_with_prompt(
1042
  return FALLBACK_MESSAGE, True
1043
 
1044
  logger.warning(
1045
- "generate_wiki_response exhausted %d iterations without a PASS",
 
1046
  MAX_ITERATIONS,
1047
  )
1048
  conversation.record_event(
 
18
  )
19
  from agent.conversation_history import ConversationHistory
20
  from agent.documents import build_documents_block, render_documents
21
+ from agent.skill_decorators import needs_conversation, needs_documents, needs_request_id
22
  from agent.skills.pediatry_wiki.scripts.browse_category import (
23
  TOP_INDEX_PATH,
24
  browse_category,
 
320
  use_verbatim_check: bool = False,
321
  judge_prompt_override: str | None = None,
322
  echo_stream: bool = False,
323
+ request_id: str | None = None,
324
  ) -> tuple[str, str, str | None, list | None]:
325
  """Return (verdict, reasoning, model_reasoning, claims). verdict is
326
  PASS / FAIL / ERROR.
 
350
  entirely (eval-only). Must take {question}/{reference}/{answer} and
351
  return the same verdict/reasoning JSON shape.
352
 
353
+ `echo_stream=True` logs the judge's reasoning/content tokens live as they
354
+ arrive (eval-only β€” off by default so production calls stay silent).
355
+
356
+ `request_id` tags every log line this call emits β€” including from a
357
+ timed-out attempt whose thread keeps running in the background and may
358
+ still log or complete later (HTTP requests can't be cancelled mid-flight;
359
+ see HFChatProvider.chat's docstring) β€” so it can be told apart from a
360
+ concurrent or later call sharing the same log stream.
361
  """
362
  rendered_docs = render_documents(documents)
363
  documents_section = (
 
373
  documents_section=documents_section,
374
  )
375
 
376
+ def _call(call_id: str):
377
  return _provider().chat(
378
  messages=[Message(role="user", content=prompt)],
379
  model_id=JUDGE_MODEL_ID,
 
384
  max_tokens=_JUDGE_MODEL_MAX_TOKENS or 65_536,
385
  stream=True, # avoids a gateway 504 on long reasoning traces
386
  echo_stream=echo_stream,
387
+ request_id=call_id,
388
  )
389
 
390
  completion = None
391
  last_timeout = False
392
  for attempt in range(1, JUDGE_CALL_MAX_ATTEMPTS_ON_TIMEOUT + 1):
393
+ call_id = f"{request_id}:judge:attempt{attempt}"
394
  try:
395
+ completion = _judge_call_pool.submit(_call, call_id).result(
396
  timeout=JUDGE_CALL_TIMEOUT_SECONDS
397
  )
398
  last_timeout = False
 
400
  except FutureTimeoutError:
401
  last_timeout = True
402
  logger.warning(
403
+ "[%s] call exceeded %ss, abandoning and retrying "
404
+ "(attempt %d/%d) β€” its thread keeps running and may still "
405
+ "log or complete later",
406
+ call_id,
407
  JUDGE_CALL_TIMEOUT_SECONDS,
408
  attempt,
409
  JUDGE_CALL_MAX_ATTEMPTS_ON_TIMEOUT,
 
414
  # the 3900-token cap); if it still exhausts its own budget and raises,
415
  # don't let that crash the whole chat turn β€” degrade to the same
416
  # graceful ERROR verdict the unparseable-output path below returns.
417
+ logger.warning("[%s] provider call failed: %s", call_id, e)
418
  return "ERROR", f"Judge provider call failed: {e}", None, None
419
  if last_timeout:
420
  # Deliberately not a graceful "ERROR" verdict: that would feed a
 
423
  # iteration for no signal, then likely timing out again on the same
424
  # slow request. Raise instead so the caller can fail once, locally.
425
  raise ProviderTimeoutExhaustedError(
426
+ f"[{request_id}] judge call timed out after "
427
+ f"{JUDGE_CALL_MAX_ATTEMPTS_ON_TIMEOUT} attempts of "
428
+ f"{JUDGE_CALL_TIMEOUT_SECONDS}s each"
429
  )
430
  if usage_sink is not None:
431
  usage_sink["prompt_tokens"] = completion.usage.prompt_tokens
 
478
 
479
 
480
  def check_attribution(
481
+ answer: str, *, usage_sink: dict | None = None, echo_stream: bool = False,
482
+ request_id: str | None = None,
483
  ) -> tuple[str, str, str | None]:
484
  """Return (verdict, reasoning, model_reasoning). FAIL means the reply
485
  attributes its answer to an external source instead of speaking directly.
 
490
  the thinking trace is worth persisting, not discarding.
491
 
492
  `usage_sink`, if given, is populated with this call's prompt_tokens /
493
+ total_tokens / attempts (see judge_answer). `echo_stream`/`request_id` β€”
494
+ see judge_answer."""
495
  prompt = ATTRIBUTION_PROMPT.format(answer=answer)
496
+ call_id = f"{request_id}:attribution"
497
  try:
498
  completion = _provider().chat(
499
  messages=[Message(role="user", content=prompt)],
 
505
  max_tokens=_JUDGE_MODEL_MAX_TOKENS or 65_536,
506
  stream=True, # avoids a gateway 504 on long reasoning traces
507
  echo_stream=echo_stream,
508
+ request_id=call_id,
509
  )
510
  except ProviderError as e:
511
  # Same reasoning as judge_answer's wrap: don't let a provider failure
 
513
  # caller treats any non-FAIL verdict as "ship the already-judge-passed
514
  # answer", so degrading to ERROR is a safe, bounded fallback instead of
515
  # crashing the whole generate_wiki_response call.
516
+ logger.warning("[%s] provider call failed: %s", call_id, e)
517
  return "ERROR", f"Attribution provider call failed: {e}", None
518
  if usage_sink is not None:
519
  usage_sink["prompt_tokens"] = completion.usage.prompt_tokens
 
591
 
592
  @needs_conversation
593
  @needs_documents
594
+ @needs_request_id
595
  def generate_wiki_response(
596
  user_query: str,
597
  conversation: ConversationHistory,
598
  documents: dict[str, str] | None = None,
599
+ request_id: str | None = None,
600
  ) -> tuple[str, bool]:
601
  """Production entry point: short-answer draft prompt, reject ledger, and
602
  the split top-level category index β€” the configuration formerly split
 
614
  # (or if) it ever finishes β€” this is prod, not eval, but nothing here
615
  # reaches the chat response itself.
616
  echo_stream=True,
617
+ request_id=request_id,
618
  )
619
 
620
 
 
632
  judge_prompt_override: str | None = None,
633
  run_attribution_check: bool = True,
634
  echo_stream: bool = False,
635
+ request_id: str | None = None,
636
  ) -> tuple[str, bool]:
637
  """Restricted mini-agent: read wiki pages β†’ answer β†’ self-correct vs. the judge.
638
 
 
682
  production wiki's category index).
683
 
684
  `echo_stream=True` streams every LLM call this makes (draft, judge,
685
+ attribution) and logs tokens live as they arrive β€” for watching a
686
+ long-running production call, not the chat response itself.
687
+
688
+ `request_id` is the current turn's id (the same id AgentClient.call()
689
+ generates as `reply_id`, injected here via `@needs_request_id`) β€” tags
690
+ every draft/judge/attribution log line this call produces, including from
691
+ a timed-out attempt whose thread keeps running in the background and may
692
+ still log or complete later, so it can be told apart from a concurrent or
693
+ later call sharing the same log stream. Not minted here: a leaf pipeline
694
+ function inventing its own id would be disconnected from the id the rest
695
+ of the system already uses to identify this turn.
696
 
697
  Records every sub-step (tool call, tool result, draft, judge verdict,
698
  improve cycle) as `pipeline_step` events on the outer ConversationHistory.
 
713
  "type": "pipeline_step",
714
  "step": "generate_wiki_response/start",
715
  "user_query": user_query,
716
+ "request_id": request_id,
717
  }
718
  )
719
 
 
755
  # would let the model repeat a discarded mistake or answer
756
  # ungrounded instead of failing loudly, so bail out instead.
757
  logger.warning(
758
+ "[%s] generate_wiki_response exceeded the context window "
759
+ "for %s (iteration %d)",
760
+ request_id,
761
  DRAFT_MODEL_ID,
762
  iteration,
763
  )
 
789
  max_tokens=_DRAFT_MODEL_MAX_TOKENS or 65_536,
790
  stream=True, # avoids a gateway 504 on long generations
791
  echo_stream=echo_stream,
792
+ request_id=f"{request_id}:draft:iter{iteration}",
793
  )
794
  span["prompt_tokens"] = completion.usage.prompt_tokens
795
  span["total_tokens"] = completion.usage.total_tokens
 
805
  # from scratch (wiping read_pages/rejected_claims and burning a
806
  # full redraft loop) instead of failing once, locally, here.
807
  logger.warning(
808
+ "[%s] generate_wiki_response: subagent provider call failed: %s",
809
+ request_id, e,
810
  )
811
  conversation.record_event(
812
  {
 
844
  wiki_root=wiki_root,
845
  )
846
  logger.info(
847
+ "[%s] subagent tool call: %s(%s) -> %s",
848
+ request_id, tc.name, tc.arguments, _clip(result),
849
  )
850
  # Only recall_topic results ground the judge β€” a category
851
  # browse is a navigation aid (which topics exist), not
 
887
  }
888
  )
889
  logger.info(
890
+ "[%s] subagent tried to answer before recalling anything "
891
+ "(iteration %d): %s",
892
+ request_id,
893
  iteration,
894
  _clip(answer),
895
  )
 
912
  }
913
  )
914
  logger.info(
915
+ "[%s] subagent draft answer (iteration %d, %d page(s) recalled): %s",
916
+ request_id,
917
  iteration,
918
  len(read_pages),
919
  _clip(answer),
 
945
  use_verbatim_check=use_verbatim_judge,
946
  judge_prompt_override=judge_prompt_override,
947
  echo_stream=echo_stream,
948
+ request_id=f"{request_id}:iter{iteration}",
949
  )
950
  )
951
  except ProviderTimeoutExhaustedError as e:
 
953
  # unlike a judge FAIL, there's nothing to feed the redraft loop.
954
  # Fail once, locally, same convention as the other no-safe-draft
955
  # exits above (context window, subagent provider_error).
956
+ logger.warning(
957
+ "[%s] generate_wiki_response: judge call timed out: %s",
958
+ request_id, e,
959
+ )
960
  conversation.record_event(
961
  {
962
  "type": "pipeline_step",
 
980
  }
981
  )
982
  logger.info(
983
+ "[%s] judge verdict (iteration %d): %s β€” %s",
984
+ request_id, iteration, verdict, _clip(judge_reasoning),
985
  )
986
  if verdict == "PASS":
987
  if not run_attribution_check:
 
1000
  with timed_block("subagent/attribution") as span:
1001
  attr_verdict, attr_reasoning, attr_model_reasoning = (
1002
  check_attribution(
1003
+ answer, usage_sink=span, echo_stream=echo_stream,
1004
+ request_id=f"{request_id}:iter{iteration}",
1005
  )
1006
  )
1007
  conversation.record_event(
 
1071
  msgs.append(Message(role="user", content=improve_prompt))
1072
  except Exception as e:
1073
  logger.exception(
1074
+ "[%s] generate_wiki_response: unexpected error, degrading to fallback: %s",
1075
+ request_id, e,
1076
  )
1077
  conversation.record_event(
1078
  {
 
1085
  return FALLBACK_MESSAGE, True
1086
 
1087
  logger.warning(
1088
+ "[%s] generate_wiki_response exhausted %d iterations without a PASS",
1089
+ request_id,
1090
  MAX_ITERATIONS,
1091
  )
1092
  conversation.record_event(
agent/skills/pediatry_wiki_minimal/scripts/generate_wiki_response_minimal.py CHANGED
@@ -18,7 +18,7 @@ generate_wiki_response, which is what SKILL.md instructs the model to call.
18
  """
19
 
20
  from agent.conversation_history import ConversationHistory
21
- from agent.skill_decorators import needs_conversation, needs_documents
22
  from agent.skills.pediatry_wiki.scripts.generate_wiki_response import (
23
  generate_wiki_response_with_prompt,
24
  )
@@ -32,10 +32,12 @@ from agent.skills.pediatry_wiki.scripts.prompts_minimal_judge import (
32
 
33
  @needs_conversation
34
  @needs_documents
 
35
  def generate_wiki_response(
36
  user_query: str,
37
  conversation: ConversationHistory,
38
  documents: dict[str, str] | None = None,
 
39
  ) -> tuple[str, bool]:
40
  return generate_wiki_response_with_prompt(
41
  user_query,
@@ -49,4 +51,5 @@ def generate_wiki_response(
49
  # (or if) it ever finishes β€” this is prod, not eval, but nothing here
50
  # reaches the chat response itself.
51
  echo_stream=True,
 
52
  )
 
18
  """
19
 
20
  from agent.conversation_history import ConversationHistory
21
+ from agent.skill_decorators import needs_conversation, needs_documents, needs_request_id
22
  from agent.skills.pediatry_wiki.scripts.generate_wiki_response import (
23
  generate_wiki_response_with_prompt,
24
  )
 
32
 
33
  @needs_conversation
34
  @needs_documents
35
+ @needs_request_id
36
  def generate_wiki_response(
37
  user_query: str,
38
  conversation: ConversationHistory,
39
  documents: dict[str, str] | None = None,
40
+ request_id: str | None = None,
41
  ) -> tuple[str, bool]:
42
  return generate_wiki_response_with_prompt(
43
  user_query,
 
51
  # (or if) it ever finishes β€” this is prod, not eval, but nothing here
52
  # reaches the chat response itself.
53
  echo_stream=True,
54
+ request_id=request_id,
55
  )
providers/hf.py CHANGED
@@ -95,14 +95,21 @@ class HFChatProvider:
95
  enable_thinking: bool | None = None,
96
  stream: bool = False,
97
  echo_stream: bool = False,
 
98
  ) -> ChatCompletion:
99
  """stream=True avoids a gateway 504 on long (e.g. high-reasoning-effort)
100
  generations β€” the router times out a single blocking request but not a
101
  continuously-flowing streamed one. The stream is fully consumed and
102
  reassembled into the same shape a non-streaming call returns; callers
103
  never see a difference beyond avoiding the timeout, unless echo_stream
104
- is also set β€” then reasoning/content tokens print to stdout live, for
105
- an interactive caller that would otherwise wait in silence."""
 
 
 
 
 
 
106
  kwargs: dict = {
107
  "model": model_id,
108
  "messages": [_msg_to_dict(m) for m in messages],
@@ -139,7 +146,7 @@ class HFChatProvider:
139
  try:
140
  completion = self._create_with_transient_retry(
141
  kwargs, rate_limit_info, server_error_info,
142
- stream=stream, echo_stream=echo_stream,
143
  )
144
  except HfHubHTTPError as e:
145
  status = e.response.status_code
@@ -157,7 +164,8 @@ class HFChatProvider:
157
  server_error_info["wait_seconds"] += wait
158
  server_error_info["retries"] += 1
159
  logger.warning(
160
- "%d (attempt %d/%d): %s β€” retrying in %.1fs",
 
161
  status,
162
  attempt,
163
  self._max_attempts,
@@ -170,7 +178,7 @@ class HFChatProvider:
170
  raise
171
  last_exc = e
172
  logger.warning(
173
- "400 (attempt %d/%d): %s", attempt, self._max_attempts, e
174
  )
175
  try:
176
  logger.warning(" raw response body: %s", e.response.text)
@@ -181,14 +189,14 @@ class HFChatProvider:
181
  continue
182
 
183
  if completion is None:
184
- logger.warning("empty completion (%d/%d)", attempt, self._max_attempts)
185
  continue
186
  if completion.choices[0]["finish_reason"] == "length":
187
- logger.warning("truncated (%d/%d)", attempt, self._max_attempts)
188
  continue
189
  msg = completion.choices[0]["message"]
190
  if not (msg.get("content") or "").strip() and not msg.get("tool_calls"):
191
- logger.warning("empty content (%d/%d)", attempt, self._max_attempts)
192
  continue
193
 
194
  result = _to_completion(completion)
@@ -199,7 +207,7 @@ class HFChatProvider:
199
  result.server_error_retries = server_error_info["retries"]
200
  return result
201
 
202
- logger.error("All %d attempts exhausted", self._max_attempts)
203
  # 429/502 raise directly above; only 400 and non-502 5xx land here.
204
  if isinstance(last_exc, HfHubHTTPError):
205
  if last_exc.response.status_code == 400:
@@ -214,7 +222,7 @@ class HFChatProvider:
214
 
215
  def _create_with_transient_retry(
216
  self, kwargs: dict, rate_limit_info: dict, server_error_info: dict,
217
- stream: bool = False, echo_stream: bool = False,
218
  ):
219
  """.create(), but 429s and 502s sleep-and-retry internally (own
220
  backoff/budget each); other 5xx propagate for chat()'s old handling.
@@ -228,7 +236,8 @@ class HFChatProvider:
228
  try:
229
  if stream:
230
  return _accumulate_stream(
231
- self._client.chat.completions.create(**kwargs), echo=echo_stream
 
232
  )
233
  return self._client.chat.completions.create(**kwargs)
234
  except HfHubHTTPError as e:
@@ -242,7 +251,8 @@ class HFChatProvider:
242
  base_wait * (2 ** (rl_attempt - 1)), _RATE_LIMIT_MAX_WAIT_SECONDS
243
  )
244
  logger.warning(
245
- "rate limited (429) β€” retrying in %.2fs (%d/%d)",
 
246
  wait,
247
  rl_attempt,
248
  _RATE_LIMIT_MAX_RETRIES,
@@ -265,8 +275,9 @@ class HFChatProvider:
265
  remaining,
266
  )
267
  logger.warning(
268
- "502 β€” retrying in %.1fs (%.1fs/%.1fs of server-error "
269
  "budget used)",
 
270
  wait,
271
  server_error_info["wait_seconds"],
272
  _SERVER_ERROR_MAX_WAIT_SECONDS,
@@ -287,13 +298,19 @@ class _StreamedCompletion:
287
  self.usage = usage
288
 
289
 
290
- def _accumulate_stream(chunks, echo: bool = False) -> _StreamedCompletion:
 
 
291
  """Reassembles a stream of ChatCompletionStreamOutput chunks into the same
292
  shape a non-streaming response has, so downstream code (_to_completion,
293
  chat()'s own post-processing) doesn't need to know the call was streamed.
294
 
295
- echo=True prints reasoning/content tokens to stdout as they arrive β€” for
296
- an interactive caller (e.g. an eval script) that would otherwise wait in
 
 
 
 
297
  silence for a slow, high-reasoning-effort call to finish."""
298
  content_parts: list[str] = []
299
  reasoning_parts: list[str] = []
@@ -314,14 +331,16 @@ def _accumulate_stream(chunks, echo: bool = False) -> _StreamedCompletion:
314
  if delta["reasoning"]:
315
  if echo:
316
  if not in_reasoning:
317
- print("\n--- reasoning (streaming) ---", flush=True)
318
  in_reasoning = True
319
  print(delta["reasoning"], end="", flush=True)
320
  reasoning_parts.append(delta["reasoning"])
321
  if delta["content"]:
322
  if echo:
323
  if not in_content:
324
- print("\n--- content (streaming) ---", flush=True)
 
 
325
  in_content = True
326
  print(delta["content"], end="", flush=True)
327
  content_parts.append(delta["content"])
@@ -335,6 +354,8 @@ def _accumulate_stream(chunks, echo: bool = False) -> _StreamedCompletion:
335
  entry["arguments"] += tc["function"]["arguments"]
336
  if choice["finish_reason"]:
337
  finish_reason = choice["finish_reason"]
 
 
338
  message = {
339
  "content": "".join(content_parts) or None,
340
  "reasoning": "".join(reasoning_parts) or None,
 
95
  enable_thinking: bool | None = None,
96
  stream: bool = False,
97
  echo_stream: bool = False,
98
+ request_id: str | None = None,
99
  ) -> ChatCompletion:
100
  """stream=True avoids a gateway 504 on long (e.g. high-reasoning-effort)
101
  generations β€” the router times out a single blocking request but not a
102
  continuously-flowing streamed one. The stream is fully consumed and
103
  reassembled into the same shape a non-streaming call returns; callers
104
  never see a difference beyond avoiding the timeout, unless echo_stream
105
+ is also set β€” then reasoning/content tokens are logged live, for an
106
+ interactive caller that would otherwise wait in silence.
107
+
108
+ `request_id` tags every log line this call emits, including ones from
109
+ a timed-out attempt whose thread the caller gave up on but which keeps
110
+ running in the background (HTTP requests can't be cancelled mid-flight)
111
+ β€” without it, that abandoned call's later output is indistinguishable
112
+ from a fresh, legitimate one sharing the same stdout/log stream."""
113
  kwargs: dict = {
114
  "model": model_id,
115
  "messages": [_msg_to_dict(m) for m in messages],
 
146
  try:
147
  completion = self._create_with_transient_retry(
148
  kwargs, rate_limit_info, server_error_info,
149
+ stream=stream, echo_stream=echo_stream, request_id=request_id,
150
  )
151
  except HfHubHTTPError as e:
152
  status = e.response.status_code
 
164
  server_error_info["wait_seconds"] += wait
165
  server_error_info["retries"] += 1
166
  logger.warning(
167
+ "[%s] %d (attempt %d/%d): %s β€” retrying in %.1fs",
168
+ request_id,
169
  status,
170
  attempt,
171
  self._max_attempts,
 
178
  raise
179
  last_exc = e
180
  logger.warning(
181
+ "[%s] 400 (attempt %d/%d): %s", request_id, attempt, self._max_attempts, e
182
  )
183
  try:
184
  logger.warning(" raw response body: %s", e.response.text)
 
189
  continue
190
 
191
  if completion is None:
192
+ logger.warning("[%s] empty completion (%d/%d)", request_id, attempt, self._max_attempts)
193
  continue
194
  if completion.choices[0]["finish_reason"] == "length":
195
+ logger.warning("[%s] truncated (%d/%d)", request_id, attempt, self._max_attempts)
196
  continue
197
  msg = completion.choices[0]["message"]
198
  if not (msg.get("content") or "").strip() and not msg.get("tool_calls"):
199
+ logger.warning("[%s] empty content (%d/%d)", request_id, attempt, self._max_attempts)
200
  continue
201
 
202
  result = _to_completion(completion)
 
207
  result.server_error_retries = server_error_info["retries"]
208
  return result
209
 
210
+ logger.error("[%s] All %d attempts exhausted", request_id, self._max_attempts)
211
  # 429/502 raise directly above; only 400 and non-502 5xx land here.
212
  if isinstance(last_exc, HfHubHTTPError):
213
  if last_exc.response.status_code == 400:
 
222
 
223
  def _create_with_transient_retry(
224
  self, kwargs: dict, rate_limit_info: dict, server_error_info: dict,
225
+ stream: bool = False, echo_stream: bool = False, request_id: str | None = None,
226
  ):
227
  """.create(), but 429s and 502s sleep-and-retry internally (own
228
  backoff/budget each); other 5xx propagate for chat()'s old handling.
 
236
  try:
237
  if stream:
238
  return _accumulate_stream(
239
+ self._client.chat.completions.create(**kwargs),
240
+ echo=echo_stream, request_id=request_id,
241
  )
242
  return self._client.chat.completions.create(**kwargs)
243
  except HfHubHTTPError as e:
 
251
  base_wait * (2 ** (rl_attempt - 1)), _RATE_LIMIT_MAX_WAIT_SECONDS
252
  )
253
  logger.warning(
254
+ "[%s] rate limited (429) β€” retrying in %.2fs (%d/%d)",
255
+ request_id,
256
  wait,
257
  rl_attempt,
258
  _RATE_LIMIT_MAX_RETRIES,
 
275
  remaining,
276
  )
277
  logger.warning(
278
+ "[%s] 502 β€” retrying in %.1fs (%.1fs/%.1fs of server-error "
279
  "budget used)",
280
+ request_id,
281
  wait,
282
  server_error_info["wait_seconds"],
283
  _SERVER_ERROR_MAX_WAIT_SECONDS,
 
298
  self.usage = usage
299
 
300
 
301
+ def _accumulate_stream(
302
+ chunks, echo: bool = False, request_id: str | None = None
303
+ ) -> _StreamedCompletion:
304
  """Reassembles a stream of ChatCompletionStreamOutput chunks into the same
305
  shape a non-streaming response has, so downstream code (_to_completion,
306
  chat()'s own post-processing) doesn't need to know the call was streamed.
307
 
308
+ echo=True logs one line via the logger announcing that reasoning/content
309
+ streaming has started β€” tagged with request_id, so a call whose caller
310
+ gave up waiting (but whose thread keeps running β€” see HFChatProvider.chat's
311
+ docstring) can still be told apart from a concurrent or later one sharing
312
+ the same log stream β€” then prints each token live to stdout as it arrives,
313
+ same as before, for an interactive caller that would otherwise wait in
314
  silence for a slow, high-reasoning-effort call to finish."""
315
  content_parts: list[str] = []
316
  reasoning_parts: list[str] = []
 
331
  if delta["reasoning"]:
332
  if echo:
333
  if not in_reasoning:
334
+ logger.info("[%s] streaming reasoning...", request_id)
335
  in_reasoning = True
336
  print(delta["reasoning"], end="", flush=True)
337
  reasoning_parts.append(delta["reasoning"])
338
  if delta["content"]:
339
  if echo:
340
  if not in_content:
341
+ if in_reasoning:
342
+ print(flush=True) # end the reasoning line first
343
+ logger.info("[%s] streaming content...", request_id)
344
  in_content = True
345
  print(delta["content"], end="", flush=True)
346
  content_parts.append(delta["content"])
 
354
  entry["arguments"] += tc["function"]["arguments"]
355
  if choice["finish_reason"]:
356
  finish_reason = choice["finish_reason"]
357
+ if echo and (in_reasoning or in_content):
358
+ print(flush=True) # close out the last open printed line
359
  message = {
360
  "content": "".join(content_parts) or None,
361
  "reasoning": "".join(reasoning_parts) or None,