diff --git a/service/Makefile b/service/Makefile index e8aeb2e..ad7fd64 100644 --- a/service/Makefile +++ b/service/Makefile @@ -55,3 +55,21 @@ db-current: .PHONY: db-history db-history: uv run --active alembic history + +# Both targets create versions/YYYYMMDD_NN_.py with revision ID YYYYMMDD_NN (next free NN for today). +MIGRATIONS_VERSIONS_DIR := src/ai_document_plugin_service/migrations/versions +NEXT_REVISION_ID = $(shell day=$$(date +%Y%m%d); \ + count=$$(find $(MIGRATIONS_VERSIONS_DIR) -maxdepth 1 -name "$${day}_*.py" | wc -l); \ + printf '%s_%02d' "$$day" $$((count + 1))) + +# Usage: make db-revision m="Add foo to bar" (empty migration) +.PHONY: db-revision +db-revision: + @test -n "$(m)" || { echo 'Usage: make db-revision m="Describe the change"'; exit 1; } + uv run --active alembic revision --rev-id "$(NEXT_REVISION_ID)" -m "$(m)" + +# Usage: make db-autogenerate m="Add foo to bar" (diffs schema.py against the DB; DB must be at head) +.PHONY: db-autogenerate +db-autogenerate: + @test -n "$(m)" || { echo 'Usage: make db-autogenerate m="Describe the change"'; exit 1; } + uv run --active alembic revision --autogenerate --rev-id "$(NEXT_REVISION_ID)" -m "$(m)" diff --git a/service/README.md b/service/README.md index d74665f..6a6a266 100644 --- a/service/README.md +++ b/service/README.md @@ -66,6 +66,9 @@ The project `Makefile` provides a few shortcuts for common development tasks: - `make db-migrate` applies all Alembic migrations to the configured database - `make db-current` shows the current Alembic revision stored in the database - `make db-history` shows available Alembic migration history +- `make db-revision m="Add foo to bar"` creates a new migration `YYYYMMDD_NN_add_foo_to_bar.py` with revision ID + `YYYYMMDD_NN`. Migrations are upgrade-only: there are no `downgrade()` functions, so undo a change with a new + migration Run these commands from the [service](/Users/hana/DSW/AI-playground/ai-document-plugin/service:1) directory. diff --git a/service/src/ai_document_plugin_service/ai/common/__init__.py b/service/src/ai_document_plugin_service/ai/common/__init__.py index 8db421b..e5b729f 100644 --- a/service/src/ai_document_plugin_service/ai/common/__init__.py +++ b/service/src/ai_document_plugin_service/ai/common/__init__.py @@ -3,6 +3,8 @@ from .logging_utils import configure_logging from .pipeline_metrics import ( PipelineMetricsCollector, + PipelineStats, + StepUsage, get_component_markdown, get_component_stats, ) @@ -12,6 +14,8 @@ 'AssignmentStats', 'Config', 'PipelineMetricsCollector', + 'PipelineStats', + 'StepUsage', 'call_with_retry', 'configure_logging', 'extract_usage_tokens', diff --git a/service/src/ai_document_plugin_service/ai/common/pipeline_metrics.py b/service/src/ai_document_plugin_service/ai/common/pipeline_metrics.py index 45dc351..00cc3a6 100644 --- a/service/src/ai_document_plugin_service/ai/common/pipeline_metrics.py +++ b/service/src/ai_document_plugin_service/ai/common/pipeline_metrics.py @@ -3,7 +3,6 @@ from dataclasses import dataclass, field from ai_document_plugin_service.ai.common.types import AssignmentStats -from ai_document_plugin_service.ai.persistence.assignment_saver_component import JsonValue @dataclass(frozen=True) @@ -12,11 +11,37 @@ class PipelineMetricStep: stats: AssignmentStats +@dataclass(frozen=True) +class StepUsage: + """LLM usage of one pipeline step in a single run.""" + + llm_calls: int + input_tokens: int + output_tokens: int + + @classmethod + def from_stats(cls, stats: AssignmentStats | None) -> 'StepUsage | None': + if stats is None: + return None + return cls( + llm_calls=stats.total_calls, + input_tokens=stats.total_input_tokens, + output_tokens=stats.total_output_tokens, + ) + + +@dataclass(frozen=True) +class PipelineStats: + """Per-step LLM usage for a single run. A step is None when it did not run (e.g. cached assignments).""" + + assignment: StepUsage | None + generation: StepUsage | None + polishing: StepUsage | None + elapsed_seconds: float + + @dataclass class PipelineMetricsCollector: - model_name: str - cost_per_mil_input: float - cost_per_mil_output: float steps: list[PipelineMetricStep] = field(default_factory=list) def add_step(self, step_name: str, stats: AssignmentStats | None) -> None: @@ -24,33 +49,32 @@ def add_step(self, step_name: str, stats: AssignmentStats | None) -> None: return self.steps.append(PipelineMetricStep(name=step_name, stats=stats)) - def get_stats(self, elapsed_seconds: float) -> JsonValue: - return self._build_summary_section(elapsed_seconds) - def log_summary(self, logger: logging.Logger) -> None: if not self.steps: logger.debug('No pipeline metrics were collected.') return - logger.debug('Token usage and cost:') + logger.debug('Token usage:') for step in self.steps: - _, _, total_cost = self._price(step.stats) logger.debug( - '%s: %s calls, %s in / %s out tokens, %.2f USD', + '%s: %s calls, %s in / %s out tokens', step.name, f'{step.stats.total_calls:,}', f'{step.stats.total_input_tokens:,}', f'{step.stats.total_output_tokens:,}', - total_cost, ) logger.debug( - 'Total: %s in / %s out tokens, %.2f USD', + 'Total: %s calls, %s in / %s out tokens', + f'{self.total_llm_calls:,}', f'{self.total_input_tokens:,}', f'{self.total_output_tokens:,}', - self.total_cost, ) + @property + def total_llm_calls(self) -> int: + return sum(step.stats.total_calls for step in self.steps) + @property def total_input_tokens(self) -> int: return sum(step.stats.total_input_tokens for step in self.steps) @@ -59,58 +83,6 @@ def total_input_tokens(self) -> int: def total_output_tokens(self) -> int: return sum(step.stats.total_output_tokens for step in self.steps) - @property - def total_cost(self) -> float: - return sum(self._price(step.stats)[2] for step in self.steps) - - @property - def total_llm_wait_ms(self) -> float: - return round(sum(step.stats.total_llm_wait_ms for step in self.steps), 3) - - @property - def total_llm_response_ms(self) -> float: - return round(sum(step.stats.total_llm_response_ms for step in self.steps), 3) - - def _price(self, stats: AssignmentStats) -> tuple[float, float, float]: - input_cost = stats.total_input_tokens * self.cost_per_mil_input / 1_000_000 - output_cost = stats.total_output_tokens * self.cost_per_mil_output / 1_000_000 - return input_cost, output_cost, input_cost + output_cost - - def _build_summary_section(self, elapsed_seconds: float) -> JsonValue: - return { - 'title': 'Pipeline token usage and cost', - 'headers': [ - 'Step', - 'LLM calls', - 'Input tokens', - 'Output tokens', - 'Cost (USD)', - ], - 'rows': [ - { - 'step': step.name, - 'llm_calls': step.stats.total_calls, - 'input_tokens': step.stats.total_input_tokens, - 'output_tokens': step.stats.total_output_tokens, - 'cost_usd': round(self._price(step.stats)[2], 2), - } - for step in self.steps - ], - 'totals': { - 'input_tokens': self.total_input_tokens, - 'output_tokens': self.total_output_tokens, - 'cost_usd': round(self.total_cost, 2), - }, - 'meta': { - 'model_name': self.model_name, - 'cost_per_mil_input': self.cost_per_mil_input, - 'cost_per_mil_output': self.cost_per_mil_output, - 'elapsed_seconds': elapsed_seconds, - 'total_llm_wait_ms': self.total_llm_wait_ms, - 'total_llm_response_ms': self.total_llm_response_ms, - }, - } - def _get_component_dict( pipeline_result: Mapping[str, object], diff --git a/service/src/ai_document_plugin_service/ai/persistence/database.py b/service/src/ai_document_plugin_service/ai/persistence/database.py index d5f91b0..f329f2d 100644 --- a/service/src/ai_document_plugin_service/ai/persistence/database.py +++ b/service/src/ai_document_plugin_service/ai/persistence/database.py @@ -8,7 +8,7 @@ from typing import Any, TypedDict, Unpack from uuid import UUID, uuid4 -from sqlalchemy import ColumnElement, Connection, Row, and_, inspect, or_ +from sqlalchemy import Column, ColumnElement, Connection, Row, and_, inspect, or_, select from sqlalchemy.dialects.postgresql import insert as postgresql_insert from sqlalchemy.engine import URL from sqlalchemy.exc import IntegrityError @@ -63,6 +63,21 @@ class GenerationUpdate(TypedDict, total=False): progress_message: str | None +class GenerationStats(TypedDict, total=False): + dmp_pre_polished: str | None + dmp_polished: str | None + assignment_llm_calls: int | None + assignment_input_tokens: int | None + assignment_output_tokens: int | None + generation_llm_calls: int | None + generation_input_tokens: int | None + generation_output_tokens: int | None + polishing_llm_calls: int | None + polishing_input_tokens: int | None + polishing_output_tokens: int | None + elapsed_seconds: float | None + + @dataclass(frozen=True) class GenerationRecord: """A raw generation (pipeline run) row.""" @@ -199,40 +214,6 @@ async def get_template( UPDATE) so it cannot change between this read and a follow-up write. """ - @abstractmethod - async def save_result( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - prepolished_markdown: str, - markdown: str, - ) -> None: - """Persist a markdown result in a database backend.""" - - @abstractmethod - async def save_stats( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - stats: JsonValue, - ) -> None: - """Persist a stats result in a database backend.""" - - @abstractmethod - async def update_result( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - markdown: str, - ) -> None: - """Persist a markdown result in a database backend.""" - @abstractmethod async def create_generation( self, @@ -260,6 +241,15 @@ async def update_generation( ``run_id``/``tenant_uuid``. """ + @abstractmethod + async def create_generation_stats( + self, + run_id: UUID, + trace_id: UUID | None, + **stats: Unpack[GenerationStats], + ) -> None: + """Create the stats row of a finished generation.""" + @abstractmethod async def get_generation( self, @@ -298,8 +288,8 @@ def __init__( self.metadata = schema.metadata self.assignment_table = schema.assignment_table self.template_table = schema.template_table - self.result_table = schema.result_table self.generation_table = schema.generation_table + self.generation_stats_table = schema.generation_stats_table self._database_verified = False logger.info( 'Initialized Postgres database client', @@ -348,6 +338,10 @@ async def _connect(self) -> AsyncIterator[AsyncConnection]: async with self.engine.begin() as connection: yield connection + def _generation_record_columns(self) -> list[Column[Any]]: + """Generation columns read into ``GenerationRecord``.""" + return [column for column in self.generation_table.c if column.name in GenerationRecord.__dataclass_fields__] + def _list_existing_tables(self, connection: Connection) -> set[str]: return set(inspect(connection).get_table_names(schema=self.schema_name)) @@ -359,7 +353,7 @@ async def _ensure_schema(self) -> None: async with self.engine.connect() as connection: existing_tables = await connection.run_sync(self._list_existing_tables) - required_tables = {'alembic_version', 'template', 'assignment', 'result', 'generation'} + required_tables = {'alembic_version', 'template', 'assignment', 'generation', 'generation_stats'} missing_tables = sorted(required_tables - existing_tables) if missing_tables: @@ -652,158 +646,6 @@ async def get_template( return None return TemplateRecord.from_row(row) - async def save_result( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - prepolished_markdown: str, - markdown: str, - ) -> None: - await self._ensure_schema() - now = datetime.now(tz=UTC) - - statement = postgresql_insert(self.result_table).values( - template_uuid=template_uuid, - knowledge_model_uuid=knowledge_model_uuid, - user_uuid=user_uuid, - tenant_uuid=tenant_uuid, - dmp_pre_polished=prepolished_markdown, - dmp=markdown, - created_at=now, - updated_at=now, - ) - - upsert_statement = statement.on_conflict_do_update( - constraint='pk_result', - set_={ - 'dmp_pre_polished': statement.excluded.dmp_pre_polished, - 'dmp': statement.excluded.dmp, - 'updated_at': statement.excluded.updated_at, - }, - ) - - async with self._connect() as connection: - await connection.execute(upsert_statement) - - logger.info( - 'Saved pipeline result', - extra={ - 'knowledge_model_uuid': knowledge_model_uuid, - 'template_uuid': str(template_uuid), - 'user_uuid': str(user_uuid), - 'tenant_uuid': str(tenant_uuid), - 'prepolished_markdown_length': len(prepolished_markdown), - 'markdown_length': len(markdown), - 'db.schema': self.schema_name, - }, - ) - - async def save_stats( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - stats: JsonValue, - ) -> None: - await self._ensure_schema() - now = datetime.now(tz=UTC) - - statement = ( - self.result_table.update() - .where( - (self.result_table.c.knowledge_model_uuid == knowledge_model_uuid) - & (self.result_table.c.template_uuid == template_uuid) - & (self.result_table.c.user_uuid == user_uuid) - & (self.result_table.c.tenant_uuid == tenant_uuid) - ) - .values(stats=stats, updated_at=now) - ) - - async with self._connect() as connection: - result = await connection.execute(statement) - - if result.rowcount == 0: - msg = 'Cannot save stats because result row does not exist yet. Save dmp and dmp_pre_polished first.' - logger.error( - 'Stats update failed because result row does not exist', - extra={ - 'knowledge_model_uuid': str(knowledge_model_uuid), - 'template_uuid': str(template_uuid), - 'user_uuid': str(user_uuid), - 'tenant_uuid': str(tenant_uuid), - 'db.schema': self.schema_name, - }, - ) - raise ValueError(msg) - - logger.info( - 'Saved pipeline stats', - extra={ - 'knowledge_model_uuid': str(knowledge_model_uuid), - 'template_uuid': str(template_uuid), - 'user_uuid': str(user_uuid), - 'tenant_uuid': str(tenant_uuid), - 'db.schema': self.schema_name, - }, - ) - - async def update_result( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - markdown: str, - ) -> None: - await self._ensure_schema() - now = datetime.now(tz=UTC) - - statement = ( - self.result_table.update() - .where( - (self.result_table.c.knowledge_model_uuid == knowledge_model_uuid) - & (self.result_table.c.template_uuid == template_uuid) - & (self.result_table.c.user_uuid == user_uuid) - & (self.result_table.c.tenant_uuid == tenant_uuid) - ) - .values( - dmp=markdown, - updated_at=now, - ) - ) - - async with self._connect() as connection: - result = await connection.execute(statement) - - if result.rowcount == 0: - msg = 'Cannot save result because result row does not exist yet. Create the row first before updating dmp.' - logger.error( - 'Result update failed because result row does not exist', - extra={ - 'knowledge_model_uuid': str(knowledge_model_uuid), - 'template_uuid': str(template_uuid), - 'user_uuid': str(user_uuid), - 'tenant_uuid': str(tenant_uuid), - 'db.schema': self.schema_name, - }, - ) - raise ValueError(msg) - - logger.info( - 'Updated stored pipeline result markdown', - extra={ - 'knowledge_model_uuid': str(knowledge_model_uuid), - 'template_uuid': str(template_uuid), - 'user_uuid': str(user_uuid), - 'tenant_uuid': str(tenant_uuid), - 'markdown_length': len(markdown), - 'db.schema': self.schema_name, - }, - ) - async def create_generation( self, questionnaire_uuid: UUID, @@ -855,7 +697,7 @@ async def update_generation( self.generation_table.update() .where((self.generation_table.c.run_id == run_id) & (self.generation_table.c.tenant_uuid == tenant_uuid)) .values(**updates, updated_at=now) - .returning(*self.generation_table.c) + .returning(*self._generation_record_columns()) ) async with self._connect() as connection: @@ -872,6 +714,18 @@ async def update_generation( ) return GenerationRecord.from_row(row) + async def create_generation_stats( + self, + run_id: UUID, + trace_id: UUID | None, + **stats: Unpack[GenerationStats], + ) -> None: + await self._ensure_schema() + statement = self.generation_stats_table.insert().values(run_id=run_id, trace_id=trace_id, **stats) + + async with self._connect() as connection: + await connection.execute(statement) + async def get_generation( self, run_id: UUID, @@ -879,7 +733,7 @@ async def get_generation( user_uuid: UUID, ) -> GenerationRecord | None: await self._ensure_schema() - statement = self.generation_table.select().where( + statement = select(*self._generation_record_columns()).where( and_( self.generation_table.c.run_id == run_id, self.generation_table.c.tenant_uuid == tenant_uuid, @@ -901,7 +755,7 @@ async def list_generations( ) -> list[GenerationRecord]: await self._ensure_schema() statement = ( - self.generation_table.select() + select(*self._generation_record_columns()) .where( and_( self.generation_table.c.questionnaire_uuid == questionnaire_uuid, diff --git a/service/src/ai_document_plugin_service/ai/persistence/saver_component.py b/service/src/ai_document_plugin_service/ai/persistence/saver_component.py index 8e20541..e69de29 100644 --- a/service/src/ai_document_plugin_service/ai/persistence/saver_component.py +++ b/service/src/ai_document_plugin_service/ai/persistence/saver_component.py @@ -1,78 +0,0 @@ -import logging -from typing import TypedDict -from uuid import UUID - -from haystack import component - -from ai_document_plugin_service.ai.persistence.database import Database - -logger = logging.getLogger(__name__) - - -class FileSaverComponentResult(TypedDict): - markdown: str - - -@component -class SaverComponent: - def __init__(self, database: Database) -> None: - self.database = database - - @component.output_types(markdown=str) - async def run_async( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - debug_markdown: str, - markdown: str, - ) -> FileSaverComponentResult: - logger.info( - 'Persisting generated markdown result', - extra={ - 'template_uuid': str(template_uuid), - 'knowledge_model_uuid': knowledge_model_uuid, - 'user_uuid': str(user_uuid), - 'tenant_uuid': str(tenant_uuid), - 'debug_markdown_length': len(debug_markdown), - 'markdown_length': len(markdown), - }, - ) - await self.database.save_result( - template_uuid, - knowledge_model_uuid, - user_uuid, - tenant_uuid, - debug_markdown, - markdown, - ) - logger.info( - 'Generated markdown result persisted', - extra={ - 'template_uuid': str(template_uuid), - 'knowledge_model_uuid': knowledge_model_uuid, - 'user_uuid': str(user_uuid), - 'tenant_uuid': str(tenant_uuid), - }, - ) - - return { - 'markdown': markdown, - } - - @component.output_types(markdown=str) - def run( - self, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - debug_markdown: str, - markdown: str, - ) -> FileSaverComponentResult: - """Async-only component; the sync pipeline entrypoint is intentionally unsupported.""" - msg = f'{type(self).__name__} is async-only; use run_async() / AsyncPipeline.run_async()' - raise NotImplementedError( - msg, - ) diff --git a/service/src/ai_document_plugin_service/ai/persistence/schema.py b/service/src/ai_document_plugin_service/ai/persistence/schema.py index 398a169..95b55f2 100644 --- a/service/src/ai_document_plugin_service/ai/persistence/schema.py +++ b/service/src/ai_document_plugin_service/ai/persistence/schema.py @@ -4,9 +4,10 @@ JSON, Column, DateTime, + Float, ForeignKey, - ForeignKeyConstraint, Index, + Integer, MetaData, Table, Text, @@ -21,8 +22,8 @@ class PersistenceSchema: metadata: MetaData assignment_table: Table template_table: Table - result_table: Table generation_table: Table + generation_stats_table: Table def create_persistence_schema(schema_name: str) -> PersistenceSchema: @@ -77,46 +78,6 @@ def create_persistence_schema(schema_name: str) -> PersistenceSchema: Column('template_uuid', UUID(as_uuid=True), ForeignKey('template.uuid'), primary_key=True, nullable=False), ) - result_table = Table( - 'result', - metadata, - Column('knowledge_model_uuid', UUID(as_uuid=True), primary_key=True), - Column('template_uuid', UUID(as_uuid=True), primary_key=True), - Column( - 'user_uuid', - UUID(as_uuid=True), - primary_key=True, - nullable=False, - ), - Column( - 'tenant_uuid', - UUID(as_uuid=True), - primary_key=True, - nullable=False, - ), - Column( - 'created_at', - DateTime(timezone=True), - nullable=False, - server_default=func.now(), - ), - Column( - 'updated_at', - DateTime(timezone=True), - nullable=False, - server_default=func.now(), - ), - Column('dmp', Text, nullable=False), - Column('dmp_pre_polished', Text, nullable=False), - Column('stats', JSON, nullable=True), - ForeignKeyConstraint(['template_uuid'], ['template.uuid'], name='fk_result_template_uuid'), - ForeignKeyConstraint( - ['knowledge_model_uuid', 'template_uuid'], - ['assignment.knowledge_model_uuid', 'assignment.template_uuid'], - name='fk_result_assignment', - ), - ) - generation_table = Table( 'generation', metadata, @@ -156,10 +117,46 @@ def create_persistence_schema(schema_name: str) -> PersistenceSchema: ), ) + # Analysis-only data, one row per generation (created together with it). + generation_stats_table = Table( + 'generation_stats', + metadata, + Column('run_id', UUID(as_uuid=True), ForeignKey('generation.run_id', ondelete='CASCADE'), primary_key=True), + # Request trace id; NULL for runs created before it was recorded. + Column('trace_id', UUID(as_uuid=True), nullable=True), + Column('dmp_pre_polished', Text, nullable=True), + # at start same as generation.result_markdown, but not editable. Older rows don't have the original value saved + Column('dmp_polished', Text, nullable=True), + # LLM usage per pipeline step. NULL when the step did not run (assignment is skipped + # when cached assignments are reused) or the run has no stats (failed or old runs). + Column('assignment_llm_calls', Integer, nullable=True), + Column('assignment_input_tokens', Integer, nullable=True), + Column('assignment_output_tokens', Integer, nullable=True), + Column('generation_llm_calls', Integer, nullable=True), + Column('generation_input_tokens', Integer, nullable=True), + Column('generation_output_tokens', Integer, nullable=True), + Column('polishing_llm_calls', Integer, nullable=True), + Column('polishing_input_tokens', Integer, nullable=True), + Column('polishing_output_tokens', Integer, nullable=True), + Column('elapsed_seconds', Float, nullable=True), + Column( + 'created_at', + DateTime(timezone=True), + nullable=False, + server_default=func.now(), + ), + Column( + 'updated_at', + DateTime(timezone=True), + nullable=False, + server_default=func.now(), + ), + ) + return PersistenceSchema( metadata=metadata, assignment_table=assignment_table, template_table=template_table, - result_table=result_table, generation_table=generation_table, + generation_stats_table=generation_stats_table, ) diff --git a/service/src/ai_document_plugin_service/ai/run_pipeline.py b/service/src/ai_document_plugin_service/ai/run_pipeline.py index 8e04f19..6a27493 100644 --- a/service/src/ai_document_plugin_service/ai/run_pipeline.py +++ b/service/src/ai_document_plugin_service/ai/run_pipeline.py @@ -4,6 +4,7 @@ import logging import time from collections.abc import Callable, Mapping +from dataclasses import dataclass from typing import TYPE_CHECKING from uuid import UUID @@ -14,6 +15,8 @@ from ai_document_plugin_service.ai.common import ( Config, PipelineMetricsCollector, + PipelineStats, + StepUsage, get_component_markdown, get_component_stats, ) @@ -27,7 +30,6 @@ DBSaver, SerializedSectionAssignment, ) -from ai_document_plugin_service.ai.persistence.saver_component import SaverComponent from ai_document_plugin_service.ai.polishing.dmp_polisher_component import DmpPolisherComponent from ai_document_plugin_service.ai.polishing.llm import SectionPolishingLLM @@ -38,21 +40,22 @@ from ai_document_plugin_service.ai.knowledgemodel.dsw_client import DSWClient from ai_document_plugin_service.ai.persistence.database import Database -# Cost per million tokens (USD) - adjust for your model -COST_PER_MIL_INPUT = 0.25 -COST_PER_MIL_OUTPUT = 2.0 - logger = logging.getLogger(__name__) ProgressCallback = Callable[[str], None] +@dataclass(frozen=True) +class PipelineOutput: + knowledge_model_uuid: UUID + markdown: str + # Generator output before polishing (the text fed into the polisher). + dmp_pre_polished: str + stats: PipelineStats + + def build_pipeline( - database: Database, - saver: DBSaver, - config: Config, - llm_client: LLMClient, - language: str, + database: Database, saver: DBSaver, config: Config, llm_client: LLMClient, language: str ) -> AsyncPipeline: pipeline = AsyncPipeline() loader_component = AssignmentLoaderComponent(database=database) @@ -61,7 +64,6 @@ def build_pipeline( assignment_saver_component = AssignmentSaverComponent(saver=saver) dmp_generator_component = DmpGeneratorComponent(SectionGenerationLLM(llm_client, config, language)) dmp_polisher_component = DmpPolisherComponent(SectionPolishingLLM(llm_client, config, language)) - saver_component = SaverComponent(database=database) # ROUTES routes: list[Route] = [ @@ -88,7 +90,6 @@ def build_pipeline( pipeline.add_component('assignment_saver_component', assignment_saver_component) pipeline.add_component('dmp_generator_component', dmp_generator_component) pipeline.add_component('dmp_polisher_component', dmp_polisher_component) - pipeline.add_component('saver_component', saver_component) # CONNECTIONS # loader_component -> router @@ -105,12 +106,8 @@ def build_pipeline( pipeline.connect('assignment_component.stats', 'assignment_saver_component.stats') # assignment_saver_component -> dmp_generator_component pipeline.connect('assignment_saver_component.assignments', 'dmp_generator_component.new_assignments') - # dmp_generator_component -> prepolished_saver_component - pipeline.connect('dmp_generator_component.debug_markdown', 'saver_component.debug_markdown') - # prepolisher_saver_component -> dmp_polisher_component + # dmp_generator_component -> dmp_polisher_component pipeline.connect('dmp_generator_component.markdown', 'dmp_polisher_component.markdown') - # dmp_polisher_component -> polished_saver_component - pipeline.connect('dmp_polisher_component.markdown', 'saver_component.markdown') return pipeline @@ -120,15 +117,11 @@ async def run_pipeline( template_uuid: UUID, template_title: str, template_data: Mapping[str, object], - user_uuid: UUID, tenant_uuid: UUID, pipeline: AsyncPipeline, - database: Database, dsw_client: DSWClient, - model_name: str, on_progress: ProgressCallback | None = None, -) -> tuple[UUID, str]: - t1 = time.time() +) -> PipelineOutput: pipeline_total_started = time.perf_counter() questionnaire_fetch_started = time.perf_counter() try: @@ -182,18 +175,11 @@ async def run_pipeline( 'template_data': template_data, 'on_progress': on_progress, }, - 'saver_component': { - 'template_uuid': template_uuid, - 'knowledge_model_uuid': knowledge_model_uuid, - 'user_uuid': user_uuid, - 'tenant_uuid': tenant_uuid, - }, }, include_outputs_from={ 'assignment_saver_component', 'dmp_generator_component', 'dmp_polisher_component', - 'saver_component', }, ) except Exception: @@ -206,10 +192,18 @@ async def run_pipeline( 'pipeline_components_finished', duration_ms=round((time.perf_counter() - pipeline_started) * 1000, 3), ) + total_time = time.perf_counter() - pipeline_total_started - result_markdown = get_component_markdown(result, 'saver_component') + result_markdown = get_component_markdown(result, 'dmp_polisher_component') if result_markdown is None: - msg = 'Missing markdown output from saver_component' + msg = 'Missing markdown output from dmp_polisher_component' + logger.error(msg, extra={'template_uuid': str(template_uuid)}) + raise RuntimeError(msg) + + # The clean generator output, not its debug_markdown (which embeds the source-question tables). + dmp_pre_polished = get_component_markdown(result, 'dmp_generator_component') + if dmp_pre_polished is None: + msg = 'Missing markdown output from dmp_generator_component' logger.error(msg, extra={'template_uuid': str(template_uuid)}) raise RuntimeError(msg) @@ -217,33 +211,12 @@ async def run_pipeline( generation_stats = get_component_stats(result, 'dmp_generator_component') polishing_stats = get_component_stats(result, 'dmp_polisher_component') - metrics_started = time.perf_counter() - try: - await write_metrics( - database, - template_uuid, - knowledge_model_uuid, - user_uuid, - tenant_uuid, - result, - model_name, - t1, - ) - except Exception: - logger.exception( - 'Failed to persist pipeline metrics', - extra={'template_uuid': str(template_uuid), 'knowledge_model_uuid': str(knowledge_model_uuid)}, - ) - raise - log_timing_event( - 'pipeline_metrics_saved', - duration_ms=round((time.perf_counter() - metrics_started) * 1000, 3), - ) + pipeline_stats = collect_stats(result, total_time) log_timing_event( 'pipeline_summary', generation_ms=generation_stats.total_duration_ms if generation_stats is not None else None, polishing_ms=polishing_stats.total_duration_ms if polishing_stats is not None else None, - total_pipeline_ms=round((time.perf_counter() - pipeline_total_started) * 1000, 3), + total_pipeline_ms=round(total_time * 1000, 3), total_llm_wait_ms=round( sum( stats.total_llm_wait_ms @@ -261,51 +234,33 @@ async def run_pipeline( 3, ), ) - return knowledge_model_uuid, result_markdown - - -async def write_metrics( - database: Database, - template_uuid: UUID, - knowledge_model_uuid: UUID, - user_uuid: UUID, - tenant_uuid: UUID, - result: Mapping[str, object], - model_name: str, - t1: float, -) -> None: - metrics = PipelineMetricsCollector( - model_name=model_name, - cost_per_mil_input=COST_PER_MIL_INPUT, - cost_per_mil_output=COST_PER_MIL_OUTPUT, - ) - metrics.add_step( - '1. Hierarchical assignment', - get_component_stats(result, 'assignment_saver_component'), - ) - metrics.add_step( - '2. DMP generator', - get_component_stats(result, 'dmp_generator_component'), - ) - metrics.add_step( - '3. DMP polisher', - get_component_stats(result, 'dmp_polisher_component'), + return PipelineOutput( + knowledge_model_uuid=knowledge_model_uuid, + markdown=result_markdown, + dmp_pre_polished=dmp_pre_polished, + stats=pipeline_stats, ) - t2 = time.time() - stats = metrics.get_stats(elapsed_seconds=t2 - t1) - await database.save_stats( - template_uuid=template_uuid, - knowledge_model_uuid=knowledge_model_uuid, - user_uuid=user_uuid, - tenant_uuid=tenant_uuid, - stats=stats, - ) +def collect_stats(result: Mapping[str, object], total_time: float) -> PipelineStats: + """Collect per-step LLM usage. Pure, so it can never fail a finished run.""" + assignment_stats = get_component_stats(result, 'assignment_saver_component') + generation_stats = get_component_stats(result, 'dmp_generator_component') + polishing_stats = get_component_stats(result, 'dmp_polisher_component') - logger.debug('Saved DMP stats') + metrics = PipelineMetricsCollector() + metrics.add_step('1. Hierarchical assignment', assignment_stats) + metrics.add_step('2. DMP generator', generation_stats) + metrics.add_step('3. DMP polisher', polishing_stats) metrics.log_summary(logger) + return PipelineStats( + assignment=StepUsage.from_stats(assignment_stats), + generation=StepUsage.from_stats(generation_stats), + polishing=StepUsage.from_stats(polishing_stats), + elapsed_seconds=total_time, + ) + def _parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser( diff --git a/service/src/ai_document_plugin_service/migrations/script.py.mako b/service/src/ai_document_plugin_service/migrations/script.py.mako index cb466e9..654c20d 100644 --- a/service/src/ai_document_plugin_service/migrations/script.py.mako +++ b/service/src/ai_document_plugin_service/migrations/script.py.mako @@ -21,7 +21,3 @@ depends_on = ${repr(depends_on)} def upgrade() -> None: ${upgrades if upgrades else "pass"} - - -def downgrade() -> None: - ${downgrades if downgrades else "pass"} diff --git a/service/src/ai_document_plugin_service/migrations/versions/20260911_01_merge_result_into_generation.py b/service/src/ai_document_plugin_service/migrations/versions/20260911_01_merge_result_into_generation.py new file mode 100644 index 0000000..fe5cadf --- /dev/null +++ b/service/src/ai_document_plugin_service/migrations/versions/20260911_01_merge_result_into_generation.py @@ -0,0 +1,145 @@ +"""Replace the result table with a per-run generation_stats table. + +Lossy migration. Creates generation_stats (1:1 with generation) for analysis-only data +and the trace id, backfills it from the result table on a best effort basis, then drops +the result table. Result rows that cannot be matched to a succeeded generation are dropped. + +Revision ID: 20260911_01 +Revises: 20260831_02 +Create Date: 2026-09-11 00:00:00 +""" + +from __future__ import annotations + +import logging + +import sqlalchemy as sa +from alembic import context, op +from sqlalchemy.dialects import postgresql + +# revision identifiers, used by Alembic. +revision = '20260911_01' +down_revision = '20260831_02' +branch_labels = None +depends_on = None + +logger = logging.getLogger(__name__) + + +def _qualified_table_reference(schema: str | None, table: str) -> str: + if schema: + return f'{schema}.{table}' + return table + + +def _qualified_column_reference(schema: str | None, table: str, column: str) -> str: + return f'{_qualified_table_reference(schema, table)}.{column}' + + +def upgrade() -> None: + schema = context.get_context().version_table_schema + + op.create_table( + 'generation_stats', + sa.Column('run_id', postgresql.UUID(as_uuid=True), nullable=False), + sa.Column('trace_id', postgresql.UUID(as_uuid=True), nullable=True), + sa.Column('dmp_pre_polished', sa.Text(), nullable=True), + sa.Column('dmp_polished', sa.Text(), nullable=True), + sa.Column('assignment_llm_calls', sa.Integer(), nullable=True), + sa.Column('assignment_input_tokens', sa.Integer(), nullable=True), + sa.Column('assignment_output_tokens', sa.Integer(), nullable=True), + sa.Column('generation_llm_calls', sa.Integer(), nullable=True), + sa.Column('generation_input_tokens', sa.Integer(), nullable=True), + sa.Column('generation_output_tokens', sa.Integer(), nullable=True), + sa.Column('polishing_llm_calls', sa.Integer(), nullable=True), + sa.Column('polishing_input_tokens', sa.Integer(), nullable=True), + sa.Column('polishing_output_tokens', sa.Integer(), nullable=True), + sa.Column('elapsed_seconds', sa.Float(), nullable=True), + sa.Column( + 'created_at', + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.func.now(), + ), + sa.Column( + 'updated_at', + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.func.now(), + ), + sa.ForeignKeyConstraint( + ['run_id'], + [_qualified_column_reference(schema, 'generation', 'run_id')], + name='fk_generation_stats_run_id', + ondelete='CASCADE', + ), + sa.PrimaryKeyConstraint('run_id', name='pk_generation_stats'), + schema=schema, + ) + + # The schema comes from Alembic config, not user input. + generation_reference = _qualified_table_reference(schema, 'generation') + generation_stats_reference = _qualified_table_reference(schema, 'generation_stats') + result_reference = _qualified_table_reference(schema, 'result') + connection = op.get_bind() + + total = connection.execute(sa.text(f'SELECT count(*) FROM {result_reference}')).scalar_one() # ruff: ignore[hardcoded-sql-expression] + + # Result data is attached to the latest succeeded generation matching its DMP. + backfill = connection.execute( + sa.text( + f""" + INSERT INTO {generation_stats_reference} ( + run_id, + dmp_pre_polished, + assignment_llm_calls, assignment_input_tokens, assignment_output_tokens, + generation_llm_calls, generation_input_tokens, generation_output_tokens, + polishing_llm_calls, polishing_input_tokens, polishing_output_tokens, + elapsed_seconds + ) + SELECT g.run_id, + r.dmp_pre_polished, + (asg.step ->> 'llm_calls')::integer, + (asg.step ->> 'input_tokens')::integer, + (asg.step ->> 'output_tokens')::integer, + (gen.step ->> 'llm_calls')::integer, + (gen.step ->> 'input_tokens')::integer, + (gen.step ->> 'output_tokens')::integer, + (pol.step ->> 'llm_calls')::integer, + (pol.step ->> 'input_tokens')::integer, + (pol.step ->> 'output_tokens')::integer, + (r.stats -> 'meta' ->> 'elapsed_seconds')::double precision + FROM {result_reference} AS r + JOIN LATERAL ( + SELECT g2.run_id FROM {generation_reference} AS g2 + WHERE g2.knowledge_model_uuid = r.knowledge_model_uuid + AND g2.template_uuid = r.template_uuid + AND g2.user_uuid = r.user_uuid + AND g2.tenant_uuid = r.tenant_uuid + AND g2.status = 'succeeded' + AND g2.result_markdown = r.dmp + ORDER BY g2.updated_at DESC + LIMIT 1 + ) AS g ON true + LEFT JOIN LATERAL ( + SELECT e.step FROM json_array_elements(r.stats -> 'rows') AS e(step) + WHERE e.step ->> 'step' = '1. Hierarchical assignment' LIMIT 1 + ) AS asg ON true + LEFT JOIN LATERAL ( + SELECT e.step FROM json_array_elements(r.stats -> 'rows') AS e(step) + WHERE e.step ->> 'step' = '2. DMP generator' LIMIT 1 + ) AS gen ON true + LEFT JOIN LATERAL ( + SELECT e.step FROM json_array_elements(r.stats -> 'rows') AS e(step) + WHERE e.step ->> 'step' = '3. DMP polisher' LIMIT 1 + ) AS pol ON true + """ # ruff: ignore[hardcoded-sql-expression] + ) + ) + logger.info( + 'Backfilled generation_stats rows from result table: %s matched, %s unmatched (dropped)', + backfill.rowcount, + total - backfill.rowcount, + ) + + op.drop_table('result', schema=schema) diff --git a/service/src/ai_document_plugin_service/service/errors.py b/service/src/ai_document_plugin_service/service/errors.py index 8b360eb..f23fa44 100644 --- a/service/src/ai_document_plugin_service/service/errors.py +++ b/service/src/ai_document_plugin_service/service/errors.py @@ -28,12 +28,7 @@ def __init__(self, detail: str) -> None: class ConflictError(ServiceError): + PIPELINE_RUN_NOT_FINISHED_MESSAGE = 'Pipeline run has not finished successfully' + def __init__(self, detail: str) -> None: super().__init__(detail, status_code=409) - - -class InternalError(ServiceError): - MISSING_KNOWLEDGE_MODEL_MESSAGE = 'Missing knowledge_model_uuid' - - def __init__(self, detail: str = 'Internal server error') -> None: - super().__init__(detail, status_code=500) diff --git a/service/src/ai_document_plugin_service/service/pipeline_service.py b/service/src/ai_document_plugin_service/service/pipeline_service.py index a2e914d..17fa823 100644 --- a/service/src/ai_document_plugin_service/service/pipeline_service.py +++ b/service/src/ai_document_plugin_service/service/pipeline_service.py @@ -14,10 +14,16 @@ log_timing_event, ) from ai_document_plugin_service.ai.common.llm_client import LLMClient +from ai_document_plugin_service.ai.common.trace_context import get_trace_uuid from ai_document_plugin_service.ai.knowledgemodel.dsw_client import DSWClient from ai_document_plugin_service.ai.persistence.assignment_saver_component import DBSaver -from ai_document_plugin_service.ai.persistence.database import Database, GenerationRecord -from ai_document_plugin_service.ai.run_pipeline import build_pipeline, run_pipeline +from ai_document_plugin_service.ai.persistence.database import ( + Database, + GenerationRecord, + GenerationStats, + GenerationUpdate, +) +from ai_document_plugin_service.ai.run_pipeline import PipelineOutput, build_pipeline, run_pipeline from ai_document_plugin_service.api.auth import AuthenticatedUser from ai_document_plugin_service.api.types import ( ErrorType, @@ -28,7 +34,7 @@ PipelineStatusResponse, PipelineSummaryResponse, ) -from ai_document_plugin_service.service.errors import InternalError, NotFoundError +from ai_document_plugin_service.service.errors import ConflictError, NotFoundError from ai_document_plugin_service.service.pipeline_queue_manager import PipelineQueueManager logger = logging.getLogger(__name__) @@ -85,6 +91,36 @@ def _generation_record_to_summary_response(record: GenerationRecord) -> Pipeline ) +def _succeeded_update(output: PipelineOutput) -> GenerationUpdate: + """The generation fields a finished run writes.""" + return { + 'status': PipelineStatus.SUCCEEDED, + 'knowledge_model_uuid': output.knowledge_model_uuid, + 'result_markdown': output.markdown, + 'progress_message': None, + } + + +def _succeeded_stats(output: PipelineOutput) -> GenerationStats: + """The stats a finished run writes. Per-step usage stays NULL for steps that did not run.""" + stats = output.stats + assignment, generation, polishing = stats.assignment, stats.generation, stats.polishing + return { + 'dmp_pre_polished': output.dmp_pre_polished, + 'dmp_polished': output.markdown, + 'assignment_llm_calls': assignment.llm_calls if assignment else None, + 'assignment_input_tokens': assignment.input_tokens if assignment else None, + 'assignment_output_tokens': assignment.output_tokens if assignment else None, + 'generation_llm_calls': generation.llm_calls if generation else None, + 'generation_input_tokens': generation.input_tokens if generation else None, + 'generation_output_tokens': generation.output_tokens if generation else None, + 'polishing_llm_calls': polishing.llm_calls if polishing else None, + 'polishing_input_tokens': polishing.input_tokens if polishing else None, + 'polishing_output_tokens': polishing.output_tokens if polishing else None, + 'elapsed_seconds': stats.elapsed_seconds, + } + + class LlmClientTenantStore: """ Manages LLM Clients for different tenants. Each tenant has its own LLM client with its own config and limits @@ -190,16 +226,8 @@ async def update_pipeline_result( if record is None: raise NotFoundError(NotFoundError.PIPELINE_RUN_MESSAGE) - if record.knowledge_model_uuid is None: - raise InternalError(InternalError.MISSING_KNOWLEDGE_MODEL_MESSAGE) - - await self.database.update_result( - template_uuid=record.template_uuid, - knowledge_model_uuid=record.knowledge_model_uuid, - user_uuid=auth.user_uuid, - tenant_uuid=auth.tenant_uuid, - markdown=save_request.result_markdown, - ) + if record.status != PipelineStatus.SUCCEEDED: + raise ConflictError(ConflictError.PIPELINE_RUN_NOT_FINISHED_MESSAGE) updated_record = await self.database.update_generation( run_id, @@ -235,14 +263,16 @@ async def _run_pipeline_job( logger.exception('Pipeline run failed', extra={'run_id': run_id, 'tenant_uuid': str(auth.tenant_uuid)}) pipeline_error = _pipeline_error_from_exception(error) try: - await self.database.update_generation( - run_id, - auth.tenant_uuid, - status=PipelineStatus.FAILED, - error_type=pipeline_error.type, - error_message=pipeline_error.message, - progress_message=None, - ) + async with self.database.transaction(): + await self.database.update_generation( + run_id, + auth.tenant_uuid, + status=PipelineStatus.FAILED, + error_type=pipeline_error.type, + error_message=pipeline_error.message, + progress_message=None, + ) + await self.database.create_generation_stats(run_id, get_trace_uuid()) except Exception: logger.exception( 'Failed to persist pipeline failure status', @@ -261,13 +291,15 @@ async def _run_pipeline( ) -> None: template = await self.database.get_template(template_uuid, auth.tenant_uuid) if template is None: - await self.database.update_generation( - run_id, - auth.tenant_uuid, - status=PipelineStatus.FAILED, - error_type=ErrorType.TEMPLATE_NOT_FOUND, - error_message=TEMPLATE_NOT_FOUND_MESSAGE, - ) + async with self.database.transaction(): + await self.database.update_generation( + run_id, + auth.tenant_uuid, + status=PipelineStatus.FAILED, + error_type=ErrorType.TEMPLATE_NOT_FOUND, + error_message=TEMPLATE_NOT_FOUND_MESSAGE, + ) + await self.database.create_generation_stats(run_id, get_trace_uuid()) logger.error( 'Pipeline run failed because template was not found', extra={'run_id': run_id, 'template_uuid': str(template_uuid), 'tenant_uuid': str(auth.tenant_uuid)}, @@ -308,35 +340,28 @@ def on_progress(message: str) -> None: ) task.add_done_callback(_log_background_update_failure) - knowledge_model_uuid, result = await run_pipeline( + output = await run_pipeline( questionnaire_uuid=questionnaire_uuid, template_uuid=template_uuid, template_title=template.title, template_data=template.content, - user_uuid=auth.user_uuid, tenant_uuid=auth.tenant_uuid, pipeline=pipeline, - database=self.database, on_progress=on_progress, - model_name=llm_client.get_model_name(), dsw_client=DSWClient(auth.token, auth.api_url), ) - log_timing_event('pipeline_generation_finished', knowledge_model_uuid=str(knowledge_model_uuid)) + log_timing_event('pipeline_generation_finished', knowledge_model_uuid=str(output.knowledge_model_uuid)) - await self.database.update_generation( - run_id, - auth.tenant_uuid, - status=PipelineStatus.SUCCEEDED, - knowledge_model_uuid=knowledge_model_uuid, - result_markdown=result, - progress_message=None, - ) + # Everything the run produced is written in one transaction, so it lands atomically. + async with self.database.transaction(): + await self.database.update_generation(run_id, auth.tenant_uuid, **_succeeded_update(output)) + await self.database.create_generation_stats(run_id, get_trace_uuid(), **_succeeded_stats(output)) logger.info( 'Pipeline run status updated to succeeded', extra={ 'run_id': run_id, 'tenant_uuid': str(auth.tenant_uuid), - 'knowledge_model_uuid': str(knowledge_model_uuid), - 'result_markdown_length': len(result), + 'knowledge_model_uuid': str(output.knowledge_model_uuid), + 'result_markdown_length': len(output.markdown), }, ) diff --git a/service/tests/common/test_pipeline_metrics.py b/service/tests/common/test_pipeline_metrics.py index 00cd06f..e2181e2 100644 --- a/service/tests/common/test_pipeline_metrics.py +++ b/service/tests/common/test_pipeline_metrics.py @@ -1,8 +1,6 @@ -from collections.abc import Mapping -from typing import cast +import time from ai_document_plugin_service.ai.common.pipeline_metrics import ( - PipelineMetricsCollector, get_component_markdown, get_component_stats, ) @@ -36,62 +34,3 @@ def test_get_component_markdown_returns_value_for_component_output() -> None: output = get_component_markdown(result, 'component_a') assert output == '# DMP' - - -def test_get_stats_returns_json_summary_with_expected_shape() -> None: - collector = PipelineMetricsCollector( - model_name='test-model', - cost_per_mil_input=0.25, - cost_per_mil_output=2.0, - ) - collector.add_step( - '1. Step', - AssignmentStats(total_calls=2, total_input_tokens=1000, total_output_tokens=200), - ) - collector.add_step( - '2. Step', - AssignmentStats(total_calls=1, total_input_tokens=500, total_output_tokens=50), - ) - - stats = collector.get_stats(elapsed_seconds=12.5) - - assert isinstance(stats, Mapping) - typed_stats = cast(dict[str, object], stats) - - assert typed_stats['title'] == 'Pipeline token usage and cost' - assert typed_stats['headers'] == [ - 'Step', - 'LLM calls', - 'Input tokens', - 'Output tokens', - 'Cost (USD)', - ] - assert typed_stats['rows'] == [ - { - 'step': '1. Step', - 'llm_calls': 2, - 'input_tokens': 1000, - 'output_tokens': 200, - 'cost_usd': 0.0, - }, - { - 'step': '2. Step', - 'llm_calls': 1, - 'input_tokens': 500, - 'output_tokens': 50, - 'cost_usd': 0.0, - }, - ] - assert typed_stats['totals'] == { - 'input_tokens': 1500, - 'output_tokens': 250, - 'cost_usd': 0.0, - } - assert typed_stats['meta'] == { - 'model_name': 'test-model', - 'cost_per_mil_input': 0.25, - 'cost_per_mil_output': 2.0, - 'elapsed_seconds': 12.5, - 'total_llm_wait_ms': 0.0, - 'total_llm_response_ms': 0.0, - }