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
14 changes: 14 additions & 0 deletions backend/.env.example
Original file line number Diff line number Diff line change
@@ -1,10 +1,14 @@
TIMEFLOW_APP_NAME=TimeFlow API
TIMEFLOW_ENVIRONMENT=development
TIMEFLOW_DATABASE_URL=postgresql+psycopg://timeapp:timeapp@127.0.0.1:5432/timeapp

# WebSocket transport limits. Durations are expressed in seconds or milliseconds as named.
TIMEFLOW_WS_HANDSHAKE_TIMEOUT_SECONDS=5.0
TIMEFLOW_WS_MAX_UNAUTHENTICATED_CONNECTIONS=100
TIMEFLOW_WS_AUDIO_QUEUE_MAX_CHUNKS=32
TIMEFLOW_WS_MAX_AUDIO_DURATION_MS=120000

# Qwen realtime ASR. Replace {WorkspaceId} locally and keep the API key out of Git.
TIMEFLOW_ALIYUN_ASR_WS_URL=wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime
TIMEFLOW_ALIYUN_ASR_API_KEY=
TIMEFLOW_ALIYUN_ASR_MODEL=qwen3-asr-flash-realtime
Expand All @@ -13,3 +17,13 @@ TIMEFLOW_ALIYUN_ASR_VAD_THRESHOLD=0.0
TIMEFLOW_ALIYUN_ASR_VAD_SILENCE_DURATION_MS=400
TIMEFLOW_ALIYUN_ASR_CONNECT_TIMEOUT_SECONDS=10
TIMEFLOW_ALIYUN_ASR_FINISH_TIMEOUT_SECONDS=10

# Qwen OpenAI-compatible LLM. Replace {WorkspaceId} locally and never commit credentials.
# Empty endpoint/key values are allowed at startup; they are required only when the adapter is used.
TIMEFLOW_OPENAI_BASE_URL=https://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/compatible-mode/v1
TIMEFLOW_OPENAI_API_KEY=
TIMEFLOW_OPENAI_MODEL=qwen-flash
TIMEFLOW_OPENAI_TIMEOUT_SECONDS=30

# Maximum serial Function Calling rounds allowed during one Agent turn.
TIMEFLOW_AGENT_MAX_TOOL_ROUNDS=4
1 change: 1 addition & 0 deletions backend/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ requires-python = ">=3.11,<3.12"
dependencies = [
"alembic>=1.16,<2",
"fastapi>=0.140.7,<1",
"openai>=2,<3",
"psycopg[binary]>=3.2,<4",
"pydantic>=2.13,<3",
"python-dotenv>=1.1,<2",
Expand Down
5 changes: 5 additions & 0 deletions backend/src/timeflow/infrastructure/external/llm/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
"""Language model provider adapters."""

from timeflow.infrastructure.external.llm.openai_compatible import OpenAICompatibleLlm

__all__ = ["OpenAICompatibleLlm"]
191 changes: 191 additions & 0 deletions backend/src/timeflow/infrastructure/external/llm/openai_compatible.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,191 @@
"""OpenAI-compatible streaming LLM adapter."""

from __future__ import annotations

import asyncio
from collections.abc import AsyncIterator, Sequence
from typing import Any, Protocol, cast

from openai import AsyncOpenAI

from timeflow.infrastructure.settings import Settings
from timeflow.intelligence.conversation.llm import (
AssistantToolCallMessage,
ChatMessage,
LlmEvent,
LlmMessage,
LlmPort,
LlmProtocolError,
LlmProviderError,
LlmStreamCompleted,
LlmUsage,
TextDelta,
ToolCallDelta,
ToolDefinition,
ToolResultMessage,
)


class _Client(Protocol):
chat: Any


def _message_payload(message: LlmMessage) -> dict[str, object]:
if isinstance(message, ChatMessage):
return {"role": message.role, "content": message.content}
if isinstance(message, AssistantToolCallMessage):
return {
"role": "assistant",
"content": message.content,
"tool_calls": [
{
"id": call.call_id,
"type": "function",
"function": {"name": call.name, "arguments": call.arguments},
}
for call in message.tool_calls
],
}
if isinstance(message, ToolResultMessage):
return {
"role": "tool",
"tool_call_id": message.tool_call_id,
"content": message.content,
}
raise TypeError(f"Unsupported LLM message: {type(message).__name__}")


def _tool_payload(definition: ToolDefinition) -> dict[str, object]:
return {
"type": "function",
"function": {
"name": definition.name,
"description": definition.description,
"parameters": dict(definition.parameters),
},
}


class OpenAICompatibleLlm(LlmPort):
"""Stream provider-neutral events from an OpenAI-compatible endpoint."""

def __init__(self, settings: Settings, client: _Client | None = None) -> None:
self._settings = settings
self._client = client or cast(
_Client,
AsyncOpenAI(
api_key=settings.openai_api_key or "not-configured",
base_url=settings.openai_base_url or None,
timeout=settings.openai_timeout_seconds,
),
)

def stream(
self,
messages: Sequence[LlmMessage],
tools: Sequence[ToolDefinition],
) -> AsyncIterator[LlmEvent]:
return self._stream(messages, tools)

async def _stream(
self,
messages: Sequence[LlmMessage],
tools: Sequence[ToolDefinition],
) -> AsyncIterator[LlmEvent]:
self._validate_settings()
stream: Any = None
usage: LlmUsage | None = None
finish_reason: str | None = None
try:
stream = await self._client.chat.completions.create(
model=self._settings.openai_model,
messages=[_message_payload(message) for message in messages],
tools=[_tool_payload(tool) for tool in tools],
stream=True,
stream_options={"include_usage": True},
extra_body={"enable_thinking": False},
parallel_tool_calls=False,
tool_choice="auto",
)
async for chunk in stream:
choices = getattr(chunk, "choices", None)
if not isinstance(choices, list):
raise LlmProtocolError("LLM chunk choices must be a list")
chunk_usage = getattr(chunk, "usage", None)
if not choices:
if chunk_usage is not None:
usage = self._parse_usage(chunk_usage)
continue
choice = choices[0]
finish_value = getattr(choice, "finish_reason", None)
if finish_value is not None and not isinstance(finish_value, str):
raise LlmProtocolError("LLM finish reason must be a string")
if finish_value is not None:
finish_reason = finish_value
delta = getattr(choice, "delta", None)
if delta is None:
raise LlmProtocolError("LLM choice delta is missing")
content = getattr(delta, "content", None)
if content is not None and not isinstance(content, str):
raise LlmProtocolError("LLM text delta must be a string")
if content:
yield TextDelta(content)
tool_calls = getattr(delta, "tool_calls", None)
if tool_calls is not None:
if not isinstance(tool_calls, list):
raise LlmProtocolError("LLM tool call deltas must be a list")
for call in tool_calls:
yield self._parse_tool_call_delta(call)
yield LlmStreamCompleted(finish_reason, usage)
except asyncio.CancelledError:
raise
except (LlmProtocolError, LlmProviderError):
raise
except Exception as exc:
raise LlmProviderError("OpenAI-compatible LLM request failed") from exc
finally:
if stream is not None:
close = getattr(stream, "close", None)
if callable(close):
await close()

def _validate_settings(self) -> None:
if not self._settings.openai_base_url:
raise LlmProviderError("OpenAI-compatible base URL is not configured")
if not self._settings.openai_api_key:
raise LlmProviderError("OpenAI-compatible API key is not configured")
if not self._settings.openai_model:
raise LlmProviderError("OpenAI-compatible model is not configured")

@staticmethod
def _parse_usage(raw_usage: object) -> LlmUsage:
values = (
getattr(raw_usage, "prompt_tokens", None),
getattr(raw_usage, "completion_tokens", None),
getattr(raw_usage, "total_tokens", None),
)
if any(
not isinstance(value, int) or isinstance(value, bool) or value < 0 for value in values
):
raise LlmProtocolError("LLM usage fields must be non-negative integers")
prompt_tokens, completion_tokens, total_tokens = cast(tuple[int, int, int], values)
return LlmUsage(prompt_tokens, completion_tokens, total_tokens)

@staticmethod
def _parse_tool_call_delta(call: object) -> ToolCallDelta:
index = getattr(call, "index", None)
call_id = getattr(call, "id", None)
function = getattr(call, "function", None)
if function is None:
raise LlmProtocolError("LLM tool call function is missing")
name = getattr(function, "name", None)
arguments = getattr(function, "arguments", None)
if not isinstance(index, int) or isinstance(index, bool):
raise LlmProtocolError("LLM tool call index must be an integer")
if call_id is not None and not isinstance(call_id, str):
raise LlmProtocolError("LLM tool call ID must be a string")
if name is not None and not isinstance(name, str):
raise LlmProtocolError("LLM tool call name must be a string")
if not isinstance(arguments, str):
raise LlmProtocolError("LLM tool call arguments must be a string")
return ToolCallDelta(index, call_id or None, name or None, arguments)
17 changes: 17 additions & 0 deletions backend/src/timeflow/infrastructure/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,11 @@ class Settings:
aliyun_asr_vad_silence_duration_ms: int = 400
aliyun_asr_connect_timeout_seconds: float = 10.0
aliyun_asr_finish_timeout_seconds: float = 10.0
openai_base_url: str = ""
openai_api_key: str = ""
openai_model: str = "qwen-flash"
openai_timeout_seconds: float = 30.0
agent_max_tool_rounds: int = 4

@classmethod
def from_environment(cls, env_file: Path | str = ".env") -> "Settings":
Expand All @@ -44,6 +49,9 @@ def from_environment(cls, env_file: Path | str = ".env") -> "Settings":
environ.get("TIMEFLOW_ALIYUN_ASR_FINISH_TIMEOUT_SECONDS", "10.0")
)

openai_timeout_seconds = float(environ.get("TIMEFLOW_OPENAI_TIMEOUT_SECONDS", "30.0"))
agent_max_tool_rounds = int(environ.get("TIMEFLOW_AGENT_MAX_TOOL_ROUNDS", "4"))

if not -1.0 <= aliyun_asr_vad_threshold <= 1.0:
raise ValueError("TIMEFLOW_ALIYUN_ASR_VAD_THRESHOLD must be between -1 and 1")
if not 200 <= aliyun_asr_vad_silence_duration_ms <= 6000:
Expand All @@ -52,6 +60,10 @@ def from_environment(cls, env_file: Path | str = ".env") -> "Settings":
)
if aliyun_asr_connect_timeout_seconds <= 0 or aliyun_asr_finish_timeout_seconds <= 0:
raise ValueError("ASR timeouts must be greater than zero")
if openai_timeout_seconds <= 0:
raise ValueError("TIMEFLOW_OPENAI_TIMEOUT_SECONDS must be greater than zero")
if agent_max_tool_rounds <= 0:
raise ValueError("TIMEFLOW_AGENT_MAX_TOOL_ROUNDS must be a positive integer")

return cls(
app_name=environ.get("TIMEFLOW_APP_NAME", "TimeFlow API"),
Expand Down Expand Up @@ -81,6 +93,11 @@ def from_environment(cls, env_file: Path | str = ".env") -> "Settings":
aliyun_asr_vad_silence_duration_ms=aliyun_asr_vad_silence_duration_ms,
aliyun_asr_connect_timeout_seconds=aliyun_asr_connect_timeout_seconds,
aliyun_asr_finish_timeout_seconds=aliyun_asr_finish_timeout_seconds,
openai_base_url=environ.get("TIMEFLOW_OPENAI_BASE_URL", ""),
openai_api_key=environ.get("TIMEFLOW_OPENAI_API_KEY", ""),
openai_model=environ.get("TIMEFLOW_OPENAI_MODEL", "qwen-flash"),
openai_timeout_seconds=openai_timeout_seconds,
agent_max_tool_rounds=agent_max_tool_rounds,
)


Expand Down
68 changes: 68 additions & 0 deletions backend/src/timeflow/intelligence/conversation/__init__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,19 @@
"""Provider-neutral conversation interfaces."""

from timeflow.intelligence.conversation.agent import (
Agent,
AgentCompleted,
AgentConversation,
AgentError,
AgentEvent,
AgentProtocolError,
AgentQuestion,
AgentTextDelta,
AgentToolError,
AgentToolRoundLimitError,
PendingQuestion,
QuestionKind,
)
from timeflow.intelligence.conversation.asr import (
AsrConnectionError,
AsrError,
Expand All @@ -10,14 +24,68 @@
TranscriptCompleted,
TranscriptPreview,
)
from timeflow.intelligence.conversation.llm import (
AssistantToolCallMessage,
ChatMessage,
LlmError,
LlmEvent,
LlmMessage,
LlmPort,
LlmProtocolError,
LlmProviderError,
LlmStreamCompleted,
LlmUsage,
TextDelta,
ToolCall,
ToolCallDelta,
ToolDefinition,
ToolResultMessage,
)
from timeflow.intelligence.conversation.tools import (
Tool,
ToolRegistry,
build_agent_tool_registry,
request_user_input_definition,
)

__all__ = [
"Agent",
"AgentCompleted",
"AgentConversation",
"AgentError",
"AgentEvent",
"AgentProtocolError",
"AgentQuestion",
"AgentTextDelta",
"AgentToolError",
"AgentToolRoundLimitError",
"AsrConnectionError",
"AsrError",
"AsrEvent",
"AsrPort",
"AsrProtocolError",
"AsrTranscriptionError",
"AssistantToolCallMessage",
"ChatMessage",
"LlmError",
"LlmEvent",
"LlmMessage",
"LlmPort",
"LlmProtocolError",
"LlmProviderError",
"LlmStreamCompleted",
"LlmUsage",
"PendingQuestion",
"QuestionKind",
"TextDelta",
"Tool",
"ToolCall",
"ToolCallDelta",
"ToolDefinition",
"ToolRegistry",
"ToolResultMessage",
"TranscriptCompleted",
"TranscriptPreview",
"build_agent_tool_registry",
"request_user_input_definition",
]
Loading
Loading