Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
1 change: 1 addition & 0 deletions .changes/observability-foundation.added.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
新增默认关闭的标准 OTLP 可观测性:保留 PostgreSQL 日志,关联跨 Peer Job、AI 调用与结果来源,并提供浏览器认证转发。
26 changes: 26 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -40,3 +40,29 @@ INKCRE_COMPOSE_PROJECT_NAME=inkcre

# Flags
SKIP_EXTENSIONS_SYNC=0

# Optional metadata-only OTLP export. Endpoints/credentials never enable it alone.
# Existing PostgreSQL logging is independent and remains enabled by default.
OBSRV__TELEMETRY_ENABLED=false
# Full per-signal HTTP(S) URLs, including the ingest path. Leave unused signals empty.
# URLs must not contain credentials, query strings or fragments.
OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=
OTEL_EXPORTER_OTLP_LOGS_ENDPOINT=
OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=
# Private server-side ingest headers, using standard OTEL key=value syntax.
# A per-signal value takes precedence over the common value.
OTEL_EXPORTER_OTLP_HEADERS=
OTEL_EXPORTER_OTLP_TRACES_HEADERS=
OTEL_EXPORTER_OTLP_LOGS_HEADERS=
OTEL_EXPORTER_OTLP_METRICS_HEADERS=

# SDK queue/export self-metrics. Compose forwards this into the process.
# For native launch, export this variable in the process environment as well.
# It does not enable application telemetry by itself.
OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED=true

# Seconds: per-signal timeout overrides this default; accepted range is (0, 30].
OTEL_EXPORTER_OTLP_TIMEOUT=10
OTEL_EXPORTER_OTLP_TRACES_TIMEOUT=
OTEL_EXPORTER_OTLP_LOGS_TIMEOUT=
OTEL_EXPORTER_OTLP_METRICS_TIMEOUT=
220 changes: 133 additions & 87 deletions app/business/agent/thread.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import asyncio
from enum import StrEnum
import inspect
import hashlib
import json
import time
import traceback
Expand All @@ -26,6 +27,7 @@
from .persistence import ThreadID, ThreadPersistenceBackend, ThreadState
from .debug import trace
from libs.obsrv.main import get_logger
from libs.obsrv.telemetry import operation


logger = get_logger().getChild("agent.thread")
Expand Down Expand Up @@ -92,88 +94,108 @@ def start_turn(self, input: UserMessage) -> asyncio.Task[TurnTermination]:
async def _run_turn(self, input: UserMessage) -> TurnTermination:
self._turn_index += 1
self._model_calls = 0
started = time.monotonic()
await trace(
"agent.turn.started",
self.id,
turn=self._turn_index,
input=input,
model=self.model,
max_model_calls=self.max_model_calls_per_turn,
)
try:
outcome = await self._execute_turn(input)
except asyncio.CancelledError:
with operation(
"agent.turn",
attributes={
"gen_ai.conversation.id": str(self.id),
"inkcre.agent.turn": self._turn_index,
},
) as observation:
span = observation.span
started = time.monotonic()
await trace(
"agent.turn.finished",
"agent.turn.started",
self.id,
turn=self._turn_index,
model_calls=self._model_calls,
outcome="cancelled",
elapsed_seconds=time.monotonic() - started,
input=input,
model=self.model,
max_model_calls=self.max_model_calls_per_turn,
)
raise
except Exception as error:
try:
outcome = await self._execute_turn(input)
except asyncio.CancelledError:
span.set_attribute("inkcre.agent.outcome", "cancelled")
await trace(
"agent.turn.finished",
self.id,
turn=self._turn_index,
model_calls=self._model_calls,
outcome="cancelled",
elapsed_seconds=time.monotonic() - started,
)
raise
except Exception as error:
span.set_attribute("inkcre.agent.outcome", "failed")
await trace(
"agent.turn.finished",
self.id,
turn=self._turn_index,
model_calls=self._model_calls,
outcome="failed",
error_type=type(error).__name__,
error=str(error),
traceback=traceback.format_exc(),
elapsed_seconds=time.monotonic() - started,
)
raise
await trace(
"agent.turn.finished",
self.id,
turn=self._turn_index,
model_calls=self._model_calls,
outcome="failed",
error_type=type(error).__name__,
error=str(error),
traceback=traceback.format_exc(),
outcome=outcome,
elapsed_seconds=time.monotonic() - started,
)
raise
await trace(
"agent.turn.finished",
self.id,
turn=self._turn_index,
model_calls=self._model_calls,
outcome=outcome,
elapsed_seconds=time.monotonic() - started,
)
return outcome
span.set_attribute("inkcre.agent.outcome", outcome.value)
span.set_attribute("inkcre.agent.model_calls", self._model_calls)
return outcome

async def _execute_turn(self, input: UserMessage) -> TurnTermination:
self._state = await self._persistence.discard_trailing_incomplete_tool_calls(self.id)
self._state = await self._persistence.append(self.id, (input,))

while True:
self._model_calls += 1
await trace(
"agent.model.started",
self.id,
turn=self._turn_index,
call=self._model_calls,
)
started = time.monotonic()
assistant = await AIManager.chat(
self._state.model,
self._state.messages,
self._state.tools,
self._state.tool_choice,
)
await trace(
"agent.model.completed",
self.id,
turn=self._turn_index,
call=self._model_calls,
response=assistant,
elapsed_seconds=time.monotonic() - started,
)
if not assistant.tool_calls:
self._state = await self._persistence.append(self.id, (assistant,))
return TurnTermination.COMPLETED
with operation(
"agent.step",
attributes={
"gen_ai.conversation.id": str(self.id),
"inkcre.agent.turn": self._turn_index,
"inkcre.agent.model_call": self._model_calls,
},
):
await trace(
"agent.model.started",
self.id,
turn=self._turn_index,
call=self._model_calls,
)
started = time.monotonic()
assistant = await AIManager.chat(
self._state.model,
self._state.messages,
self._state.tools,
self._state.tool_choice,
)
await trace(
"agent.model.completed",
self.id,
turn=self._turn_index,
call=self._model_calls,
response=assistant,
elapsed_seconds=time.monotonic() - started,
)
if not assistant.tool_calls:
self._state = await self._persistence.append(self.id, (assistant,))
return TurnTermination.COMPLETED

results = await self._execute_tool_batch(assistant)
self._state = await self._persistence.append(
self.id,
(assistant, ToolResultMessage(results=results)),
)
if self._model_calls >= self._state.max_model_calls_per_turn:
return TurnTermination.MAX_MODEL_CALLS
results = await self._execute_tool_batch(assistant)
self._state = await self._persistence.append(
self.id,
(assistant, ToolResultMessage(results=results)),
)
if self._model_calls >= self._state.max_model_calls_per_turn:
return TurnTermination.MAX_MODEL_CALLS

async def _execute_tool_batch(
self,
Expand All @@ -189,37 +211,61 @@ async def _execute_tool_batch(
return tuple(await asyncio.gather(*tasks))

async def _execute_tool_call(self, call: ToolCall) -> ToolResult:
await trace(
"agent.tool.started",
self.id,
turn=self._turn_index,
call=self._model_calls,
tool_call=call,
)
started = time.monotonic()
try:
result = await self._invoke_tool_call(call)
except asyncio.CancelledError:
with operation(
"agent.tool",
attributes={
"gen_ai.operation.name": "execute_tool",
"gen_ai.conversation.id": str(self.id),
"inkcre.agent.turn": self._turn_index,
"inkcre.agent.model_call": self._model_calls,
# Provider-generated IDs may contain content; retain correlation as a digest.
"inkcre.agent.tool_call.id_hash": hashlib.sha256(
call.id.encode(errors="replace")
).hexdigest(),
},
) as observation:
span = observation.span
tool = self._tools.get(call.tool)
if tool is not None:
span.set_attribute("gen_ai.tool.name", tool.definition.id)
await trace(
"agent.tool.cancelled",
"agent.tool.started",
self.id,
turn=self._turn_index,
call=self._model_calls,
tool_call=call,
)
started = time.monotonic()
try:
result = await self._invoke_tool_call(call)
except asyncio.CancelledError:
span.set_attribute("inkcre.agent.outcome", "cancelled")
await trace(
"agent.tool.cancelled",
self.id,
turn=self._turn_index,
call=self._model_calls,
tool_call_id=call.id,
tool=call.tool,
elapsed_seconds=time.monotonic() - started,
)
raise
await trace(
"agent.tool.completed",
self.id,
turn=self._turn_index,
call=self._model_calls,
tool_call_id=call.id,
tool=call.tool,
result=result,
elapsed_seconds=time.monotonic() - started,
)
raise
await trace(
"agent.tool.completed",
self.id,
turn=self._turn_index,
call=self._model_calls,
tool=call.tool,
result=result,
elapsed_seconds=time.monotonic() - started,
)
return result
span.set_attribute(
"inkcre.agent.outcome", "error" if result.is_error else "completed"
)
if result.is_error:
span.set_attribute("error.type", "ToolResultError")
observation.outcome = "error"
return result

async def _invoke_tool_call(self, call: ToolCall) -> ToolResult:
tool = self._tools.get(call.tool)
Expand Down
23 changes: 23 additions & 0 deletions app/business/ai/dialects/alibaba_model_studio.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,12 @@
from collections.abc import Sequence
import json
import typing
import time

from openai import AsyncStream
from openai.types.chat import ChatCompletionChunk
from opentelemetry import trace
from opentelemetry.trace import INVALID_SPAN
import pydantic

from app.schemas.ai import (
Expand All @@ -21,8 +24,11 @@
VideoContentPart,
)

from libs.obsrv.telemetry import is_enabled

from ..contracts import AIOutputContractError
from ..main import AIManager
from ..telemetry import record_response_usage
from .openai_compatible import OpenAICompatibleConfig, OpenAICompatibleDialect


Expand Down Expand Up @@ -109,6 +115,12 @@ async def chat(
if tool_choice is not None:
arguments["tool_choice"] = self._tool_choice_param(tool_choice)

span = trace.get_current_span() if is_enabled() else INVALID_SPAN
span.set_attribute("gen_ai.request.stream", True)
started = time.monotonic()
usage = None
finish_reasons: list[str] = []
first_chunk = True
try:
stream = typing.cast(
AsyncStream[ChatCompletionChunk],
Expand All @@ -117,8 +129,18 @@ async def chat(
text_parts: list[str] = []
calls: dict[int, dict[str, str]] = {}
async for chunk in stream:
if first_chunk:
span.set_attribute(
"gen_ai.response.time_to_first_chunk", time.monotonic() - started
)
first_chunk = False
# Usage-only terminal chunks have no choices. Totals replace earlier totals.
if chunk.usage is not None:
usage = chunk.usage
if not chunk.choices:
continue
if chunk.choices[0].finish_reason is not None:
finish_reasons = [chunk.choices[0].finish_reason]
delta = chunk.choices[0].delta
if isinstance(delta.content, str):
text_parts.append(delta.content)
Expand All @@ -132,6 +154,7 @@ async def chat(
if call.function.arguments:
current["arguments"] += call.function.arguments
finally:
record_response_usage("chat", usage, finish_reasons)
await client.close()

parsed_calls: list[ToolCall] = []
Expand Down
Loading
Loading