fix(migrations): support partial monitoring schemas and align regression fixtures

This commit is contained in:
dadachann
2026-09-11 14:21:44 +08:00
parent 8bc1411d07
commit 5fc5c242ee
6 changed files with 41 additions and 42 deletions
@@ -20,6 +20,8 @@ _KEY = ['workspace_uuid', 'bot_id', 'session_id']
def upgrade() -> None: def upgrade() -> None:
conn = op.get_bind() conn = op.get_bind()
inspector = sa.inspect(conn) inspector = sa.inspect(conn)
if _TABLE not in inspector.get_table_names():
return
pk = inspector.get_pk_constraint(_TABLE) pk = inspector.get_pk_constraint(_TABLE)
if pk['constrained_columns'] == _KEY: if pk['constrained_columns'] == _KEY:
return return
@@ -89,6 +91,8 @@ def upgrade() -> None:
def downgrade() -> None: def downgrade() -> None:
conn = op.get_bind() conn = op.get_bind()
if _TABLE not in sa.inspect(conn).get_table_names():
return
collisions = conn.execute( collisions = conn.execute(
sa.text('SELECT 1 FROM monitoring_sessions GROUP BY workspace_uuid, session_id HAVING COUNT(*) > 1 LIMIT 1') sa.text('SELECT 1 FROM monitoring_sessions GROUP BY workspace_uuid, session_id HAVING COUNT(*) > 1 LIMIT 1')
).first() ).first()
+8 -1
View File
@@ -9,7 +9,7 @@ Run: uv run pytest tests/integration/api/test_monitoring.py -q
from __future__ import annotations from __future__ import annotations
import pytest import pytest
from unittest.mock import MagicMock, AsyncMock, Mock from unittest.mock import MagicMock, AsyncMock, Mock, patch
from types import SimpleNamespace from types import SimpleNamespace
from tests.factories import FakeApp from tests.factories import FakeApp
@@ -280,13 +280,20 @@ class TestMonitoringAllDataEndpoint:
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_get_all_data_success(self, quart_test_client): async def test_get_all_data_success(self, quart_test_client):
"""GET /api/v1/monitoring/data returns all data.""" """GET /api/v1/monitoring/data returns all data."""
traffic = {'series': [], 'truncated': False}
with patch(
'langbot.pkg.api.http.controller.groups.monitoring.get_traffic_series',
new=AsyncMock(return_value=traffic),
) as get_traffic:
response = await quart_test_client.get( response = await quart_test_client.get(
'/api/v1/monitoring/data', headers={'Authorization': 'Bearer test_token'} '/api/v1/monitoring/data', headers={'Authorization': 'Bearer test_token'}
) )
get_traffic.assert_awaited_once()
assert response.status_code == 200 assert response.status_code == 200
data = await response.get_json() data = await response.get_json()
assert 'overview' in data['data'] assert 'overview' in data['data']
assert data['data']['traffic'] == traffic
@pytest.mark.usefixtures('mock_circular_import_chain') @pytest.mark.usefixtures('mock_circular_import_chain')
@@ -193,6 +193,22 @@ async def create_legacy_resource_schema(engine, *, instance_uuid: str) -> None:
sa.Column('message_id', sa.String(255), nullable=True), sa.Column('message_id', sa.String(255), nullable=True),
) )
# Include historical monitoring columns consumed by later migrations.
for table_name in ('monitoring_messages', 'monitoring_sessions'):
table = monitoring_tables[table_name]
for name, value in (('bot_name', 'bot'), ('pipeline_id', 'pipeline-1'), ('pipeline_name', 'pipeline')):
table.append_column(sa.Column(name, sa.String(255), nullable=False, default=value))
for name in ('platform', 'user_id', 'user_name'):
table.append_column(sa.Column(name, sa.String(255)))
if table_name == 'monitoring_messages':
table.append_column(sa.Column('bot_id', sa.String(255), nullable=False, default='bot-1'))
table.append_column(sa.Column('role', sa.String(50)))
else:
table.append_column(sa.Column('message_count', sa.Integer, nullable=False, default=1))
table.append_column(
sa.Column('start_time', sa.DateTime, nullable=False, default=datetime.datetime(2026, 1, 1))
)
now = datetime.datetime(2026, 1, 1) now = datetime.datetime(2026, 1, 1)
async with engine.begin() as conn: async with engine.begin() as conn:
await conn.run_sync(metadata.create_all) await conn.run_sync(metadata.create_all)
@@ -17,7 +17,6 @@ from sqlalchemy import text
from sqlalchemy.ext.asyncio import create_async_engine from sqlalchemy.ext.asyncio import create_async_engine
from langbot.pkg.entity.persistence.base import Base from langbot.pkg.entity.persistence.base import Base
from langbot.pkg.entity.persistence.monitoring import MonitoringMessage, MonitoringSession
from langbot.pkg.persistence import mgr as persistence_mgr # noqa: F401 -- register all ORM tables from langbot.pkg.persistence import mgr as persistence_mgr # noqa: F401 -- register all ORM tables
from langbot.pkg.persistence.alembic_runner import ( from langbot.pkg.persistence.alembic_runner import (
run_alembic_downgrade, run_alembic_downgrade,
@@ -100,17 +99,6 @@ class TestSQLiteMigrationBaseline:
class TestSQLiteMigrationUpgrade: class TestSQLiteMigrationUpgrade:
"""Tests for upgrade to head workflow.""" """Tests for upgrade to head workflow."""
@pytest.fixture(autouse=True)
async def existing_monitoring_tables(self, sqlite_engine):
# Historical instances already have monitoring tables. Partial fixtures
# below omit unrelated tables, but later session migrations require these.
async with sqlite_engine.begin() as conn:
await conn.run_sync(
lambda sync: Base.metadata.create_all(
sync, tables=[MonitoringMessage.__table__, MonitoringSession.__table__]
)
)
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_upgrade_from_published_space_launch_head_to_merged_head(self, sqlite_engine): async def test_upgrade_from_published_space_launch_head_to_merged_head(self, sqlite_engine):
"""A database released at the production-only 0016 head must remain upgradable.""" """A database released at the production-only 0016 head must remain upgradable."""
@@ -292,6 +280,15 @@ class TestSQLiteMigrationUpgrade:
class TestSQLiteMigrationFreshDatabase: class TestSQLiteMigrationFreshDatabase:
"""Tests for fresh database workflow.""" """Tests for fresh database workflow."""
@pytest.mark.asyncio
async def test_bot_scoped_sessions_skips_absent_table(self, sqlite_engine):
"""A partial schema needs no session key migration in either direction."""
await run_alembic_stamp(sqlite_engine, '0022_codex_credentials')
await run_alembic_upgrade(sqlite_engine, '0023_bot_scoped_sessions')
assert await get_alembic_current(sqlite_engine) == '0023_bot_scoped_sessions'
await run_alembic_downgrade(sqlite_engine, '0022_codex_credentials')
assert await get_alembic_current(sqlite_engine) == '0022_codex_credentials'
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_fresh_db_upgrade_from_scratch(self, tmp_path): async def test_fresh_db_upgrade_from_scratch(self, tmp_path):
""" """
@@ -589,31 +589,6 @@ class TestPostgreSQLResourceTenancyMigration:
instance_uuid='postgres-resource-migration-test', instance_uuid='postgres-resource-migration-test',
) )
async with postgres_engine.begin() as conn: async with postgres_engine.begin() as conn:
# The shared tenancy fixture is intentionally lean. Restore the
# monitoring columns present since 5d9f6ec7 (user_name: 89064a9d)
# before exercising later migrations that read their contents.
for table_name in ('monitoring_messages', 'monitoring_sessions'):
columns = {
'bot_name': "VARCHAR(255) NOT NULL DEFAULT 'bot'",
'pipeline_id': "VARCHAR(255) NOT NULL DEFAULT 'pipeline-1'",
'pipeline_name': "VARCHAR(255) NOT NULL DEFAULT 'pipeline'",
'platform': 'VARCHAR(255)',
'user_id': 'VARCHAR(255)',
'user_name': 'VARCHAR(255)',
}
if table_name == 'monitoring_messages':
columns.update(
bot_id="VARCHAR(255) NOT NULL DEFAULT 'bot-1'",
role='VARCHAR(50)',
)
else:
columns.update(
message_count='INTEGER NOT NULL DEFAULT 1',
start_time="TIMESTAMP NOT NULL DEFAULT '2026-01-01'",
)
for column_name, definition in columns.items():
await conn.execute(text(f'ALTER TABLE {table_name} ADD COLUMN {column_name} {definition}'))
await conn.execute(text(f'ALTER TABLE {table_name} ALTER COLUMN {column_name} DROP DEFAULT'))
await conn.execute(text('UPDATE users SET "user" = \'Straße@Example.COM\'')) await conn.execute(text('UPDATE users SET "user" = \'Straße@Example.COM\''))
await conn.execute( await conn.execute(
text('INSERT INTO users ("user", password) VALUES (:email, :password)'), text('INSERT INTO users ("user", password) VALUES (:email, :password)'),
@@ -142,7 +142,7 @@ async def test_legacy_sqlite_resources_are_backfilled_and_contracted(tmp_path):
assert pk_columns == { assert pk_columns == {
'binary_storages': ('workspace_uuid', 'unique_key'), 'binary_storages': ('workspace_uuid', 'unique_key'),
'plugin_settings': ('workspace_uuid', 'plugin_author', 'plugin_name'), 'plugin_settings': ('workspace_uuid', 'plugin_author', 'plugin_name'),
'monitoring_sessions': ('workspace_uuid', 'session_id'), 'monitoring_sessions': ('workspace_uuid', 'bot_id', 'session_id'),
} }
pipeline_run_foreign_keys = await _inspect( pipeline_run_foreign_keys = await _inspect(