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
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -257,7 +257,7 @@ line-length = 88
extend-exclude = ["tests/profiles/syntax_error.py"]

[tool.ruff.lint]
extend-select = ["C", "FA", "FLY", "I", "ISC", "TID251", "UP006", "UP007", "UP008", "UP010", "UP012", "UP017", "UP024", "UP035", "UP040", "UP041", "RUF100", "B009", "B010", "B011", "BLE", "DTZ005", "D202", "PLR1733", "RUF022", "PLW1510", "TC004", "PIE790", "PERF402", "FURB129", "RET501", "LOG"]
extend-select = ["C", "FA", "FLY", "I", "ISC", "TID251", "UP006", "UP007", "UP008", "UP010", "UP012", "UP017", "UP024", "UP035", "UP040", "UP041", "RUF100", "B009", "B010", "B011", "BLE", "DTZ005", "D202", "PLR1733", "RUF022", "PLW1510", "TC004", "PIE790", "PERF402", "FURB129", "RET501", "LOG", "D413"]
ignore = ["UP047", "UP045", "S", "C901", "RUF", "SIM", "B017", "TRY004", "TRY201", "TRY203", "TRY401", "B008", "EXE001", "G201", "PYI034", "PYI064"]

[tool.ruff.lint.flake8-tidy-imports.banned-api]
Expand Down
4 changes: 4 additions & 0 deletions src/a2a_storage/context_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ async def get(self, context_id: str) -> Optional[str]:

Returns:
The OGX conversation ID, or None if not found.

"""

@abstractmethod
Expand All @@ -32,6 +33,7 @@ async def set(self, context_id: str, conversation_id: str) -> None:
Args:
context_id: The A2A context ID.
conversation_id: The OGX conversation ID.

"""

@abstractmethod
Expand All @@ -40,6 +42,7 @@ async def delete(self, context_id: str) -> None:

Args:
context_id: The A2A context ID to delete.

"""

@abstractmethod
Expand All @@ -55,4 +58,5 @@ def ready(self) -> bool:

Returns:
True if the store is initialized and ready, False otherwise.

"""
4 changes: 4 additions & 0 deletions src/a2a_storage/in_memory_context_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ async def get(self, context_id: str) -> Optional[str]:

Returns:
The OGX conversation ID, or None if not found.

"""
async with self._lock:
conversation_id = self._contexts.get(context_id)
Expand All @@ -52,6 +53,7 @@ async def set(self, context_id: str, conversation_id: str) -> None:
Args:
context_id: The A2A context ID.
conversation_id: The OGX conversation ID.

"""
async with self._lock:
self._contexts[context_id] = conversation_id
Expand All @@ -66,6 +68,7 @@ async def delete(self, context_id: str) -> None:

Args:
context_id: The A2A context ID to delete.

"""
async with self._lock:
if context_id in self._contexts:
Expand All @@ -89,5 +92,6 @@ def ready(self) -> bool:

Returns:
True, as in-memory store is always ready after construction.

"""
return self._initialized
5 changes: 5 additions & 0 deletions src/a2a_storage/postgres_context_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ def __init__(
Args:
engine: SQLAlchemy async engine connected to the PostgreSQL database.
create_table: If True, create the table on initialization.

"""
logger.debug("Initializing PostgresA2AContextStore")
self._engine = engine
Expand Down Expand Up @@ -75,6 +76,7 @@ async def get(self, context_id: str) -> Optional[str]:

Returns:
The OGX conversation ID, or None if not found.

"""
await self._ensure_initialized()

Expand All @@ -99,6 +101,7 @@ async def set(self, context_id: str, conversation_id: str) -> None:
Args:
context_id: The A2A context ID.
conversation_id: The OGX conversation ID.

"""
await self._ensure_initialized()

Expand All @@ -124,6 +127,7 @@ async def delete(self, context_id: str) -> None:

Args:
context_id: The A2A context ID to delete.

"""
await self._ensure_initialized()

Expand All @@ -139,5 +143,6 @@ def ready(self) -> bool:

Returns:
True if the store is initialized, False otherwise.

"""
return self._initialized
5 changes: 5 additions & 0 deletions src/a2a_storage/sqlite_context_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ def __init__(
Args:
engine: SQLAlchemy async engine connected to the SQLite database.
create_table: If True, create the table on initialization.

"""
logger.debug("Initializing SQLiteA2AContextStore")
self._engine = engine
Expand Down Expand Up @@ -74,6 +75,7 @@ async def get(self, context_id: str) -> Optional[str]:

Returns:
The OGX conversation ID, or None if not found.

"""
await self._ensure_initialized()

Expand All @@ -98,6 +100,7 @@ async def set(self, context_id: str, conversation_id: str) -> None:
Args:
context_id: The A2A context ID.
conversation_id: The OGX conversation ID.

"""
await self._ensure_initialized()

Expand Down Expand Up @@ -125,6 +128,7 @@ async def delete(self, context_id: str) -> None:

Args:
context_id: The A2A context ID to delete.

"""
await self._ensure_initialized()

Expand All @@ -140,5 +144,6 @@ def ready(self) -> bool:

Returns:
True if the store is initialized, False otherwise.

"""
return self._initialized
3 changes: 3 additions & 0 deletions src/a2a_storage/storage_factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ async def create_task_store(cls, config: A2AStateConfiguration) -> TaskStore:

Returns:
TaskStore implementation (InMemoryTaskStore or DatabaseTaskStore).

"""
if cls._task_store is not None:
return cls._task_store
Expand Down Expand Up @@ -86,6 +87,7 @@ async def create_context_store(

Returns:
A2AContextStore implementation.

"""
if cls._context_store is not None:
return cls._context_store
Expand Down Expand Up @@ -131,6 +133,7 @@ async def _get_or_create_engine(cls, config: A2AStateConfiguration) -> AsyncEngi

Returns:
SQLAlchemy AsyncEngine.

"""
if cls._engine is not None:
return cls._engine
Expand Down
5 changes: 5 additions & 0 deletions src/app/database.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ def get_engine() -> Engine:
Raises:
RuntimeError: If the database engine has not been initialized; call
initialize_database() first.

"""
if engine is None:
raise RuntimeError(
Expand All @@ -45,6 +46,7 @@ def create_tables() -> None:
Raises:
RuntimeError: If the global database engine is not initialized (call
initialize_database() first).

"""
Base.metadata.create_all(get_engine())

Expand All @@ -60,6 +62,7 @@ def get_session() -> Session:
Raises:
RuntimeError: If the database has not been initialized; call
initialize_database() first.

"""
if session_local is None:
raise RuntimeError(
Expand All @@ -85,6 +88,7 @@ def _create_sqlite_engine(config: SQLiteDatabaseConfiguration, **kwargs: Any) ->
------
FileNotFoundError: If the parent directory of `config.db_path` does not exist.
RuntimeError: If engine creation fails.

"""
if not Path(config.db_path).parent.exists():
raise FileNotFoundError(
Expand Down Expand Up @@ -122,6 +126,7 @@ def _create_postgres_engine(
Raises:
------
RuntimeError: If engine creation fails or if creating the specified schema fails.

"""
postgres_url = (
f"postgresql://{config.user}:{config.password.get_secret_value()}@"
Expand Down
17 changes: 17 additions & 0 deletions src/app/endpoints/a2a.py
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ async def _get_task_store() -> TaskStore:

Returns:
TaskStore instance based on configuration.

"""
global _TASK_STORE # pylint: disable=global-statement
if _TASK_STORE is None:
Expand All @@ -116,6 +117,7 @@ async def _get_context_store() -> A2AContextStore:

Returns:
A2AContextStore instance based on configuration.

"""
global _CONTEXT_STORE # pylint: disable=global-statement
if _CONTEXT_STORE is None:
Expand All @@ -141,6 +143,7 @@ def _build_a2a_parts_from_agent_result(

Returns:
List of A2A Part objects for the final artifact.

"""
if run_result is not None:
final_text = run_result.response.text or "".join(accumulated_text)
Expand All @@ -157,6 +160,7 @@ def _record_model_span(span: trace.Span, model_id: str) -> None:
Parameters:
span: The active OpenTelemetry span.
model_id: Full model identifier in "provider/model" format.

"""
provider_id, bare_model_id = extract_provider_and_model_from_model_id(model_id)
set_span_attributes(
Expand All @@ -183,6 +187,7 @@ def _record_execution_span(
run_result: Completed agent run result, or None.
compacted: Whether the turn used compacted conversation context.
inference_time: Request processing duration in seconds.

"""
if tool_call_names:
set_span_attributes(
Expand Down Expand Up @@ -234,6 +239,7 @@ async def _persist_compacted_a2a_turn(
the request was served in compacted mode.
agent: The pydantic-ai agent whose model captured the output items.
task_id: A2A task identifier, used for error reporting.

"""
if not compaction.compacted or compaction.original_input is None:
return
Expand Down Expand Up @@ -275,6 +281,7 @@ def process_event(

Args:
event: The event to process

"""
if isinstance(event, TaskStatusUpdateEvent):
if event.status.state == TaskState.failed:
Expand Down Expand Up @@ -335,6 +342,7 @@ def __init__(
auth_token: Authentication token for the request
mcp_headers: MCP headers for context propagation
request_headers: Incoming HTTP request headers for allowlist propagation

"""
self.auth_token: str = auth_token
self.mcp_headers: McpHeaders = mcp_headers or {}
Expand All @@ -352,6 +360,7 @@ async def execute(
Args:
context: The request context containing user input and metadata
event_queue: Queue for sending response events

"""
if not context.message:
raise ValueError("A2A request must have a message")
Expand Down Expand Up @@ -406,6 +415,7 @@ async def _process_task_streaming( # pylint: disable=too-many-locals,too-many-s
task_updater: Task updater for sending events
task_id: The task ID to use for this execution
context_id: The context ID to use for this execution

"""
if not task_id or not context_id:
raise ValueError("Task ID and Context ID are required")
Expand Down Expand Up @@ -628,6 +638,7 @@ async def _convert_stream_to_events(

Yields:
A2A events (TaskStatusUpdateEvent or TaskArtifactUpdateEvent)

"""
if not task_id or not context_id:
raise ValueError("Task ID and Context ID are required")
Expand Down Expand Up @@ -686,6 +697,7 @@ def _dispatch_agent_event( # pylint: disable=too-many-arguments,too-many-positi

Returns:
A2A TaskStatusUpdateEvent, or None if the event is not mapped.

"""
_ = artifact_id

Expand Down Expand Up @@ -752,6 +764,7 @@ def _text_status_event(

Returns:
TaskStatusUpdateEvent with the text delta.

"""
return TaskStatusUpdateEvent(
task_id=task_id,
Expand Down Expand Up @@ -781,6 +794,7 @@ async def cancel(

Raises:
NotImplementedError: Task cancellation is not currently supported

"""
logger.info("Cancellation requested but not currently supported")
raise NotImplementedError("Task cancellation not currently supported")
Expand All @@ -798,6 +812,7 @@ def get_lightspeed_agent_card() -> AgentCard:

Returns:
AgentCard: The agent card describing Lightspeed's capabilities.

"""
# Get base URL from configuration or construct it
service_config = configuration.service_configuration
Expand Down Expand Up @@ -922,6 +937,7 @@ async def _create_a2a_app(

Returns:
A2A Starlette ASGI application

"""
agent_executor = A2AAgentExecutor(
auth_token=auth_token,
Expand Down Expand Up @@ -1061,6 +1077,7 @@ async def _handle_a2a_jsonrpc( # pylint: disable=too-many-locals,too-many-state

Returns:
JSON-RPC response or streaming response

"""
logger.debug("A2A endpoint called: %s %s", request.method, request.url.path)

Expand Down
3 changes: 3 additions & 0 deletions src/app/endpoints/conversations_v1.py
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,7 @@ async def get_conversation_endpoint_handler( # pylint: disable=too-many-locals,
Returns:
ConversationResponse: Structured response containing the conversation
ID and simplified chat history

"""
with tracer.start_as_current_span("conversations_v1.get") as span:
check_configuration_loaded(configuration)
Expand Down Expand Up @@ -356,6 +357,7 @@ async def delete_conversation_endpoint_handler(

Returns:
ConversationDeleteResponse: Response indicating the result of the deletion operation

"""
with tracer.start_as_current_span("conversations_v1.delete") as span:
check_configuration_loaded(configuration)
Expand Down Expand Up @@ -477,6 +479,7 @@ async def update_conversation_endpoint_handler( # pylint: disable=too-many-stat

Returns:
ConversationUpdateResponse: Response indicating the result of the update operation

"""
with tracer.start_as_current_span("conversations_v1.update") as span:
check_configuration_loaded(configuration)
Expand Down
1 change: 1 addition & 0 deletions src/app/endpoints/conversations_v2.py
Original file line number Diff line number Diff line change
Expand Up @@ -275,6 +275,7 @@ def build_conversation_turn_from_cache_entry(entry: CacheEntry) -> ConversationT

Returns:
ConversationTurn object with messages, tool_calls, tool_results, and timestamps

"""
# Create Message objects for user and assistant
messages = [
Expand Down
Loading
Loading