436 lines
15 KiB
Python
436 lines
15 KiB
Python
"""PostgreSQL adapter for the approved booking Markdown fixture."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
from .md_fixture import (
|
|
PROCESSOR_NAME,
|
|
PROCESSOR_VERSION,
|
|
RESULT_SCHEMA_VERSION,
|
|
FixtureDocument,
|
|
build_fixture_parse_result,
|
|
canonical_json_bytes,
|
|
parse_fixture_file,
|
|
rule_set_sha256,
|
|
)
|
|
|
|
|
|
TARGET_DATABASE = "booking_test"
|
|
EXPECTED_FIXTURE_ROWS = 867
|
|
EXPECTED_FIXTURE_WORKSHEETS = 6
|
|
EXPECTED_FIXTURE_GROUP_CODES = 348
|
|
PROJECT_ROOT = Path(__file__).resolve().parents[1]
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class DatabaseConfig:
|
|
host: str
|
|
port: int
|
|
user: str
|
|
password: str
|
|
database: str
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ImportResult:
|
|
status: str
|
|
source_batch_id: int
|
|
source_rows: int
|
|
current_parses: int
|
|
room_items: int
|
|
distinct_group_codes: int
|
|
|
|
|
|
def _artifact_id(
|
|
cursor,
|
|
*,
|
|
fixture_path: Path,
|
|
content_sha256: str,
|
|
byte_size: int,
|
|
) -> int:
|
|
cursor.execute(
|
|
"""
|
|
SELECT id, byte_size
|
|
FROM ingestion.artifacts
|
|
WHERE artifact_kind = 'booking_source_md'
|
|
AND sha256 = %s
|
|
""",
|
|
(content_sha256,),
|
|
)
|
|
existing = cursor.fetchone()
|
|
if existing is not None:
|
|
if int(existing[1]) != byte_size:
|
|
raise RuntimeError("existing booking artifact metadata conflicts")
|
|
return int(existing[0])
|
|
|
|
try:
|
|
object_key = fixture_path.resolve().relative_to(PROJECT_ROOT).as_posix()
|
|
except ValueError:
|
|
raise RuntimeError("booking fixture must be inside the project root") from None
|
|
cursor.execute(
|
|
"""
|
|
INSERT INTO ingestion.artifacts (
|
|
artifact_kind, storage_provider, bucket_alias, object_key,
|
|
original_filename, sha256, byte_size, mime_type
|
|
)
|
|
VALUES (
|
|
'booking_source_md', 'local_fixture', 'workspace', %s,
|
|
%s, %s, %s, 'text/markdown'
|
|
)
|
|
RETURNING id
|
|
""",
|
|
(object_key, fixture_path.name, content_sha256, byte_size),
|
|
)
|
|
return int(cursor.fetchone()[0])
|
|
|
|
|
|
def _existing_result(
|
|
cursor,
|
|
*,
|
|
artifact_id: int,
|
|
expected: FixtureDocument,
|
|
) -> ImportResult | None:
|
|
cursor.execute(
|
|
"""
|
|
SELECT id, batch_status, source_rows, accepted_rows, failed_rows
|
|
FROM booking.source_batches
|
|
WHERE source_artifact_id = %s
|
|
""",
|
|
(artifact_id,),
|
|
)
|
|
batch = cursor.fetchone()
|
|
if batch is None:
|
|
return None
|
|
batch_id = int(batch[0])
|
|
if (
|
|
batch[1] != "accepted"
|
|
or int(batch[2]) != len(expected.rows)
|
|
or int(batch[3]) != len(expected.rows)
|
|
or int(batch[4]) != 0
|
|
):
|
|
raise RuntimeError("existing booking source batch is incomplete")
|
|
cursor.execute(
|
|
"""
|
|
SELECT
|
|
count(*) AS source_rows,
|
|
count(current_parse.source_row_id) AS current_parses,
|
|
count(item.id) AS room_items,
|
|
count(DISTINCT source.group_code_key) AS distinct_group_codes
|
|
FROM booking.source_rows AS source
|
|
LEFT JOIN booking.current_row_parses AS current_parse
|
|
ON current_parse.source_row_id = source.id
|
|
LEFT JOIN booking.room_items AS item
|
|
ON item.parse_version_id = current_parse.parse_version_id
|
|
WHERE source.source_batch_id = %s
|
|
""",
|
|
(batch_id,),
|
|
)
|
|
counts = tuple(int(value) for value in cursor.fetchone())
|
|
expected_counts = (
|
|
len(expected.rows),
|
|
len(expected.rows),
|
|
len(expected.rows),
|
|
expected.distinct_group_code_count,
|
|
)
|
|
if counts != expected_counts:
|
|
raise RuntimeError("existing booking fixture rows are inconsistent")
|
|
return ImportResult(
|
|
status="already_applied_and_verified",
|
|
source_batch_id=batch_id,
|
|
source_rows=counts[0],
|
|
current_parses=counts[1],
|
|
room_items=counts[2],
|
|
distinct_group_codes=counts[3],
|
|
)
|
|
|
|
|
|
def import_fixture(
|
|
config: DatabaseConfig,
|
|
fixture_path: Path,
|
|
contract_path: Path,
|
|
) -> ImportResult:
|
|
if config.database != TARGET_DATABASE:
|
|
raise RuntimeError("database target is invalid")
|
|
fixture_path = fixture_path.resolve()
|
|
contract_path = contract_path.resolve()
|
|
document = parse_fixture_file(fixture_path)
|
|
if (
|
|
len(document.rows) != EXPECTED_FIXTURE_ROWS
|
|
or document.worksheet_count != EXPECTED_FIXTURE_WORKSHEETS
|
|
or document.distinct_group_code_count
|
|
!= EXPECTED_FIXTURE_GROUP_CODES
|
|
):
|
|
raise RuntimeError("approved booking fixture baseline changed")
|
|
|
|
fixture_bytes = fixture_path.read_bytes()
|
|
fixture_sha256 = hashlib.sha256(fixture_bytes).hexdigest()
|
|
rule_sha256 = rule_set_sha256(contract_path)
|
|
run_key = f"booking-fixture:{fixture_sha256[:24]}"
|
|
delivery_payload = {
|
|
"source_kind": "expected_fixture",
|
|
"source_rows": len(document.rows),
|
|
"worksheet_count": document.worksheet_count,
|
|
"distinct_group_codes": document.distinct_group_code_count,
|
|
"source_file": document.source_file,
|
|
"mapping_file": document.mapping_file,
|
|
"contract": f"booking-row-parse-result/{RESULT_SCHEMA_VERSION}",
|
|
}
|
|
delivery_json = canonical_json_bytes(delivery_payload).decode("utf-8")
|
|
delivery_sha256 = hashlib.sha256(
|
|
delivery_json.encode("utf-8")
|
|
).hexdigest()
|
|
|
|
try:
|
|
import psycopg
|
|
except ImportError:
|
|
raise RuntimeError("psycopg is required for PostgreSQL import") from None
|
|
|
|
connection = psycopg.connect(
|
|
host=config.host,
|
|
port=config.port,
|
|
user=config.user,
|
|
password=config.password,
|
|
dbname=config.database,
|
|
connect_timeout=5,
|
|
options=(
|
|
"-c statement_timeout=120000 "
|
|
"-c lock_timeout=5000 "
|
|
"-c idle_in_transaction_session_timeout=180000 "
|
|
"-c application_name=arr_booking_fixture_import"
|
|
),
|
|
)
|
|
try:
|
|
with connection.cursor() as cursor:
|
|
cursor.execute("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE")
|
|
cursor.execute(
|
|
"SELECT current_database(), current_setting('transaction_read_only')"
|
|
)
|
|
if cursor.fetchone() != (TARGET_DATABASE, "off"):
|
|
raise RuntimeError("database target/write guard failed")
|
|
cursor.execute(
|
|
"SELECT pg_advisory_xact_lock(hashtextextended(%s, 0))",
|
|
("arr_booking_fixture_import",),
|
|
)
|
|
cursor.execute("SELECT to_regclass('booking.source_rows')")
|
|
if cursor.fetchone()[0] is None:
|
|
raise RuntimeError("ARR MVP v1 booking schema is not installed")
|
|
|
|
artifact_id = _artifact_id(
|
|
cursor,
|
|
fixture_path=fixture_path,
|
|
content_sha256=fixture_sha256,
|
|
byte_size=len(fixture_bytes),
|
|
)
|
|
existing = _existing_result(
|
|
cursor,
|
|
artifact_id=artifact_id,
|
|
expected=document,
|
|
)
|
|
if existing is not None:
|
|
connection.rollback()
|
|
return existing
|
|
|
|
cursor.execute(
|
|
"""
|
|
INSERT INTO ingestion.processing_runs (
|
|
run_key, pipeline_type, source_artifact_id, run_status,
|
|
requested_processor_version, requested_rule_set_sha256,
|
|
delivered_processor_version, delivered_rule_set_sha256,
|
|
result_schema_version, delivery_sha256, delivery_json
|
|
)
|
|
VALUES (
|
|
%s, 'booking_source_import', %s, 'validating',
|
|
%s, %s, %s, %s, %s, %s, %s::jsonb
|
|
)
|
|
RETURNING id
|
|
""",
|
|
(
|
|
run_key,
|
|
artifact_id,
|
|
PROCESSOR_VERSION,
|
|
rule_sha256,
|
|
PROCESSOR_VERSION,
|
|
rule_sha256,
|
|
RESULT_SCHEMA_VERSION,
|
|
delivery_sha256,
|
|
delivery_json,
|
|
),
|
|
)
|
|
processing_run_id = int(cursor.fetchone()[0])
|
|
cursor.execute(
|
|
"""
|
|
INSERT INTO booking.source_batches (
|
|
source_artifact_id, source_kind, source_format_version,
|
|
batch_status, source_rows, accepted_rows, failed_rows
|
|
)
|
|
VALUES (
|
|
%s, 'expected_fixture', 'md/1.0',
|
|
'processing', %s, 0, 0
|
|
)
|
|
RETURNING id
|
|
""",
|
|
(artifact_id, len(document.rows)),
|
|
)
|
|
source_batch_id = int(cursor.fetchone()[0])
|
|
|
|
for row in document.rows:
|
|
row_sha256 = row.sha256()
|
|
parse_result = build_fixture_parse_result(
|
|
row,
|
|
rule_sha256=rule_sha256,
|
|
)
|
|
result_bytes = canonical_json_bytes(parse_result)
|
|
result_sha256 = hashlib.sha256(result_bytes).hexdigest()
|
|
cursor.execute(
|
|
"""
|
|
WITH source AS (
|
|
INSERT INTO booking.source_rows (
|
|
source_batch_id, source_worksheet, source_row_no,
|
|
group_code_raw, type_of_room_raw, no_of_rooms,
|
|
source_row_sha256
|
|
)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING id
|
|
),
|
|
parsed AS (
|
|
INSERT INTO booking.parse_versions (
|
|
source_row_id, version_no, parse_status,
|
|
result_schema_version, processor_name,
|
|
processor_version, rule_set_sha256,
|
|
input_row_sha256, result_sha256,
|
|
result_json, validated_at
|
|
)
|
|
SELECT
|
|
source.id, 1, 'accepted', %s, %s, %s, %s,
|
|
%s, %s, %s::jsonb, now()
|
|
FROM source
|
|
RETURNING id, source_row_id
|
|
),
|
|
item AS (
|
|
INSERT INTO booking.room_items (
|
|
parse_version_id, item_no, room_type_raw,
|
|
room_type_code, quantity, unit_price,
|
|
currency_code, price_token_raw, source_fragment
|
|
)
|
|
SELECT
|
|
parsed.id, 1, %s, %s, %s,
|
|
NULL, NULL, NULL, %s
|
|
FROM parsed
|
|
RETURNING id
|
|
)
|
|
SELECT source_row_id, id
|
|
FROM parsed
|
|
""",
|
|
(
|
|
source_batch_id,
|
|
row.worksheet,
|
|
row.row_no,
|
|
row.group_code_raw,
|
|
row.type_of_room_raw,
|
|
row.no_of_rooms,
|
|
row_sha256,
|
|
RESULT_SCHEMA_VERSION,
|
|
PROCESSOR_NAME,
|
|
PROCESSOR_VERSION,
|
|
rule_sha256,
|
|
row_sha256,
|
|
result_sha256,
|
|
result_bytes.decode("utf-8"),
|
|
row.type_of_room_raw,
|
|
row.type_of_room_raw,
|
|
row.no_of_rooms,
|
|
row.type_of_room_raw,
|
|
),
|
|
)
|
|
|
|
cursor.execute(
|
|
"""
|
|
UPDATE booking.source_batches
|
|
SET batch_status = 'accepted',
|
|
accepted_rows = source_rows,
|
|
failed_rows = 0,
|
|
validated_at = now(),
|
|
finished_at = now()
|
|
WHERE id = %s
|
|
""",
|
|
(source_batch_id,),
|
|
)
|
|
cursor.execute(
|
|
"""
|
|
INSERT INTO booking.current_row_parses (
|
|
source_row_id, parse_version_id
|
|
)
|
|
SELECT source.id, parsed.id
|
|
FROM booking.source_rows AS source
|
|
JOIN booking.parse_versions AS parsed
|
|
ON parsed.source_row_id = source.id
|
|
AND parsed.version_no = 1
|
|
AND parsed.parse_status = 'accepted'
|
|
WHERE source.source_batch_id = %s
|
|
""",
|
|
(source_batch_id,),
|
|
)
|
|
cursor.execute(
|
|
"""
|
|
UPDATE ingestion.processing_runs
|
|
SET run_status = 'accepted',
|
|
updated_at = now(),
|
|
validated_at = now(),
|
|
finished_at = now()
|
|
WHERE id = %s
|
|
""",
|
|
(processing_run_id,),
|
|
)
|
|
cursor.execute(
|
|
"""
|
|
INSERT INTO ingestion.outbox_events (
|
|
event_key, aggregate_type, aggregate_id,
|
|
event_type, payload
|
|
)
|
|
VALUES (
|
|
%s, 'booking_source_batch', %s,
|
|
'booking.source_batch.accepted', %s::jsonb
|
|
)
|
|
""",
|
|
(
|
|
f"booking-source-batch:{source_batch_id}:accepted",
|
|
source_batch_id,
|
|
json.dumps(
|
|
{
|
|
"source_batch_id": source_batch_id,
|
|
"source_rows": len(document.rows),
|
|
"distinct_group_codes":
|
|
document.distinct_group_code_count,
|
|
},
|
|
separators=(",", ":"),
|
|
),
|
|
),
|
|
)
|
|
result = _existing_result(
|
|
cursor,
|
|
artifact_id=artifact_id,
|
|
expected=document,
|
|
)
|
|
if result is None:
|
|
raise RuntimeError("booking fixture post-write validation failed")
|
|
connection.commit()
|
|
except Exception:
|
|
connection.rollback()
|
|
raise
|
|
finally:
|
|
connection.close()
|
|
|
|
return ImportResult(
|
|
status="imported_and_verified",
|
|
source_batch_id=result.source_batch_id,
|
|
source_rows=result.source_rows,
|
|
current_parses=result.current_parses,
|
|
room_items=result.room_items,
|
|
distinct_group_codes=result.distinct_group_codes,
|
|
)
|