Files

204 lines
7.6 KiB
Python

"""Isolated PostgreSQL stage-one migration/restore acceptance; fixed test URL only."""
import asyncio
import os
import subprocess
from pathlib import Path
from alembic import command
from alembic.config import Config
from cryptography.fernet import Fernet
from sqlalchemy import text
from sqlalchemy.ext.asyncio import create_async_engine
URL = "postgresql+asyncpg://postgres:research-test-only@127.0.0.1:18436/wq_research_stage1_test"
os.environ.update(
DATABASE_URL=URL, ADMIN_PASSWORD="research-acceptance-only", ENCRYPTION_KEY=Fernet.generate_key().decode()
)
async def sql(statement):
engine = create_async_engine(URL)
async with engine.begin() as db:
result = await db.execute(text(statement))
rows = result.fetchall() if result.returns_rows else None
await engine.dispose()
return rows
async def acceptance():
import httpx
from app.config import Settings
from app.main import create_app
from app.worldquant import WqClient
from tests.catalog_fake import catalog_response
from tests.research_metadata_fake import response as metadata_response
from tests.test_backtests import execute, setup, start
from tests.test_catalog import prepare, sync
from tests.test_research_workspace import expansion, template
def upstream(request):
if request.url.path == "/authentication":
return httpx.Response(201, json={})
if request.url.path == "/users/self":
return httpx.Response(200, json={"id": "PG_RESEARCH_USER"})
return metadata_response(request) or catalog_response(request) or httpx.Response(404)
settings = Settings(_env_file=None, enable_runner=False, public_origin="http://testserver")
app = create_app(settings, WqClient(settings, transport=httpx.MockTransport(upstream)))
async with app.router.lifespan_context(app):
async with httpx.AsyncClient(
transport=httpx.ASGITransport(app), base_url="http://testserver", headers={"X-WQ-Request": "1"}
) as client:
assert (
await client.post(
"/api/v1/auth/login", json={"username": "admin", "password": "research-acceptance-only"}
)
).status_code == 200
await client.put(
"/api/v1/account/credentials",
json={"email": "test@example.com", "password": "synthetic-only"},
)
job = (await client.post("/api/v1/account/connect")).json()
await app.state.runner.execute(job["id"])
catalog = (client, app.state.runner, {})
await sync(catalog)
version = (await sync(catalog, "TEST_FIN"))["id"]
fixed = (await prepare(client, version)).json()
assert (await client.post("/api/v1/catalog/operators/refresh")).status_code == 200
assert (await client.post("/api/v1/catalog/setting-options/refresh")).status_code == 200
saved = (
await client.post("/api/v1/research/assets", json={"kind": "template", "content": template()})
).json()
a, b = await asyncio.gather(
*(
client.put(
"/api/v1/research/assets/" + saved["id"],
json={"kind": "template", "version": 1, "content": {**template(), "name": name}},
)
for name in ["A", "B"]
)
)
assert sorted([a.status_code, b.status_code]) == [200, 409]
old = (await client.get("/api/v1/research/assets/" + saved["id"] + "?version=1")).json()
assert old["name"] == template()["name"]
body = expansion(fixed["id"], asset_id=saved["id"], version=1)
body.pop("template")
experiment = (await client.post("/api/v1/research/experiments", json=body)).json()
assert len(experiment["candidates"]) == 2
fake, lane = await setup(app)
preview = (
await client.post("/api/v1/research/experiments/" + experiment["id"] + "/preview", json={})
).json()
first, second = await asyncio.gather(
start(client, preview, "research-confirm"), start(client, preview, "research-confirm")
)
assert first["backtest_run_id"] == second["backtest_run_id"]
await execute(app, lane, first["backtest_run_id"])
current = (await client.get("/api/v1/research/experiments/" + experiment["id"])).json()
assert current["backtest_run_ids"] == [first["backtest_run_id"]]
print(
"PASS PostgreSQL: versions, concurrent CAS, fixed input → experiment → idempotent backtest → provenance"
)
if __name__ == "__main__":
config = Config("alembic.ini")
if asyncio.run(sql("SELECT tablename FROM pg_tables WHERE schemaname='public'")):
raise RuntimeError("Dedicated acceptance database must be empty")
command.upgrade(config, "0005")
asyncio.run(
sql(
"INSERT INTO alphas (id,hidden,settings,is_metrics,os_metrics,checks,synced_at,raw) VALUES ('OLD_RESEARCH',false,'{}','{}','{}','[]',now(),'{}')"
)
)
asyncio.run(
sql(
"INSERT INTO research (alpha_id,note,tags,favorite,state,updated_at,version) VALUES ('OLD_RESEARCH','preserve note','[]',false,'inbox',now(),7)"
)
)
command.upgrade(config, "head")
command.check(config)
assert asyncio.run(sql("SELECT note,version FROM research WHERE alpha_id='OLD_RESEARCH'")) == [
("preserve note", 7)
]
asyncio.run(acceptance())
dump = Path("/tmp/wq-research-stage1.dump")
with dump.open("wb") as output:
subprocess.run(
[
"docker",
"exec",
"wq-research-acceptance-pg",
"pg_dump",
"-U",
"postgres",
"-Fc",
"wq_research_stage1_test",
],
stdout=output,
check=True,
)
subprocess.run(
[
"docker",
"exec",
"wq-research-acceptance-pg",
"createdb",
"-U",
"postgres",
"wq_research_restore_stage1",
],
check=True,
)
with dump.open("rb") as input_file:
subprocess.run(
[
"docker",
"exec",
"-i",
"wq-research-acceptance-pg",
"pg_restore",
"-U",
"postgres",
"-d",
"wq_research_restore_stage1",
],
stdin=input_file,
check=True,
)
query = "SELECT (SELECT count(*) FROM research_revisions),(SELECT count(*) FROM research_experiments),(SELECT count(*) FROM backtest_runs),(SELECT note FROM research WHERE alpha_id='OLD_RESEARCH')"
original = subprocess.check_output(
[
"docker",
"exec",
"wq-research-acceptance-pg",
"psql",
"-U",
"postgres",
"-d",
"wq_research_stage1_test",
"-Atc",
query,
]
)
restored = subprocess.check_output(
[
"docker",
"exec",
"wq-research-acceptance-pg",
"psql",
"-U",
"postgres",
"-d",
"wq_research_restore_stage1",
"-Atc",
query,
]
)
assert original == restored
print(
"PASS PostgreSQL 17: 0005 → 0006, schema check, old notes preserved, pg_dump/pg_restore artifacts and provenance counts match"
)