Files
wyndham-ARR/booking_ingestion/excel_review_postgres.py
2026-07-31 15:11:42 +08:00

945 lines
35 KiB
Python

"""Durable PostgreSQL review workflow for Booking Excel extraction drafts."""
from __future__ import annotations
import hashlib
import json
import uuid
from collections import OrderedDict
from datetime import datetime
from typing import Any, Callable, Optional
from arr_ingestion.contracts import ArtifactRef
from .excel import (
PROCESSOR_NAME,
PROCESSOR_VERSION,
RESULT_SCHEMA_VERSION,
REVIEW_CONFIRMED,
XLSX_MIME,
ExcelDocument,
ExcelRow,
ExtractedRoomItem,
build_excel_parse_result,
canonical_json_bytes,
rule_set_sha256,
)
from .excel_postgres import (
ARTIFACT_BUCKET_ALIAS,
ARTIFACT_STORAGE_PROVIDER,
BookingExcelRepositoryError,
BookingSourceSummary,
DatabaseConfig,
PostgresBookingExcelRepository,
)
from .excel_review import (
BookingDraftItem,
BookingDraftPage,
BookingDraftSummary,
validate_draft_id,
validate_review_item_id,
validate_review_item_ids,
validate_review_quantity,
validate_review_room_type,
)
REVIEW_SCHEMA_VERSION = "1.0"
REVIEW_LOCK_KEY = "arr_booking_excel_review"
class PostgresBookingReviewRepository(PostgresBookingExcelRepository):
"""Persist editable extraction drafts and activate reviewed rows atomically."""
def __init__(
self,
config: DatabaseConfig,
connect: Optional[Callable[[str], Any]] = None,
) -> None:
super().__init__(config, connect)
@staticmethod
def _review_ready(cursor: Any) -> None:
cursor.execute(
"""
SELECT
to_regclass('booking.extraction_drafts'),
to_regclass('booking.extraction_draft_items')
"""
)
row = cursor.fetchone()
if not row or row[0] is None or row[1] is None:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DATABASE_NOT_READY",
"Booking Excel 复核数据库迁移尚未就绪",
)
def _begin_review(self, cursor: Any, *, read_only: bool = False) -> None:
self._begin(cursor, read_only=read_only)
self._review_ready(cursor)
@staticmethod
def _lock(cursor: Any) -> None:
cursor.execute(
"SELECT pg_advisory_xact_lock(hashtextextended(%s, 0))",
(REVIEW_LOCK_KEY,),
)
@staticmethod
def _draft_summary_from_row(row: object) -> Optional[BookingDraftSummary]:
if row is None:
return None
values = tuple(row) # type: ignore[arg-type]
return BookingDraftSummary(
draft_id=str(values[0]),
filename=str(values[1]),
status=str(values[2]),
source_rows=int(values[3]),
worksheet_count=int(values[4]),
distinct_group_codes=int(values[5]),
extracted_items=int(values[6]),
confirmed_items=int(values[7]),
pending_items=int(values[8]),
deleted_items=int(values[9]),
confirmed_room_quantity=int(values[10]),
created_at=values[11] if isinstance(values[11], datetime) else None,
updated_at=values[12] if isinstance(values[12], datetime) else None,
)
@staticmethod
def _draft_summary(cursor: Any, draft_id: str) -> BookingDraftSummary:
cursor.execute(
"""
SELECT
draft.draft_id,
artifact.original_filename,
draft.draft_status,
draft.source_rows,
draft.worksheet_count,
draft.distinct_group_codes,
draft.extracted_items,
count(item.id) FILTER (
WHERE item.review_status = 'confirmed'
)::bigint,
count(item.id) FILTER (
WHERE item.review_status = 'pending'
)::bigint,
count(item.id) FILTER (
WHERE item.review_status = 'deleted'
)::bigint,
COALESCE(sum(item.quantity) FILTER (
WHERE item.review_status = 'confirmed'
), 0)::bigint,
draft.created_at,
draft.updated_at
FROM booking.extraction_drafts AS draft
JOIN ingestion.artifacts AS artifact
ON artifact.id = draft.source_artifact_id
LEFT JOIN booking.extraction_draft_items AS item
ON item.draft_id = draft.draft_id
WHERE draft.draft_id = %s
GROUP BY draft.draft_id, artifact.original_filename
""",
(draft_id,),
)
summary = PostgresBookingReviewRepository._draft_summary_from_row(
cursor.fetchone()
)
if summary is None:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DRAFT_NOT_FOUND",
"找不到 Excel 提取草稿",
404,
)
return summary
@staticmethod
def _latest_reviewing_draft_id(cursor: Any) -> Optional[str]:
cursor.execute(
"""
SELECT draft_id
FROM booking.extraction_drafts
WHERE draft_status = 'reviewing'
ORDER BY updated_at DESC, draft_id DESC
LIMIT 1
"""
)
row = cursor.fetchone()
return str(row[0]) if row is not None else None
@staticmethod
def _draft_item_from_row(row: object) -> BookingDraftItem:
values = tuple(row) # type: ignore[arg-type]
return BookingDraftItem(
item_id=int(values[0]),
draft_id=str(values[1]),
worksheet=str(values[2]),
row_no=int(values[3]),
item_no=int(values[4]),
tour_code=str(values[5]),
room_type_raw=str(values[6]),
room_type_code=str(values[7]) if values[7] else None,
quantity=int(values[8]),
review_status=str(values[9]),
automatic=bool(values[10]),
source_fragment=str(values[11] or ""),
updated_at=values[12] if isinstance(values[12], datetime) else None,
)
@classmethod
def _draft_page(
cls,
cursor: Any,
draft_id: str,
limit: int,
offset: int,
) -> BookingDraftPage:
summary = cls._draft_summary(cursor, draft_id)
cursor.execute(
"""
SELECT count(*)::bigint
FROM booking.extraction_draft_items
WHERE draft_id = %s
AND review_status <> 'deleted'
""",
(draft_id,),
)
total = int(cursor.fetchone()[0])
cursor.execute(
"""
SELECT
id,
draft_id,
source_worksheet,
source_row_no,
item_no,
group_code_raw,
room_type_raw,
room_type_code,
quantity,
review_status,
automatic,
source_fragment,
updated_at
FROM booking.extraction_draft_items
WHERE draft_id = %s
AND review_status <> 'deleted'
ORDER BY
CASE review_status WHEN 'pending' THEN 0 ELSE 1 END,
source_worksheet,
source_row_no,
item_no
LIMIT %s OFFSET %s
""",
(draft_id, limit, offset),
)
return BookingDraftPage(
summary=summary,
items=tuple(cls._draft_item_from_row(row) for row in cursor.fetchall()),
total=total,
limit=limit,
offset=offset,
)
@staticmethod
def _artifact_id(cursor: Any, source: ArtifactRef) -> int:
"""Register one source object; tolerate same-byte re-extraction."""
cursor.execute(
"""
SELECT id, byte_size, mime_type
FROM ingestion.artifacts
WHERE artifact_kind = 'booking_excel'
AND sha256 = %s
ORDER BY id
LIMIT 1
FOR SHARE
""",
(source.sha256,),
)
existing = cursor.fetchone()
if existing is not None:
if int(existing[1]) != source.byte_size or str(existing[2]) != source.mime_type:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_ARTIFACT_CONFLICT",
"Excel 工件身份冲突",
409,
)
return int(existing[0])
cursor.execute(
"""
INSERT INTO ingestion.artifacts (
artifact_kind,
storage_provider,
bucket_alias,
object_key,
original_filename,
sha256,
byte_size,
mime_type
)
VALUES ('booking_excel', %s, %s, %s, %s, %s, %s, %s)
RETURNING id
""",
(
ARTIFACT_STORAGE_PROVIDER,
ARTIFACT_BUCKET_ALIAS,
source.object_key,
source.original_filename,
source.sha256,
source.byte_size,
source.mime_type,
),
)
return int(cursor.fetchone()[0])
def current_draft(
self,
*,
limit: int = 50,
offset: int = 0,
) -> Optional[BookingDraftPage]:
connection = self._open()
try:
with connection.cursor() as cursor:
self._begin_review(cursor, read_only=True)
draft_id = self._latest_reviewing_draft_id(cursor)
if draft_id is None:
return None
return self._draft_page(cursor, draft_id, limit, offset)
except BookingExcelRepositoryError:
raise
except Exception:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DATABASE_UNAVAILABLE",
"Excel 提取草稿暂时无法读取",
) from None
finally:
try:
connection.rollback()
finally:
connection.close()
def create_draft(
self,
draft_id: str,
source: ArtifactRef,
document: ExcelDocument,
*,
limit: int = 50,
) -> BookingDraftPage:
draft_id = validate_draft_id(draft_id)
if (
source.role != "booking_source"
or source.file_kind != "booking_excel"
or source.mime_type != XLSX_MIME
or source.sha256 != document.source_sha256
or source.byte_size != document.source_byte_size
):
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_ARTIFACT_INVALID",
"Excel 工件身份无效",
422,
)
connection = self._open()
try:
with connection.cursor() as cursor:
self._begin_review(cursor)
self._lock(cursor)
artifact_id = self._artifact_id(cursor, source)
cursor.execute(
"""
SELECT draft_id, draft_status
FROM booking.extraction_drafts
WHERE source_artifact_id = %s
FOR UPDATE
""",
(artifact_id,),
)
prior = cursor.fetchone()
if prior is not None and str(prior[1]) == "activated":
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_SOURCE_ALREADY_ACTIVATED",
"这份 Excel 已经启用,无需重复提取",
409,
)
if prior is not None:
cursor.execute(
"DELETE FROM booking.extraction_draft_items WHERE draft_id = %s",
(str(prior[0]),),
)
cursor.execute(
"DELETE FROM booking.extraction_drafts WHERE draft_id = %s",
(str(prior[0]),),
)
cursor.execute(
"""
UPDATE booking.extraction_drafts
SET draft_status = 'superseded', updated_at = now()
WHERE draft_status = 'reviewing'
"""
)
cursor.execute(
"""
INSERT INTO booking.extraction_drafts (
draft_id,
source_artifact_id,
draft_status,
processor_version,
rule_set_sha256,
source_sha256,
source_byte_size,
source_rows,
worksheet_count,
distinct_group_codes,
extracted_items
)
VALUES (%s, %s, 'reviewing', %s, %s, %s, %s, %s, %s, %s, %s)
""",
(
draft_id,
artifact_id,
PROCESSOR_VERSION,
rule_set_sha256(),
document.source_sha256,
document.source_byte_size,
len(document.rows),
document.worksheet_count,
document.distinct_group_code_count,
document.extracted_item_count,
),
)
for row in document.rows:
for item in row.room_items:
cursor.execute(
"""
INSERT INTO booking.extraction_draft_items (
draft_id,
source_worksheet,
source_row_no,
item_no,
group_code_raw,
hotel_raw,
room_type_raw,
room_type_code,
quantity,
review_status,
automatic,
source_fragment
)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
""",
(
draft_id,
row.worksheet,
row.row_no,
item.item_no,
row.group_code_raw,
row.hotel_raw,
item.room_type_raw,
item.room_type_code,
item.quantity,
"pending" if item.needs_review else "confirmed",
not item.needs_review,
item.source_fragment,
),
)
page = self._draft_page(cursor, draft_id, limit, 0)
connection.commit()
return page
except BookingExcelRepositoryError:
connection.rollback()
raise
except Exception:
connection.rollback()
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DATABASE_WRITE_FAILED",
"Excel 已提取,但复核草稿未能保存",
) from None
finally:
connection.close()
def update_item(
self,
draft_id: str,
item_id: int,
room_type_code: str,
quantity: int,
) -> BookingDraftItem:
draft_id = validate_draft_id(draft_id)
item_id = validate_review_item_id(item_id)
room_type_code = validate_review_room_type(room_type_code)
quantity = validate_review_quantity(quantity)
connection = self._open()
try:
with connection.cursor() as cursor:
self._begin_review(cursor)
self._lock(cursor)
cursor.execute(
"""
UPDATE booking.extraction_draft_items AS item
SET room_type_code = %s,
quantity = %s,
review_status = 'confirmed',
automatic = false,
updated_at = now()
FROM booking.extraction_drafts AS draft
WHERE item.id = %s
AND item.draft_id = %s
AND item.review_status <> 'deleted'
AND draft.draft_id = item.draft_id
AND draft.draft_status = 'reviewing'
RETURNING
item.id,
item.draft_id,
item.source_worksheet,
item.source_row_no,
item.item_no,
item.group_code_raw,
item.room_type_raw,
item.room_type_code,
item.quantity,
item.review_status,
item.automatic,
item.source_fragment,
item.updated_at
""",
(room_type_code, quantity, item_id, draft_id),
)
row = cursor.fetchone()
if row is None:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_REVIEW_ITEM_NOT_FOUND",
"找不到可编辑的提取记录",
404,
)
cursor.execute(
"UPDATE booking.extraction_drafts SET updated_at = now() WHERE draft_id = %s",
(draft_id,),
)
item = self._draft_item_from_row(row)
connection.commit()
return item
except BookingExcelRepositoryError:
connection.rollback()
raise
except Exception:
connection.rollback()
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DATABASE_WRITE_FAILED",
"提取记录未能保存",
) from None
finally:
connection.close()
def delete_item(self, draft_id: str, item_id: int) -> None:
item_id = validate_review_item_id(item_id)
self.delete_items(draft_id, (item_id,))
def delete_items(self, draft_id: str, item_ids: object) -> None:
draft_id = validate_draft_id(draft_id)
validated_ids = validate_review_item_ids(item_ids)
connection = self._open()
try:
with connection.cursor() as cursor:
self._begin_review(cursor)
self._lock(cursor)
cursor.execute(
"""
UPDATE booking.extraction_draft_items AS item
SET review_status = 'deleted',
automatic = false,
updated_at = now()
FROM booking.extraction_drafts AS draft
WHERE item.id = ANY(%s::bigint[])
AND item.draft_id = %s
AND item.review_status <> 'deleted'
AND draft.draft_id = item.draft_id
AND draft.draft_status = 'reviewing'
RETURNING item.id
""",
(list(validated_ids), draft_id),
)
deleted_ids = {int(row[0]) for row in cursor.fetchall()}
if deleted_ids != set(validated_ids):
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_REVIEW_ITEM_NOT_FOUND",
"部分所选记录已不存在或不可删除,请刷新后重试",
404,
)
cursor.execute(
"UPDATE booking.extraction_drafts SET updated_at = now() WHERE draft_id = %s",
(draft_id,),
)
connection.commit()
except BookingExcelRepositoryError:
connection.rollback()
raise
except Exception:
connection.rollback()
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DATABASE_WRITE_FAILED",
"所选提取记录未能删除",
) from None
finally:
connection.close()
def discard_draft(self, draft_id: str) -> None:
draft_id = validate_draft_id(draft_id)
connection = self._open()
try:
with connection.cursor() as cursor:
self._begin_review(cursor)
self._lock(cursor)
cursor.execute(
"""
UPDATE booking.extraction_drafts
SET draft_status = 'superseded', updated_at = now()
WHERE draft_id = %s
AND draft_status = 'reviewing'
RETURNING draft_id
""",
(draft_id,),
)
if cursor.fetchone() is None:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DRAFT_NOT_FOUND",
"找不到可放弃的 Excel 提取草稿",
404,
)
connection.commit()
except BookingExcelRepositoryError:
connection.rollback()
raise
except Exception:
connection.rollback()
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DATABASE_WRITE_FAILED",
"Excel 提取草稿未能放弃",
) from None
finally:
connection.close()
def activate_draft(self, draft_id: str) -> BookingSourceSummary:
draft_id = validate_draft_id(draft_id)
connection = self._open()
try:
with connection.cursor() as cursor:
self._begin_review(cursor)
self._lock(cursor)
cursor.execute(
"""
SELECT draft.source_artifact_id, draft.source_sha256
FROM booking.extraction_drafts AS draft
WHERE draft.draft_id = %s
AND draft.draft_status = 'reviewing'
FOR UPDATE OF draft
""",
(draft_id,),
)
draft = cursor.fetchone()
if draft is None:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DRAFT_NOT_FOUND",
"找不到可确认的 Excel 提取草稿",
404,
)
cursor.execute(
"""
SELECT count(*)::bigint
FROM booking.extraction_draft_items
WHERE draft_id = %s
AND review_status = 'pending'
""",
(draft_id,),
)
if int(cursor.fetchone()[0]) > 0:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_REVIEW_REQUIRED",
"仍有待确认房型,请全部确认或删除后再启用",
409,
)
cursor.execute(
"""
SELECT
source_worksheet,
source_row_no,
group_code_raw,
hotel_raw,
room_type_raw,
room_type_code,
quantity,
source_fragment
FROM booking.extraction_draft_items
WHERE draft_id = %s
AND review_status = 'confirmed'
ORDER BY source_worksheet, source_row_no, item_no
""",
(draft_id,),
)
grouped: OrderedDict[
tuple[str, int, str, str], list[ExtractedRoomItem]
] = OrderedDict()
for values in cursor.fetchall():
key = (str(values[0]), int(values[1]), str(values[2]), str(values[3]))
items = grouped.setdefault(key, [])
items.append(
ExtractedRoomItem(
item_no=len(items) + 1,
room_type_raw=str(values[4]),
room_type_code=str(values[5]),
quantity=int(values[6]),
review_status=REVIEW_CONFIRMED,
source_fragment=str(values[7] or ""),
)
)
if not grouped:
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_REVIEW_EMPTY",
"没有可启用的房型记录",
409,
)
rows = tuple(
ExcelRow(
worksheet=key[0],
row_no=key[1],
group_code_raw=key[2],
hotel_raw=key[3],
room_items=tuple(items),
)
for key, items in grouped.items()
)
artifact_id = int(draft[0])
delivery = {
"draft_id": draft_id,
"source_sha256": str(draft[1]),
"source_rows": len(rows),
"room_items": sum(len(row.room_items) for row in rows),
"room_quantity": sum(row.no_of_rooms for row in rows),
}
delivery_bytes = canonical_json_bytes(delivery)
delivery_sha256 = hashlib.sha256(delivery_bytes).hexdigest()
run_key = "bookingjob-" + uuid.uuid4().hex
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_set_sha256(),
PROCESSOR_VERSION,
rule_set_sha256(),
RESULT_SCHEMA_VERSION,
delivery_sha256,
delivery_bytes.decode("utf-8"),
),
)
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, 'booking_excel', %s, 'processing', %s, 0, 0)
RETURNING id
""",
(artifact_id, f"xlsx-reviewed/{REVIEW_SCHEMA_VERSION}", len(rows)),
)
source_batch_id = int(cursor.fetchone()[0])
for row in rows:
row_sha256 = row.sha256()
result = build_excel_parse_result(row)
result_bytes = canonical_json_bytes(result)
result_sha256 = hashlib.sha256(result_bytes).hexdigest()
cursor.execute(
"""
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
""",
(
source_batch_id,
row.worksheet,
row.row_no,
row.group_code_raw,
row.hotel_raw,
row.no_of_rooms,
row_sha256,
),
)
source_row_id = int(cursor.fetchone()[0])
cursor.execute(
"""
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
)
VALUES (%s, 1, 'accepted', %s, %s, %s, %s, %s, %s, %s::jsonb, now())
RETURNING id
""",
(
source_row_id,
RESULT_SCHEMA_VERSION,
PROCESSOR_NAME,
PROCESSOR_VERSION,
rule_set_sha256(),
row_sha256,
result_sha256,
result_bytes.decode("utf-8"),
),
)
parse_version_id = int(cursor.fetchone()[0])
for item in row.room_items:
cursor.execute(
"""
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
)
VALUES (%s, %s, %s, %s, %s, NULL, NULL, NULL, %s)
""",
(
parse_version_id,
item.item_no,
item.room_type_raw,
item.room_type_code,
item.quantity,
item.source_fragment,
),
)
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,),
)
self._activate(cursor, 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(
"""
UPDATE booking.extraction_drafts
SET draft_status = 'activated',
activated_source_batch_id = %s,
activated_at = now(),
updated_at = now()
WHERE draft_id = %s
""",
(source_batch_id, draft_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(rows),
"draft_id": draft_id,
},
separators=(",", ":"),
),
),
)
summary = self._batch_summary(
cursor,
source_batch_id,
"reviewed_and_activated",
)
connection.commit()
return summary
except BookingExcelRepositoryError:
connection.rollback()
raise
except Exception:
connection.rollback()
raise BookingExcelRepositoryError(
"BOOKING_EXCEL_DATABASE_WRITE_FAILED",
"Excel 复核结果未能启用",
) from None
finally:
connection.close()