mirror of
https://github.com/langbot-app/LangBot.git
synced 2026-08-09 20:50:58 +00:00
393 lines
14 KiB
Python
393 lines
14 KiB
Python
from __future__ import annotations
|
|
|
|
import datetime
|
|
from types import SimpleNamespace
|
|
|
|
import pytest
|
|
import sqlalchemy
|
|
from sqlalchemy.ext.asyncio import create_async_engine
|
|
|
|
from langbot.pkg.api.http.authz import WorkspaceRequiredError
|
|
from langbot.pkg.api.http.context import ExecutionContext
|
|
from langbot.pkg.api.http.service.monitoring import MonitoringService
|
|
from langbot.pkg.entity.persistence.base import Base
|
|
from langbot.pkg.entity.persistence.monitoring import MonitoringLLMCall, MonitoringMessage
|
|
from langbot.pkg.entity.persistence.workspace import Workspace
|
|
from langbot.pkg.persistence.mgr import PersistenceManager
|
|
|
|
|
|
pytestmark = pytest.mark.asyncio
|
|
|
|
WORKSPACE_A = '00000000-0000-0000-0000-00000000000a'
|
|
WORKSPACE_B = '00000000-0000-0000-0000-00000000000b'
|
|
|
|
|
|
def _context(workspace_uuid: str) -> ExecutionContext:
|
|
return ExecutionContext(
|
|
instance_uuid='instance',
|
|
workspace_uuid=workspace_uuid,
|
|
placement_generation=3,
|
|
bot_uuid='same-bot',
|
|
pipeline_uuid='same-pipeline',
|
|
)
|
|
|
|
|
|
class _PersistenceManager:
|
|
def __init__(self, engine):
|
|
self.engine = engine
|
|
|
|
async def execute_async(self, *args, **kwargs):
|
|
async with self.engine.connect() as connection:
|
|
result = await connection.execute(*args, **kwargs)
|
|
await connection.commit()
|
|
return result
|
|
|
|
def get_db_engine(self):
|
|
return self.engine
|
|
|
|
@staticmethod
|
|
def serialize_model(model, data, masked_columns=None):
|
|
return {
|
|
column.name: (
|
|
getattr(data, column.name).isoformat()
|
|
if isinstance(getattr(data, column.name), datetime.datetime)
|
|
else getattr(data, column.name)
|
|
)
|
|
for column in model.__table__.columns
|
|
if column.name not in (masked_columns or [])
|
|
}
|
|
|
|
|
|
@pytest.fixture
|
|
async def service(tmp_path):
|
|
engine = create_async_engine(f'sqlite+aiosqlite:///{tmp_path / "monitoring.db"}')
|
|
async with engine.begin() as connection:
|
|
await connection.run_sync(Base.metadata.create_all)
|
|
await connection.execute(
|
|
sqlalchemy.insert(Workspace),
|
|
[
|
|
{
|
|
'uuid': WORKSPACE_A,
|
|
'instance_uuid': 'instance',
|
|
'name': 'A',
|
|
'slug': 'a',
|
|
'source': 'cloud_projection',
|
|
},
|
|
{
|
|
'uuid': WORKSPACE_B,
|
|
'instance_uuid': 'instance',
|
|
'name': 'B',
|
|
'slug': 'b',
|
|
'source': 'cloud_projection',
|
|
},
|
|
],
|
|
)
|
|
application = SimpleNamespace(
|
|
persistence_mgr=_PersistenceManager(engine),
|
|
instance_config=SimpleNamespace(data={'database': {'use': 'sqlite'}}),
|
|
)
|
|
yield MonitoringService(application)
|
|
await engine.dispose()
|
|
|
|
|
|
async def _record_message(service, context, content):
|
|
return await service.record_message(
|
|
context,
|
|
bot_id='same-bot',
|
|
bot_name='Same Bot',
|
|
pipeline_id='same-pipeline',
|
|
pipeline_name='Same Pipeline',
|
|
message_content=content,
|
|
session_id='same-session',
|
|
)
|
|
|
|
|
|
async def test_monitoring_write_without_execution_context_fails_closed(service):
|
|
with pytest.raises(WorkspaceRequiredError):
|
|
await _record_message(service, None, 'unscoped')
|
|
|
|
|
|
async def test_same_session_and_resource_ids_do_not_collide(service):
|
|
context_a = _context(WORKSPACE_A)
|
|
context_b = _context(WORKSPACE_B)
|
|
message_a = await _record_message(service, context_a, 'tenant-a')
|
|
message_b = await _record_message(service, context_b, 'tenant-b')
|
|
await service.record_session_start(
|
|
context_a,
|
|
session_id='same-session',
|
|
bot_id='same-bot',
|
|
bot_name='Same Bot',
|
|
pipeline_id='same-pipeline',
|
|
pipeline_name='Same Pipeline',
|
|
)
|
|
await service.record_session_start(
|
|
context_b,
|
|
session_id='same-session',
|
|
bot_id='same-bot',
|
|
bot_name='Same Bot',
|
|
pipeline_id='same-pipeline',
|
|
pipeline_name='Same Pipeline',
|
|
)
|
|
|
|
messages_a, total_a = await service.get_messages(context_a)
|
|
messages_b, total_b = await service.get_messages(context_b)
|
|
assert total_a == total_b == 1
|
|
assert messages_a[0]['message_content'] == 'tenant-a'
|
|
assert messages_b[0]['message_content'] == 'tenant-b'
|
|
assert (await service.get_message_details(context_b, message_a))['found'] is False
|
|
assert (await service.get_message_details(context_a, message_b))['found'] is False
|
|
|
|
|
|
async def test_tool_call_inherits_context_from_connection_message_row(service):
|
|
context = _context(WORKSPACE_A)
|
|
message_id = await _record_message(service, context, 'tool context')
|
|
|
|
await service.record_tool_call(
|
|
context,
|
|
tool_name='search',
|
|
tool_source='native',
|
|
duration=12,
|
|
message_id=message_id,
|
|
)
|
|
|
|
tool_calls, total = await service.get_tool_calls(context)
|
|
assert total == 1
|
|
assert tool_calls[0]['bot_id'] == 'same-bot'
|
|
assert tool_calls[0]['pipeline_id'] == 'same-pipeline'
|
|
assert tool_calls[0]['session_id'] == 'same-session'
|
|
assert tool_calls[0]['message_id'] == message_id
|
|
|
|
|
|
async def test_feedback_upsert_and_cancel_are_workspace_scoped(service):
|
|
context_a = _context(WORKSPACE_A)
|
|
context_b = _context(WORKSPACE_B)
|
|
await service.record_feedback(context_a, feedback_id='same-feedback', feedback_type=1)
|
|
await service.record_feedback(context_b, feedback_id='same-feedback', feedback_type=2)
|
|
|
|
stats_a = await service.get_feedback_stats(context_a)
|
|
stats_b = await service.get_feedback_stats(context_b)
|
|
assert stats_a['total_likes'] == 1
|
|
assert stats_a['total_dislikes'] == 0
|
|
assert stats_b['total_likes'] == 0
|
|
assert stats_b['total_dislikes'] == 1
|
|
|
|
await service.record_feedback(context_a, feedback_id='same-feedback', feedback_type=3)
|
|
assert (await service.get_feedback_stats(context_a))['total_feedback'] == 0
|
|
assert (await service.get_feedback_stats(context_b))['total_feedback'] == 1
|
|
|
|
|
|
async def test_monitoring_queries_and_detail_views_are_strictly_bounded(service):
|
|
context = _context(WORKSPACE_A)
|
|
service.ap.instance_config.data['monitoring'] = {
|
|
'query_limits': {
|
|
'page_rows': 2,
|
|
'export_rows': 2,
|
|
'detail_rows': 2,
|
|
'timeseries_buckets': 2,
|
|
'max_offset': 10,
|
|
}
|
|
}
|
|
await service.record_session_start(
|
|
context,
|
|
session_id='same-session',
|
|
bot_id='same-bot',
|
|
bot_name='Same Bot',
|
|
pipeline_id='same-pipeline',
|
|
pipeline_name='Same Pipeline',
|
|
)
|
|
message_ids = [await _record_message(service, context, f'message-{index}') for index in range(4)]
|
|
for index in range(3):
|
|
await service.record_llm_call(
|
|
context,
|
|
bot_id='same-bot',
|
|
bot_name='Same Bot',
|
|
pipeline_id='same-pipeline',
|
|
pipeline_name='Same Pipeline',
|
|
session_id='same-session',
|
|
model_name='model',
|
|
input_tokens=1,
|
|
output_tokens=2,
|
|
duration=10,
|
|
message_id=message_ids[0],
|
|
)
|
|
await service.record_tool_call(
|
|
context,
|
|
tool_name=f'tool-{index}',
|
|
tool_source='native',
|
|
duration=5,
|
|
session_id='same-session',
|
|
message_id=message_ids[0],
|
|
)
|
|
await service.record_error(
|
|
context,
|
|
bot_id='same-bot',
|
|
bot_name='Same Bot',
|
|
pipeline_id='same-pipeline',
|
|
pipeline_name='Same Pipeline',
|
|
error_type='Failure',
|
|
error_message=f'error-{index}',
|
|
session_id='same-session',
|
|
message_id=message_ids[0],
|
|
)
|
|
|
|
page, total = await service.get_messages(context, limit=100000, offset=-5)
|
|
exported = await service.export_messages(context, limit=100000)
|
|
session_detail = await service.get_session_analysis(context, 'same-session')
|
|
message_detail = await service.get_message_details(context, message_ids[0])
|
|
|
|
assert total == 4
|
|
assert len(page) == 2
|
|
assert len(exported) == 2
|
|
assert session_detail['message_stats']['total'] == 4
|
|
assert session_detail['llm_stats']['total_calls'] == 3
|
|
assert session_detail['tool_stats']['total_calls'] == 3
|
|
assert len(session_detail['tool_calls']) == 2
|
|
assert len(session_detail['errors']) == 2
|
|
assert session_detail['detail_truncated'] == {
|
|
'tool_calls': True,
|
|
'errors': True,
|
|
}
|
|
assert message_detail['llm_stats']['total_calls'] == 3
|
|
assert len(message_detail['llm_calls']) == 2
|
|
assert len(message_detail['errors']) == 2
|
|
assert message_detail['detail_truncated'] == {
|
|
'llm_calls': True,
|
|
'errors': True,
|
|
}
|
|
|
|
service.ap.instance_config.data['monitoring']['query_limits'] = {
|
|
'page_rows': 999999,
|
|
'export_rows': 999999,
|
|
'detail_rows': 999999,
|
|
'timeseries_buckets': 999999,
|
|
'max_offset': 99999999,
|
|
}
|
|
assert service.normalize_page_window(999999, 99999999) == (5000, 10000000)
|
|
assert service.normalize_export_limit(999999) == 50000
|
|
assert service._detail_limit() == 10000
|
|
assert service._timeseries_bucket_limit() == 10000
|
|
|
|
|
|
async def test_token_statistics_aggregate_and_limit_groups_in_database(service):
|
|
context = _context(WORKSPACE_A)
|
|
service.ap.instance_config.data['monitoring'] = {
|
|
'query_limits': {
|
|
'page_rows': 1,
|
|
'timeseries_buckets': 2,
|
|
}
|
|
}
|
|
first_hour = datetime.datetime(2026, 7, 28, 10, 0)
|
|
rows = [
|
|
{
|
|
'id': f'llm-{index}',
|
|
'workspace_uuid': WORKSPACE_A,
|
|
'timestamp': first_hour + datetime.timedelta(hours=hour, minutes=index),
|
|
'model_name': model,
|
|
'input_tokens': input_tokens,
|
|
'output_tokens': output_tokens,
|
|
'total_tokens': input_tokens + output_tokens,
|
|
'duration': 100,
|
|
'cost': 0.01,
|
|
'status': 'success',
|
|
'bot_id': 'same-bot',
|
|
'bot_name': 'Same Bot',
|
|
'pipeline_id': 'same-pipeline',
|
|
'pipeline_name': 'Same Pipeline',
|
|
'session_id': 'same-session',
|
|
}
|
|
for index, (hour, model, input_tokens, output_tokens) in enumerate(
|
|
[
|
|
(0, 'small-model', 1, 2),
|
|
(1, 'large-model', 3, 4),
|
|
(2, 'large-model', 5, 6),
|
|
(2, 'large-model', 7, 8),
|
|
]
|
|
)
|
|
]
|
|
await service.ap.persistence_mgr.execute_async(sqlalchemy.insert(MonitoringLLMCall), rows)
|
|
|
|
stats = await service.get_token_statistics(context, bucket='hour')
|
|
|
|
assert stats['summary']['total_calls'] == 4
|
|
assert stats['summary']['total_tokens'] == 36
|
|
assert stats['by_model_truncated'] is True
|
|
assert [model['model_name'] for model in stats['by_model']] == ['large-model']
|
|
assert stats['timeseries_truncated'] is True
|
|
assert stats['timeseries'] == [
|
|
{
|
|
'bucket': '2026-07-28 11:00',
|
|
'input_tokens': 3,
|
|
'output_tokens': 4,
|
|
'total_tokens': 7,
|
|
'calls': 1,
|
|
},
|
|
{
|
|
'bucket': '2026-07-28 12:00',
|
|
'input_tokens': 12,
|
|
'output_tokens': 14,
|
|
'total_tokens': 26,
|
|
'calls': 2,
|
|
},
|
|
]
|
|
|
|
|
|
async def test_cleanup_commits_sqlite_delete_before_vacuum(tmp_path):
|
|
engine = create_async_engine(
|
|
f'sqlite+aiosqlite:///{tmp_path / "monitoring-cleanup.db"}',
|
|
connect_args={'timeout': 0.1},
|
|
)
|
|
application = SimpleNamespace(
|
|
instance_config=SimpleNamespace(data={'database': {'use': 'sqlite'}}),
|
|
)
|
|
manager = PersistenceManager(application)
|
|
manager.db = SimpleNamespace(get_engine=lambda: engine)
|
|
application.persistence_mgr = manager
|
|
try:
|
|
async with engine.begin() as connection:
|
|
await connection.run_sync(Base.metadata.create_all)
|
|
await connection.execute(
|
|
sqlalchemy.insert(Workspace).values(
|
|
uuid=WORKSPACE_A,
|
|
instance_uuid='instance',
|
|
name='A',
|
|
slug='a',
|
|
source='cloud_projection',
|
|
)
|
|
)
|
|
await connection.execute(
|
|
sqlalchemy.insert(MonitoringMessage),
|
|
[
|
|
{
|
|
'id': f'expired-message-{index}',
|
|
'workspace_uuid': WORKSPACE_A,
|
|
'timestamp': datetime.datetime.now(datetime.timezone.utc).replace(tzinfo=None)
|
|
- datetime.timedelta(days=30),
|
|
'bot_id': 'bot',
|
|
'bot_name': 'Bot',
|
|
'pipeline_id': 'pipeline',
|
|
'pipeline_name': 'Pipeline',
|
|
'message_content': 'expired',
|
|
'session_id': 'session',
|
|
'status': 'success',
|
|
'level': 'info',
|
|
}
|
|
for index in range(5)
|
|
],
|
|
)
|
|
|
|
deleted = await MonitoringService(application).cleanup_expired_records(
|
|
_context(WORKSPACE_A),
|
|
retention_days=1,
|
|
batch_size=2,
|
|
max_batches_per_table=1,
|
|
)
|
|
|
|
assert deleted['monitoring_messages'] == 2
|
|
async with engine.connect() as connection:
|
|
remaining = await connection.scalar(
|
|
sqlalchemy.select(sqlalchemy.func.count()).select_from(MonitoringMessage)
|
|
)
|
|
assert remaining == 3
|
|
finally:
|
|
await engine.dispose()
|