Skip to content
This repository was archived by the owner on Sep 23, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Next Next commit
feat(telemetry): align events with TS schema, add trace_id and missin…
…g events

Align the Python telemetry surface with the TS rewrite's event registry
(agent-core-v2 events.ts), and attach the KFC x-trace-id response header
to request-scoped events so they can join with server request logs.

kosong:
- Kimi provider captures x-trace-id via with_raw_response (stream and
  non-stream); trace_id carried on APIStatusError (openai + httpx paths)
- StreamedMessage/GenerateResult/StepResult thread trace_id through
- New on_trace_id callback fired when response headers arrive (before
  streaming, so events emitted mid-stream see the current request)

kimi-cli:
- Two-level current-trace holder (ContextVar per turn task + root mirror
  for UI) feeding trace_id into api_error, compaction_*, turn_interrupted,
  turn_ended, cancel, tool_call, tool_call_dedup_detected, tool_call_repeat,
  permission_approval_result, question_answered/dismissed
- compaction_finished/failed renamed to TS property names
  (source/tokens_before/tokens_after/input_tokens/output_tokens) and gain
  round/thinking_effort/compacted_count
- api_error gains retryable (TS isRetryableGenerateError table, incl.
  408/409/529), provider_type/protocol and overloaded(529); the non-TS
  'api' error_type folds into 'other'
- tool_call gains tool_call_id, cancelled outcome and trace_id;
  error_type becomes the TS enum ('error'|'cancelled') with the exception
  class moved to error_class
- tool_call_dedup_detected re-added at same-step/cross-step detection
- turn_interrupted gains interrupt_reason; cancel gains from
  (streaming/compacting); question_answered gains answered
- New events: turn_ended (unconditional at turn end) and
  permission_approval_result (all approval decision points; existing
  tool_approved/tool_rejected kept for compatibility)
  • Loading branch information
7Sageer committed Jul 15, 2026
commit b749c3efe89dabd634573112b857cc2fdc77366a
6 changes: 6 additions & 0 deletions packages/kosong/src/kosong/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@ async def step(
*,
on_message_part: Callback[[StreamedMessagePart], None] | None = None,
on_tool_result: Callable[[ToolResult], None] | None = None,
on_trace_id: Callback[[str | None], None] | None = None,
) -> "StepResult":
"""
Run one agent "step". In one step, the function generates LLM response based on the given
Expand Down Expand Up @@ -162,6 +163,7 @@ async def on_tool_call(tool_call: ToolCall):
history,
on_message_part=on_message_part,
on_tool_call=on_tool_call,
on_trace_id=on_trace_id,
)
except (ChatProviderError, asyncio.CancelledError):
# cancel all the futures to avoid hanging tasks
Expand All @@ -177,6 +179,7 @@ async def on_tool_call(tool_call: ToolCall):
result.usage,
tool_calls,
tool_result_futures,
trace_id=result.trace_id,
)


Expand All @@ -197,6 +200,9 @@ class StepResult:
_tool_result_futures: dict[str, ToolResultFuture]
"""@private The futures of the results of the spawned tool calls."""

trace_id: str | None = None
"""The ``x-trace-id`` response header of the request, if the provider exposes it."""

async def tool_results(self) -> list[ToolResult]:
"""All the tool results returned by corresponding tool calls."""
if not self._tool_result_futures:
Expand Down
10 changes: 10 additions & 0 deletions packages/kosong/src/kosong/_generate.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ async def generate(
*,
on_message_part: Callback[[StreamedMessagePart], None] | None = None,
on_tool_call: Callback[[ToolCall], None] | None = None,
on_trace_id: Callback[[str | None], None] | None = None,
) -> "GenerateResult":
"""
Generate one message based on the given context.
Expand All @@ -34,6 +35,8 @@ async def generate(
history: The message history to use for generation.
on_message_part: An optional callback to be called for each raw message part.
on_tool_call: An optional callback to be called for each complete tool call.
on_trace_id: An optional callback fired with the request's ``x-trace-id``
response header as soon as it is available (before streaming starts).

Returns:
A tuple of the generated message and the token usage (if available).
Expand All @@ -51,6 +54,10 @@ async def generate(

logger.trace("Generating with history: {history}", history=history)
stream = await chat_provider.generate(system_prompt, tools, history)
if on_trace_id:
# getattr for robustness against third-party StreamedMessage
# implementations that predate the trace_id property.
await callback(on_trace_id, getattr(stream, "trace_id", None))
async for part in stream:
logger.trace("Received part: {part}", part=part)
if on_message_part:
Expand Down Expand Up @@ -92,6 +99,7 @@ async def generate(
id=stream.id,
message=message,
usage=stream.usage,
trace_id=getattr(stream, "trace_id", None),
)


Expand All @@ -105,6 +113,8 @@ class GenerateResult:
"""The generated message."""
usage: TokenUsage | None
"""The token usage of the generated message."""
trace_id: str | None = None
"""The ``x-trace-id`` response header of the request, if the provider exposes it."""


def _message_append(message: Message, part: StreamedMessagePart) -> None:
Expand Down
25 changes: 23 additions & 2 deletions packages/kosong/src/kosong/chat_provider/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,15 @@ def usage(self) -> TokenUsage | None:
"""The token usage of the streamed message."""
...

@property
def trace_id(self) -> str | None:
"""The ``x-trace-id`` response header of the underlying request, if exposed.

Only providers talking to the KFC inference service (e.g. Kimi) return a
value; all other providers return None.
"""
...


class TokenUsage(BaseModel):
"""Token usage statistics."""
Expand Down Expand Up @@ -156,11 +165,20 @@ class APIStatusError(ChatProviderError):

status_code: int
request_id: str | None
trace_id: str | None

def __init__(self, status_code: int, message: str, *, request_id: str | None = None):
def __init__(
self,
status_code: int,
message: str,
*,
request_id: str | None = None,
trace_id: str | None = None,
):
super().__init__(message)
self.status_code = status_code
self.request_id = request_id
self.trace_id = trace_id


class APIEmptyResponseError(ChatProviderError):
Expand All @@ -183,5 +201,8 @@ def convert_httpx_error(error: httpx.HTTPError) -> ChatProviderError:
return APIConnectionError(str(error))
if isinstance(error, httpx.HTTPStatusError):
req_id = error.response.headers.get("x-request-id")
return APIStatusError(error.response.status_code, str(error), request_id=req_id)
trace_id = error.response.headers.get("x-trace-id")
return APIStatusError(
error.response.status_code, str(error), request_id=req_id, trace_id=trace_id
)
return ChatProviderError(f"HTTP error: {error}")
4 changes: 4 additions & 0 deletions packages/kosong/src/kosong/chat_provider/chaos.py
Original file line number Diff line number Diff line change
Expand Up @@ -250,6 +250,10 @@ def id(self) -> str | None:
def usage(self) -> TokenUsage | None:
return self._wrapped.usage

@property
def trace_id(self) -> str | None:
return self._wrapped.trace_id

def _should_corrupt_tool_call(self) -> bool:
probability = self._config.corrupt_tool_call_probability
return probability > 0 and self._rng.random() < probability
Expand Down
4 changes: 4 additions & 0 deletions packages/kosong/src/kosong/chat_provider/echo/echo.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,3 +123,7 @@ def id(self) -> str | None:
@property
def usage(self) -> TokenUsage | None:
return self._usage

@property
def trace_id(self) -> str | None:
return None
Original file line number Diff line number Diff line change
Expand Up @@ -101,3 +101,7 @@ def id(self) -> str | None:
@property
def usage(self) -> TokenUsage | None:
return self._usage

@property
def trace_id(self) -> str | None:
return None
21 changes: 18 additions & 3 deletions packages/kosong/src/kosong/chat_provider/kimi.py
Original file line number Diff line number Diff line change
Expand Up @@ -166,15 +166,20 @@ async def generate(
generation_kwargs.pop("max_completion_tokens", None)

try:
response = await self.client.chat.completions.create(
raw_response = await self.client.chat.completions.with_raw_response.create(
model=self.model,
messages=messages,
tools=(_convert_tool(tool) for tool in tools),
stream=self.stream,
stream_options={"include_usage": True} if self.stream else omit,
**generation_kwargs,
)
return KimiStreamedMessage(response)
# The promise resolves as soon as response headers arrive (before the
# stream body), so the trace id is available even mid-stream.
# Note: LegacyAPIResponse.parse() is sync in openai SDK 2.x; it will
# become a coroutine in the next major version.
trace_id = raw_response.headers.get("x-trace-id")
return KimiStreamedMessage(raw_response.parse(), trace_id=trace_id)
except (OpenAIError, httpx.HTTPError) as e:
raise convert_error(e) from e

Expand Down Expand Up @@ -372,13 +377,19 @@ def _convert_tool(tool: Tool) -> ChatCompletionToolParam:
class KimiStreamedMessage:
"""The streamed message of the Kimi chat provider."""

def __init__(self, response: ChatCompletion | AsyncStream[ChatCompletionChunk]):
def __init__(
self,
response: ChatCompletion | AsyncStream[ChatCompletionChunk],
*,
trace_id: str | None = None,
):
if isinstance(response, ChatCompletion):
self._iter = self._convert_non_stream_response(response)
else:
self._iter = self._convert_stream_response(response)
self._id: str | None = None
self._usage: CompletionUsage | None = None
self._trace_id = trace_id

def __aiter__(self) -> AsyncIterator[StreamedMessagePart]:
return self
Expand All @@ -390,6 +401,10 @@ async def __anext__(self) -> StreamedMessagePart:
def id(self) -> str | None:
return self._id

@property
def trace_id(self) -> str | None:
return self._trace_id

@property
def usage(self) -> TokenUsage | None:
if self._usage:
Expand Down
4 changes: 4 additions & 0 deletions packages/kosong/src/kosong/chat_provider/mock.py
Original file line number Diff line number Diff line change
Expand Up @@ -78,3 +78,7 @@ def id(self) -> str:
@property
def usage(self) -> TokenUsage | None:
return None

@property
def trace_id(self) -> str | None:
return None
5 changes: 4 additions & 1 deletion packages/kosong/src/kosong/chat_provider/openai_common.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,10 @@ def convert_error(error: OpenAIError | httpx.HTTPError) -> ChatProviderError:
match error:
case openai.APIStatusError():
req_id = error.response.headers.get("x-request-id")
return APIStatusError(error.status_code, error.message, request_id=req_id)
trace_id = error.response.headers.get("x-trace-id")
return APIStatusError(
error.status_code, error.message, request_id=req_id, trace_id=trace_id
)
case openai.APITimeoutError():
return APITimeoutError(error.message)
case openai.APIConnectionError():
Expand Down
4 changes: 4 additions & 0 deletions packages/kosong/src/kosong/contrib/chat_provider/anthropic.py
Original file line number Diff line number Diff line change
Expand Up @@ -525,6 +525,10 @@ async def __anext__(self) -> StreamedMessagePart:
def id(self) -> str | None:
return self._id

@property
def trace_id(self) -> str | None:
return None

@property
def usage(self) -> TokenUsage | None:
# https://docs.claude.com/en/docs/build-with-claude/prompt-caching#tracking-cache-performance
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,10 @@ async def __anext__(self) -> StreamedMessagePart:
def id(self) -> str | None:
return self._id

@property
def trace_id(self) -> str | None:
return None

@property
def usage(self) -> TokenUsage | None:
if self._usage is None:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -242,6 +242,10 @@ async def __anext__(self) -> StreamedMessagePart:
def id(self) -> str | None:
return self._id

@property
def trace_id(self) -> str | None:
return None

@property
def usage(self) -> TokenUsage | None:
if self._usage:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -474,6 +474,10 @@ async def __anext__(self) -> StreamedMessagePart:
def id(self) -> str | None:
return self._id

@property
def trace_id(self) -> str | None:
return None

@property
def usage(self) -> TokenUsage | None:
if self._usage:
Expand Down
123 changes: 123 additions & 0 deletions packages/kosong/tests/test_trace_id.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
"""Verify ``x-trace-id`` response header capture and propagation.

The KFC inference service returns a trace id in the ``x-trace-id`` response
header. It must be captured by the Kimi provider (both success and error
paths) and propagated through ``generate``/``step`` results and the
``on_trace_id`` early callback.
"""

import httpx
import pytest
import respx

import kosong
from kosong.chat_provider import APIStatusError
from kosong.chat_provider.kimi import Kimi
from kosong.tooling.empty import EmptyToolset

TRACE_ID = "trace-abc-123"
URL = "https://api.moonshot.ai/v1/chat/completions"


def _completion_payload() -> dict[str, object]:
return {
"id": "chatcmpl-test",
"object": "chat.completion",
"model": "test-model",
"choices": [
{
"index": 0,
"message": {"role": "assistant", "content": "ok"},
"finish_reason": "stop",
}
],
"usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2},
}


@pytest.mark.asyncio
async def test_kimi_captures_trace_id_header():
with respx.mock:
respx.post(URL).mock(
return_value=httpx.Response(
200, json=_completion_payload(), headers={"x-trace-id": TRACE_ID}
)
)
provider = Kimi(model="test-model", api_key="token", stream=False)
stream = await provider.generate("", [], [])
assert stream.trace_id == TRACE_ID
async for _ in stream:
pass


@pytest.mark.asyncio
async def test_kimi_trace_id_none_without_header():
with respx.mock:
respx.post(URL).mock(return_value=httpx.Response(200, json=_completion_payload()))
provider = Kimi(model="test-model", api_key="token", stream=False)
stream = await provider.generate("", [], [])
assert stream.trace_id is None


@pytest.mark.asyncio
async def test_kimi_streaming_captures_trace_id():
"""Streaming path: with_raw_response resolves at headers, parse() yields the stream."""
sse = (
'data: {"id":"chatcmpl-x","object":"chat.completion.chunk","created":1,'
'"model":"test-model","choices":[{"index":0,'
'"delta":{"role":"assistant","content":"hi"},"finish_reason":null}]}\n\n'
'data: {"id":"chatcmpl-x","object":"chat.completion.chunk","created":1,'
'"model":"test-model","choices":[{"index":0,"delta":{},"finish_reason":"stop"}],'
'"usage":{"prompt_tokens":1,"completion_tokens":1,"total_tokens":2}}\n\n'
"data: [DONE]\n\n"
)
with respx.mock:
respx.post(URL).mock(
return_value=httpx.Response(
200,
headers={"content-type": "text/event-stream", "x-trace-id": TRACE_ID},
content=sse,
)
)
provider = Kimi(model="test-model", api_key="token", stream=True)
stream = await provider.generate("", [], [])
assert stream.trace_id == TRACE_ID
parts = [part async for part in stream]
assert parts


@pytest.mark.asyncio
async def test_api_status_error_carries_trace_id():
with respx.mock:
respx.post(URL).mock(
return_value=httpx.Response(
500,
json={"error": {"message": "boom", "type": "server_error"}},
headers={"x-trace-id": TRACE_ID},
)
)
provider = Kimi(model="test-model", api_key="token", stream=False)
with pytest.raises(APIStatusError) as exc_info:
await provider.generate("", [], [])
assert exc_info.value.trace_id == TRACE_ID


@pytest.mark.asyncio
async def test_step_result_and_on_trace_id_callback():
seen: list[str | None] = []
with respx.mock:
respx.post(URL).mock(
return_value=httpx.Response(
200, json=_completion_payload(), headers={"x-trace-id": TRACE_ID}
)
)
provider = Kimi(model="test-model", api_key="token", stream=False)
result = await kosong.step(
chat_provider=provider,
system_prompt="",
toolset=EmptyToolset(),
history=[],
on_trace_id=seen.append,
)
assert result.trace_id == TRACE_ID
assert seen == [TRACE_ID]
Loading
Loading