diff --git a/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md b/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md
index b89ea55..9abeffd 100644
--- a/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md
+++ b/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md
@@ -12,7 +12,7 @@
## Scope
-- 修复生产 OHIP 查询、环境内任务恢复及接口字段人工完善;保留 XML、原费率筛选、定价、日报/月报规则,并落实用户确认的取消预订排除规则。
+- 修复生产 OHIP 查询、环境内任务恢复及接口字段人工完善;保留 XML、原费率筛选、定价、日报/月报规则,并落实用户确认的取消预订及 PM 房型排除规则。
## Intent And Constraints
@@ -255,3 +255,18 @@ Read: memory-index, project-positioning, current-state latest September sections
- User asks how PM was identified and what characteristics support it. Same-task ownership/feature/codex/base2417b1a/owned checkout/no peers and retained Project Context Loaded verified; Planning gate Passed for read-only source-evidence inspection. No new source request, report/rule change or hotel configuration assumption is authorized by this question.
- Reopened the saved original search/detail responses for all10 API0 price-unmatched rows. In every row search.roomStay.roomType, detail.roomStay.currentRoomInfo.roomType and detail.roomStay.roomRates[the9/16day].roomType explicitly equal PM. This is the same3-way agreement used by source_fields.agreed_room_type, not an inference from room-number prefix, zero amount or roomTypeCharged. Two examples: confirmation300638606/room9002 and300638238/room9004, all3 PM fields agree; both also have amount0 and adults0/children0. These latter features are auxiliary, never a standalone pseudo-room rule.
- Broader source count23 PM records (19 CheckedOut,4 Cancelled) should not be confused with the10 actionable unmatched PM records discussed here. Nothing establishes that PM is the only pseudo code in this hotel, or that every9xxx room/0-price reservation is pseudo. The official Oracle default-PM meaning remains supported by the links above, while actual report eligibility still requires the original report criterion. No code/data/price edits, Oracle calls, test suite, service action or actual generation occurred. Explain explicit roomType evidence separately from the not-yet-settled decision to exclude it from the report.
+
+## Same-task Follow-up: User-approved PM Exclusion
+
+- User explicitly confirms PM must not be included and requests updating processing rules. This resolves the PM inclusion-policy question above; exclusion is now authorized for both XML and OHIP input, before required-field checks/deduplication/pricing. Only exact trimmed/case-insensitive ROOM_CATEGORY_LABEL=PM is approved; do not infer from0 price,9xxx rooms, absent type, roomTypeCharged or other pseudo codes. Cancellation remains the first mutually exclusive exclusion. Original source records and prior staff decisions remain auditable.
+- Concurrent task gate Passed: same20261008-production-review-9e7b feature/codex/owned checkout/codex/arr-production-review/base2417b1a, no peers; primary unknown docs untouched. Project Context Loaded: retained required memory-index/positioning/current-state/decision-index/ADR004,006,007/architecture/domain/evidence/reflection/commitments/stale context plus active task record/read-before-planning/planning-gate. New user authorization supersedes only the prior undecided-PM boundary. Relevant modules: shared processor/independent validator/schema, Finance validation/SQL contracts, upstream field review, local versioned runtime. Planning gate Passed.
+- Plan: add explicit excluded_pm audit outcome and processing identity4.4.0 across shared XML/data rules and independent checks; extend ingestion/Finance schema and counts without altering prior published versions; prevent PM source-field prompts; regression-check both input paths, exclusions before validation/dedup/pricing, zero-priced physical-room retention and legacy contracts. Replay9/16 immutable capture offline, then safely activate the local rule and supersede its obsolete price-review job while retaining any decisions on still-relevant keys. Do not refetch Oracle, fabricate prices, overwrite completed real reports or silently reprocess published10/7. Required local maintenance will be rehearsed/backup/revision-checked; record material preservation limitations before any unsafe change.
+
+- Implementation outcome: active processor4.4.0 shares the same explicit `is_pm_record` rule for XML and direct OHIP data. Existing text normalization runs before trimmed/case-insensitive PM comparison. Cancellation wins; PM then precedes rate/required fields/deduplication/pricing, and `excluded_pm` / `ROOM_TYPE_PM_EXCLUDED` retains every original/audit record with null price/channel facts. No price lookup, whitelist, normal zero-price or formula rule changed. Source review now uses the same normalized records for issues/counts/edit guards; prior PM decisions remain in immutable audit, with a manually completed room-type value still correctable before finalization. Chinese/English/Thai exclusion counts are visible. Package archives match their source.
+- Finance outcome: added migration021 and explicit startup checks for both input paths; new independent PM count reconciles source rows and validates the excluded record shape.4.4 requires all8 outcome counters;4.3/4.2/retired v3 replay retain their old contracts. Upgrade leaves historical versions/records intact with the new count defaulting to0; rollback refuses committed PM-exclusion facts, existing mutation guards and role access remain.
+- Verification: processor PM8 + cancellation9 distinct cases passed; ingestion/Finance61 distinct tests passed with disposable local PostgreSQL (no skips), including4.3→021 historical row preservation, both input paths, daily/monthly results, rejected-batch audit, idempotent retry, rollback and role checks. Root source-review/cancellation/executor run passed41 tests with1 PG class skipped before opt-in; the subsequently opted-in DirectFieldReviewPostgresTests and all6 final PM review cases passed7/7. Root package/runtime run passed28/28, JavaScript review passed34/34. Independent review identified differing emoji normalization in review counters/edit guard, fixed by sharing processor normalization and covered by a new PM😀 case; manual room-type correction is also covered. Preliminary runner name mistakes were corrected and are not represented as case failures or as passing tests. Git whitespace, node syntax, documentation structure and task ownership checks passed.
+- Offline real9/16 replay: unchanged209-row OHIP source yields6 excluded cancellations,19 excluded PM,34 excluded rates,115 priced candidates and35 price-unmatched records in exactly1 key LIAN TAI/WHO2/1000. The unchanged184-row original XML yields0 cancellations/PM,34 rate exclusions and the same115+35 candidates/one key. This is evidence for this date after the explicit new PM policy, not a claim that the separately deferred universal API/XML population question is solved. No actual price or source value was invented, and no Oracle call occurred.
+- Safe local activation: verified no active acquisition/processing queue, gracefully stopped the existing8875 service, cold-backed up the entire119MB private instance to `policy44-service-backup-20261008-142628`, then ran the privately rehearsed explicit upgrade. Schema021, instance hash, launcher and unfinished9/17 identity/receipt were advanced.9/17 remains editing revision4 with0 pending fields,18 cancellations and11 PM exclusions; no staff finalization occurred. Replayed9/16 from its original checkpoint using a source whose fetch method raises if called, creating a new immutable4.4 job/request; superseded/cancelled its old5-key review through the normal revision-checked cancellation method and preserved its1 user-saved0 decision in the old case/audit. No matching still-required price was saved, so the new1-key review remains open0/1.
+- Local preservation evidence: existing Finance snapshot before/after matches exactly (1 prior version and58 audit records, comparing all old columns); newly added excluded_pm_rows remains0 for that historical version. Verified25 original handoff/review files unchanged. This includes preservation of completed10/7 rather than automatic reprocessing. New9/16 request `d89e8c15f9c94bd1a63411fa88e935d7`, job `arrbatch-517912e100e4f0dd6afecf1410b3e028f31eeecd85a0d26d`; old job remains separately auditable/cancelled. Service restarted with all local health flags ready. No production hotel writes, remote deployment/push or primary-checkout integration occurred.
+- Product verification: authenticated API and CUA confirm9/16 one pending price key, Oracle amount1000,35 records/35 rooms/35 nights, blank processed price, disabled final generation; original XML upload remains present. Browser result tab retained as deliverable, screenshot `production-validation-20261007/pm-excluded-sept16-price-review-20261008.png`. Cold backup, raw replay, policy intent/completed receipt and upgrade script remain private outside Git. Actual9/16 daily/monthly generation awaits the user's price decision;0 remains accepted when explicitly confirmed.
+- Promotion/follow-up: integration owner should promote the approved cancellation→PM→rate→validation→dedup→pricing order and schema021 deployment requirement. Retain the distinction between19 PM scope exclusions and10 removed unmatched PM reservations/4 removed price keys. User can now fill the sole LIAN TAI/WHO2/1000 processed-price key, save and confirm generation through the existing daily/monthly flow. Do not transfer the old PM0 value onto this different key or silently reprocess prior published reports.
diff --git a/arr-opera-daily-ingest.skill b/arr-opera-daily-ingest.skill
index 0ca58d9..55fb430 100644
Binary files a/arr-opera-daily-ingest.skill and b/arr-opera-daily-ingest.skill differ
diff --git a/arr-opera-daily-ingest.zip b/arr-opera-daily-ingest.zip
index 0ca58d9..55fb430 100644
Binary files a/arr-opera-daily-ingest.zip and b/arr-opera-daily-ingest.zip differ
diff --git a/arr-opera-daily-ingest/SKILL.md b/arr-opera-daily-ingest/SKILL.md
index c563223..19ad0ce 100644
--- a/arr-opera-daily-ingest/SKILL.md
+++ b/arr-opera-daily-ingest/SKILL.md
@@ -63,7 +63,7 @@ Read `result.json` only after the process exits.
If Python or `openpyxl` is unavailable, stop with an infrastructure failure. Do not switch to a different spreadsheet implementation.
-## Direct data entry (processor 4.3.0)
+## Direct data entry (processor 4.4.0)
Use `--data-json` instead of `--xml` for the frozen OHIP data source. Both inputs use the same classification, pricing, review and daily workbook rules. JSON is never converted into XML. XML keeps result schema 4.0; direct data uses schema 5.0 with `ingestion_mode=ohip_data` and an explicit `source_data` / `ohip_json` artifact. The independent validator accepts the same source switch. See `references/data-result.schema.json` and `references/data-structured-result.schema.json`. Trace remains blank for the 15-field data input.
@@ -71,3 +71,7 @@ Use `--data-json` instead of `--xml` for the frozen OHIP data source. Both input
## Cancelled reservation scope (processor 4.3.0)
Exclude explicitly cancelled reservations before rate-code, required-field, pricing and room/date deduplication checks. Preserve every source row in the audit with `excluded_cancelled` and `RESERVATION_CANCELLED_EXCLUDED`; never request a room for an excluded reservation. Both XML and OHIP data use the same predicate. XML recognizes only explicit CXL/CANCELLED/CANCELED status values; unknown abbreviations and conflicting status fields are not inferred. Retired direct-MCP v3 replay retains its original scope.
+
+## PM room type scope (processor 4.4.0)
+
+After cancelled reservation exclusion, exclude records whose explicit `ROOM_CATEGORY_LABEL`, trimmed and compared case-insensitively, equals `PM`. Apply this before whitelist, required-field, deduplication and pricing checks in both XML and OHIP paths. Preserve every row as `excluded_pm` with `ROOM_TYPE_PM_EXCLUDED`; do not request fields or prices for it. Do not infer PM from room number, amount, occupancy, charged room type or other pseudo codes. Cancelled PM rows count only as cancelled. Retired direct-MCP v3 replay keeps its original scope.
diff --git a/arr-opera-daily-ingest/references/business-rules.md b/arr-opera-daily-ingest/references/business-rules.md
index 50c9ca8..1866f67 100644
--- a/arr-opera-daily-ingest/references/business-rules.md
+++ b/arr-opera-daily-ingest/references/business-rules.md
@@ -4,7 +4,7 @@
1. Validate the invocation and parse one fixed `RES_DETAIL` XML.
2. Require exactly one XML business date and keep one audit record for every `G_RESERVATION` in XML order.
-3. First exclude explicitly cancelled reservations as `excluded_cancelled` with `RESERVATION_CANCELLED_EXCLUDED`; preserve the original row and source sequence. Other records require `RATE_CODE` so whitelist membership is knowable.
+3. First exclude explicitly cancelled reservations as `excluded_cancelled` with `RESERVATION_CANCELLED_EXCLUDED`, then exclude an explicit `ROOM_CATEGORY_LABEL` of `PM` as `excluded_pm` with `ROOM_TYPE_PM_EXCLUDED`. Preserve the original row and source sequence. Other records require `RATE_CODE` so whitelist membership is knowable.
4. Classify a trimmed, uppercased rate code outside the whitelist as `excluded_rate_code`.
5. Validate every whitelist candidate; classify invalid rows as `validation_failed`.
6. Deduplicate valid candidates by `DISP_ROOM_NO + ARRIVAL`; retain the first XML occurrence and classify later occurrences as `duplicate` pointing to the first source sequence.
@@ -96,6 +96,12 @@ Do not access OSS, embed credentials, generate/update a monthly workbook, query/
## Cancellation exclusion approved 2026-10-08
-The active processor4.3.0 excludes cancelled reservations whether or not a room is assigned, before all business-field checks and room/date deduplication. OHIP uses `reservation_status`, retained from matching search/detail responses. XML uses the explicitly named `RESV_STATUS`, `RESERVATION_STATUS`, `SHORT_RESV_STATUS` fields. Trimmed, case-insensitive `CXL`, `CANCELLED`, `CANCELED` identify cancellation. A cancellation marker conflicting with a different nonempty status is not sufficient to exclude. Missing/unknown statuses (including CA/CD), NoShow and reservation types are not interpreted as cancellation. The retired direct-MCP v3 replay keeps its old rules.
+Processors4.3.0 and later exclude cancelled reservations whether or not a room is assigned, before all business-field checks and room/date deduplication. OHIP uses `reservation_status`, retained from matching search/detail responses. XML uses the explicitly named `RESV_STATUS`, `RESERVATION_STATUS`, `SHORT_RESV_STATUS` fields. Trimmed, case-insensitive `CXL`, `CANCELLED`, `CANCELED` identify cancellation. A cancellation marker conflicting with a different nonempty status is not sufficient to exclude. Missing/unknown statuses (including CA/CD), NoShow and reservation types are not interpreted as cancellation. The retired direct-MCP v3 replay keeps its old rules.
-Every source row stays in the structured audit. `outcome_counts.excluded_cancelled` is always present for4.3.0, including zero; it is separate from `removed_by_rate_code`. Excluded cancellations have no pricing, amount or channel facts and never participate in daily or monthly totals. Original source files and captured HTTP responses remain read-only.
+Every source row stays in the structured audit. `outcome_counts.excluded_cancelled` is always present for4.3.0 and later, including zero; it is separate from `removed_by_rate_code`. Excluded cancellations have no pricing, amount or channel facts and never participate in daily or monthly totals. Original source files and captured HTTP responses remain read-only.
+
+## PM room-type scope
+
+Processor4.4.0 excludes records only when the explicit `ROOM_CATEGORY_LABEL`, after trimming and case normalization, equals `PM`. This rule is shared by XML and direct OHIP data. Cancellation takes precedence: a cancelled PM record is counted only as `excluded_cancelled`. PM exclusion precedes rate-code/required-field checks, room/date deduplication and pricing, so a PM row cannot consume a real room's duplicate key or produce a manual-price issue. A room number, zero price, guest count, missing room type, or a similar label such as `PM1` never proves PM.
+
+The source record and source sequence remain audited with `excluded_pm` and exactly `ROOM_TYPE_PM_EXCLUDED`; the original room-type value retains its case after normal text trimming. `outcome_counts.excluded_pm` is always present for4.4.0, including zero, and remains separate from cancelled/rate-code exclusions. PM rows have no price, total, KB, channel or pricing-method facts and do not enter daily/monthly totals or manual pricing. The frozen direct-MCP v3 path retains its original scope and never emits the new outcome.
diff --git a/arr-opera-daily-ingest/references/data-structured-result.schema.json b/arr-opera-daily-ingest/references/data-structured-result.schema.json
index d92e503..b3dfbee 100644
--- a/arr-opera-daily-ingest/references/data-structured-result.schema.json
+++ b/arr-opera-daily-ingest/references/data-structured-result.schema.json
@@ -17,7 +17,7 @@
"activation_eligible": { "type": "boolean" },
"ingestion_mode": { "type": "string", "const": "ohip_data" },
"business_date": { "type": ["string", "null"], "format": "date" },
- "processor_version": { "type": "string", "const": "4.3.0" },
+ "processor_version": { "type": "string", "const": "4.4.0" },
"rule_set_sha256": { "$ref": "#/$defs/sha256" },
"source_rows": { "type": "integer", "minimum": 0 },
"removed_by_rate_code": { "type": "integer", "minimum": 0 },
@@ -32,11 +32,12 @@
"manually_priced_rows": { "type": "integer", "minimum": 0 },
"outcome_counts": {
"type": "object", "additionalProperties": false,
- "required": ["candidate", "duplicate", "excluded_cancelled", "excluded_rate_code", "price_unmatched", "retained", "validation_failed"],
+ "required": ["candidate", "duplicate", "excluded_cancelled", "excluded_pm", "excluded_rate_code", "price_unmatched", "retained", "validation_failed"],
"properties": {
"candidate": { "type": "integer", "minimum": 0 },
"duplicate": { "type": "integer", "minimum": 0 },
"excluded_cancelled": { "type": "integer", "minimum": 0 },
+ "excluded_pm": { "type": "integer", "minimum": 0 },
"excluded_rate_code": { "type": "integer", "minimum": 0 },
"price_unmatched": { "type": "integer", "minimum": 0 },
"retained": { "type": "integer", "minimum": 0 },
@@ -148,7 +149,7 @@
"properties": {
"source_sequence": { "type": "integer", "minimum": 1 }, "source_location": { "type": "string", "minLength": 1 },
"source_worksheet": { "type": "null" }, "source_row_no": { "type": "null" },
- "outcome": { "type": "string", "enum": ["retained", "candidate", "excluded_rate_code", "duplicate", "validation_failed", "price_unmatched", "excluded_cancelled"] },
+ "outcome": { "type": "string", "enum": ["retained", "candidate", "excluded_rate_code", "duplicate", "validation_failed", "price_unmatched", "excluded_cancelled", "excluded_pm"] },
"decision_codes": { "type": "array", "uniqueItems": true, "items": { "type": "string", "minLength": 1 } },
"duplicate_of_source_sequence": { "type": ["integer", "null"], "minimum": 1 },
"adults": { "type": ["integer", "null"] }, "children": { "type": ["integer", "null"] }, "block_code": { "type": "string" },
@@ -175,6 +176,15 @@
"kb_amount": { "type": "null" }, "channel_key": { "type": "null" }, "pricing_method": { "type": "null" }
} }
},
+ {
+ "if": { "properties": { "outcome": { "const": "excluded_pm" } }, "required": ["outcome"] },
+ "then": { "properties": {
+ "room_category_label": { "pattern": "^[Pp][Mm]$" },
+ "decision_codes": { "const": ["ROOM_TYPE_PM_EXCLUDED"] },
+ "real_price": { "type": "null" }, "total_price": { "type": "null" },
+ "kb_amount": { "type": "null" }, "channel_key": { "type": "null" }, "pricing_method": { "type": "null" }
+ } }
+ },
{
"if": { "properties": { "outcome": { "const": "duplicate" } }, "required": ["outcome"] },
"then": { "properties": { "duplicate_of_source_sequence": { "type": "integer", "minimum": 1 } } },
diff --git a/arr-opera-daily-ingest/references/field-contracts.md b/arr-opera-daily-ingest/references/field-contracts.md
index 09f4de0..c9ce460 100644
--- a/arr-opera-daily-ingest/references/field-contracts.md
+++ b/arr-opera-daily-ingest/references/field-contracts.md
@@ -107,7 +107,7 @@ When trimmed `RES_COMMENT` is empty, `group_code_key` is null and `booking_sourc
The complete record field list and conditional nullability rules are authoritative in [structured-result.schema.json](structured-result.schema.json).
-## Direct OHIP data input (processor 4.3.0)
+## Direct OHIP data input (processor 4.4.0)
`--data-json` accepts frozen `arr-ohip-data/v1` field observations with complete collection, exact hotel/date context, unique reservation IDs and continuous source order. Classification, whitelist, room/arrival deduplication, pricing, zero-price exceptions, nights, channel assignment and formulas are shared with XML. An unresolved field on a retained candidate fails the batch; it is not treated as a missing price. Non-whitelist rows are excluded before business-field validation, as in XML.
@@ -115,3 +115,5 @@ The reservation comment uses the first nonempty note in the supplied order after
An OHIP row may carry a nonempty string `reservation_status` outside the15 business `fields`. New acquisition retains it from independently matched search/detail states. For pre-upgrade pending captures, an internal saved-response evidence receipt can append statuses to a derived source without changing original bytes or staff decisions. Cancelled rows are excluded before candidate validation. A missing status in an old source is not silently interpreted as cancelled.
+
+The shared scope now excludes an explicit `ROOM_CATEGORY_LABEL` equal to `PM` (trimmed, case-insensitive), after cancellation and before business-field validation. PM rows retain source/audit fields with `excluded_pm` and null pricing/channel facts. Neither zero price nor room number identifies PM.
diff --git a/arr-opera-daily-ingest/references/structured-output.md b/arr-opera-daily-ingest/references/structured-output.md
index d261c02..93195be 100644
--- a/arr-opera-daily-ingest/references/structured-output.md
+++ b/arr-opera-daily-ingest/references/structured-output.md
@@ -31,7 +31,7 @@ The payload includes:
- `removed_by_rate_code`
- `removed_as_duplicates`
- `output_rows`, `candidate_rows`, `review_required_rows`, and `review_issue_count`
-- seven-outcome reconciliation (processor4.3.0)
+- eight-outcome reconciliation (processor4.4.0)
- grouped, privacy-minimized `review_issues` with fixed-table candidate-price comparisons
- `review_case_id`, `manual_override_sha256`, and `manually_priced_rows` on final manual replay only
- channel counts
@@ -49,6 +49,8 @@ Allowed outcomes:
- `retained`
- `candidate`
+- `excluded_cancelled`
+- `excluded_pm`
- `excluded_rate_code`
- `duplicate`
- `validation_failed`
@@ -93,3 +95,5 @@ Processor 4.2.0 adds explicit direct-data schema 5.0 (`ohip_data`, `source_data`
Processor4.3.0 preserves result schema4.0/5.0 and adds `excluded_cancelled` to `outcome_counts` and record outcomes. It requires `RESERVATION_CANCELLED_EXCLUDED` with no price/channel facts; all source rows still reconcile. Old processor results retain their original outcome contract; the new outcome cannot be declared under an old processor identity. Finance requires migration020.
+
+Processor4.4.0 keeps schema4.0/5.0 and adds `excluded_pm` with `ROOM_TYPE_PM_EXCLUDED` for an explicit, trimmed, case-insensitive `ROOM_CATEGORY_LABEL = PM`. Cancellation is checked first; PM exclusion is checked before rate-code/required-field validation, deduplication and pricing. The count is required even when zero and is reconciled separately from cancelled and rate-code exclusions. PM audit rows preserve source order and room type but have null `real_price`, `total_price`, `kb_amount`, `channel_key` and `pricing_method`; they never enter daily/monthly facts or manual-price issues. Finance requires migration021. Old processor identities and the frozen v3 replay cannot declare this outcome.
diff --git a/arr-opera-daily-ingest/references/structured-result.schema.json b/arr-opera-daily-ingest/references/structured-result.schema.json
index e531d23..3db8918 100644
--- a/arr-opera-daily-ingest/references/structured-result.schema.json
+++ b/arr-opera-daily-ingest/references/structured-result.schema.json
@@ -17,7 +17,7 @@
"activation_eligible": { "type": "boolean" },
"ingestion_mode": { "type": "string", "const": "opera_xml" },
"business_date": { "type": ["string", "null"], "format": "date" },
- "processor_version": { "type": "string", "const": "4.3.0" },
+ "processor_version": { "type": "string", "const": "4.4.0" },
"rule_set_sha256": { "$ref": "#/$defs/sha256" },
"source_rows": { "type": "integer", "minimum": 0 },
"removed_by_rate_code": { "type": "integer", "minimum": 0 },
@@ -32,11 +32,12 @@
"manually_priced_rows": { "type": "integer", "minimum": 0 },
"outcome_counts": {
"type": "object", "additionalProperties": false,
- "required": ["candidate", "duplicate", "excluded_cancelled", "excluded_rate_code", "price_unmatched", "retained", "validation_failed"],
+ "required": ["candidate", "duplicate", "excluded_cancelled", "excluded_pm", "excluded_rate_code", "price_unmatched", "retained", "validation_failed"],
"properties": {
"candidate": { "type": "integer", "minimum": 0 },
"duplicate": { "type": "integer", "minimum": 0 },
"excluded_cancelled": { "type": "integer", "minimum": 0 },
+ "excluded_pm": { "type": "integer", "minimum": 0 },
"excluded_rate_code": { "type": "integer", "minimum": 0 },
"price_unmatched": { "type": "integer", "minimum": 0 },
"retained": { "type": "integer", "minimum": 0 },
@@ -148,7 +149,7 @@
"properties": {
"source_sequence": { "type": "integer", "minimum": 1 }, "source_location": { "type": "string", "minLength": 1 },
"source_worksheet": { "type": "null" }, "source_row_no": { "type": "null" },
- "outcome": { "type": "string", "enum": ["retained", "candidate", "excluded_rate_code", "duplicate", "validation_failed", "price_unmatched", "excluded_cancelled"] },
+ "outcome": { "type": "string", "enum": ["retained", "candidate", "excluded_rate_code", "duplicate", "validation_failed", "price_unmatched", "excluded_cancelled", "excluded_pm"] },
"decision_codes": { "type": "array", "uniqueItems": true, "items": { "type": "string", "minLength": 1 } },
"duplicate_of_source_sequence": { "type": ["integer", "null"], "minimum": 1 },
"adults": { "type": ["integer", "null"] }, "children": { "type": ["integer", "null"] }, "block_code": { "type": "string" },
@@ -175,6 +176,15 @@
"kb_amount": { "type": "null" }, "channel_key": { "type": "null" }, "pricing_method": { "type": "null" }
} }
},
+ {
+ "if": { "properties": { "outcome": { "const": "excluded_pm" } }, "required": ["outcome"] },
+ "then": { "properties": {
+ "room_category_label": { "pattern": "^[Pp][Mm]$" },
+ "decision_codes": { "const": ["ROOM_TYPE_PM_EXCLUDED"] },
+ "real_price": { "type": "null" }, "total_price": { "type": "null" },
+ "kb_amount": { "type": "null" }, "channel_key": { "type": "null" }, "pricing_method": { "type": "null" }
+ } }
+ },
{
"if": { "properties": { "outcome": { "const": "duplicate" } }, "required": ["outcome"] },
"then": { "properties": { "duplicate_of_source_sequence": { "type": "integer", "minimum": 1 } } },
diff --git a/arr-opera-daily-ingest/scripts/process_daily.py b/arr-opera-daily-ingest/scripts/process_daily.py
index cbf4026..ba964bc 100755
--- a/arr-opera-daily-ingest/scripts/process_daily.py
+++ b/arr-opera-daily-ingest/scripts/process_daily.py
@@ -27,7 +27,7 @@ from openpyxl.utils import get_column_letter
RESULT_VERSION = "4.0"
-PROCESSOR_VERSION = "4.3.0"
+PROCESSOR_VERSION = "4.4.0"
DATA_RESULT_VERSION = "5.0"
STRUCTURED_RESULT_SCHEMA_VERSION = "4.0"
# Retired direct-MCP/callback compatibility. This is deliberately an internal,
@@ -66,6 +66,7 @@ RULE_SET_PATHS = (
FINAL_OUTCOMES = {
"retained",
"excluded_cancelled",
+ "excluded_pm",
"excluded_rate_code",
"duplicate",
"validation_failed",
@@ -1115,6 +1116,11 @@ def is_cancelled_record(record: Mapping[str, Any]) -> bool:
and text_or_blank(record.get("_RESERVATION_STATUS")).upper() in CANCELLED_RESERVATION_STATUSES)
+def is_pm_record(record: Mapping[str, Any]) -> bool:
+ """Exclude only an explicitly supplied PM room type, never infer it from other fields."""
+ return text_or_blank(record.get("ROOM_CATEGORY_LABEL")).upper() == "PM"
+
+
def classify_normalized_records(all_records, business_date, *, apply_scope=True):
candidates: List[Dict[str, Any]] = []
errors: List[ErrorItem] = []
@@ -1124,6 +1130,10 @@ def classify_normalized_records(all_records, business_date, *, apply_scope=True)
record["_OUTCOME"] = "excluded_cancelled"
append_decision(record, "RESERVATION_CANCELLED_EXCLUDED")
continue
+ if apply_scope and is_pm_record(record):
+ record["_OUTCOME"] = "excluded_pm"
+ append_decision(record, "ROOM_TYPE_PM_EXCLUDED")
+ continue
normalized_rate = record["_NORMALIZED_RATE_CODE"]
if not normalized_rate:
error = xml_record_error(
@@ -1787,7 +1797,7 @@ def build_legacy_direct_structured_result(
for outcome in sorted(LEGACY_DIRECT_FINAL_OUTCOMES)
}
if (
- any(record.get("_OUTCOME") == "candidate" for record in records)
+ any(record.get("_OUTCOME") not in LEGACY_DIRECT_FINAL_OUTCOMES for record in records)
or any(record.get("_PRICING_METHOD") == "manual_review" for record in records)
or legacy_counts["validation_failed"]
or legacy_counts["price_unmatched"]
@@ -1869,13 +1879,13 @@ def build_legacy_direct_failed_structured_result(
outcome: sum(1 for record in records if record.get("_OUTCOME") == outcome)
for outcome in sorted(LEGACY_DIRECT_FINAL_OUTCOMES)
}
- if any(record.get("_OUTCOME") == "candidate" for record in records):
+ if any(record.get("_OUTCOME") not in LEGACY_DIRECT_FINAL_OUTCOMES for record in records):
raise ProcessingFailure(
[
ErrorItem(
"LEGACY_DIRECT_OUTPUT_INVALID",
"legacy_direct",
- "旧direct_mcp失败输出不得包含候选复核行",
+ "旧direct_mcp失败输出不得包含冻结v3契约之外的记录状态",
)
]
)
@@ -2077,7 +2087,7 @@ def validate_structured_completeness(payload: Dict[str, Any]) -> None:
failures.append("source_sequence必须从1开始连续且保持XML顺序")
counts = payload.get("outcome_counts", {})
if not isinstance(counts, dict) or set(counts) != FINAL_OUTCOMES:
- failures.append("outcome_counts必须覆盖七种固定outcome")
+ failures.append("outcome_counts必须覆盖八种固定outcome")
counts = {}
for outcome in FINAL_OUTCOMES:
actual = sum(1 for record in records if record.get("outcome") == outcome)
@@ -2104,6 +2114,7 @@ def validate_structured_completeness(payload: Dict[str, Any]) -> None:
+ payload.get("removed_as_duplicates", 0)
+ payload.get("output_rows", 0)
+ counts.get("excluded_cancelled", 0)
+ + counts.get("excluded_pm", 0)
):
failures.append("成功payload的源行分解不平衡")
if payload.get("review_required_rows") or payload.get("review_issue_count"):
@@ -2121,7 +2132,7 @@ def validate_structured_completeness(payload: Dict[str, Any]) -> None:
failures.append("待复核payload必须包含缺价行")
if payload.get("candidate_rows", 0) + payload.get("review_required_rows", 0) + counts.get(
"excluded_rate_code", 0
- ) + counts.get("duplicate", 0) + counts.get("excluded_cancelled", 0) != source_rows:
+ ) + counts.get("duplicate", 0) + counts.get("excluded_cancelled", 0) + counts.get("excluded_pm", 0) != source_rows:
failures.append("待复核payload的源行分解不平衡")
issues = payload.get("review_issues")
if not isinstance(issues, list) or payload.get("review_issue_count") != len(issues):
@@ -2218,6 +2229,14 @@ def validate_structured_completeness(payload: Dict[str, Any]) -> None:
if any(record.get(field) is not None for field in (
"real_price", "total_price", "kb_amount", "channel_key", "pricing_method")):
failures.append("取消排除记录不得参与定价或渠道归属")
+ if record.get("outcome") == "excluded_pm":
+ if text_or_blank(record.get("room_category_label")).upper() != "PM":
+ failures.append("PM排除记录必须有明确的PM房型")
+ if decisions != ["ROOM_TYPE_PM_EXCLUDED"]:
+ failures.append("PM排除记录必须保留明确的PM房型排除原因")
+ if any(record.get(field) is not None for field in (
+ "real_price", "total_price", "kb_amount", "channel_key", "pricing_method")):
+ failures.append("PM排除记录不得参与定价或渠道归属")
if record.get("normalized_rate_code") != (
text_or_blank(record.get("rate_code")).upper() or None
):
diff --git a/arr-opera-daily-ingest/scripts/validate_daily.py b/arr-opera-daily-ingest/scripts/validate_daily.py
index 9693d0c..5deaedd 100755
--- a/arr-opera-daily-ingest/scripts/validate_daily.py
+++ b/arr-opera-daily-ingest/scripts/validate_daily.py
@@ -535,23 +535,6 @@ def validate_structured_result_contract(
f"structured-result.json 的 {field} 应为 {expected!r}",
)
)
- expected_outcome_counts = {
- "duplicate": removed_duplicates,
- "excluded_cancelled": source_rows - removed_by_rate - removed_duplicates - len(expected_records),
- "excluded_rate_code": removed_by_rate,
- "candidate": 0,
- "price_unmatched": 0,
- "retained": len(expected_records),
- "validation_failed": 0,
- }
- if payload.get("outcome_counts") != expected_outcome_counts:
- errors.append(
- validation_error(
- "OUTPUT_STRUCTURED_OUTCOME_MISMATCH",
- f"structured outcome计数应为 {expected_outcome_counts}",
- )
- )
-
artifacts = payload.get("artifacts")
if isinstance(artifacts, dict):
contract = core.source_contract(xml_path)
@@ -621,6 +604,23 @@ def validate_structured_result_contract(
)
)
return
+ expected_outcome_counts = {
+ "duplicate": removed_duplicates,
+ "excluded_cancelled": sum(record.get("_OUTCOME") == "excluded_cancelled" for record in all_records),
+ "excluded_pm": sum(record.get("_OUTCOME") == "excluded_pm" for record in all_records),
+ "excluded_rate_code": removed_by_rate,
+ "candidate": 0,
+ "price_unmatched": 0,
+ "retained": len(expected_records),
+ "validation_failed": 0,
+ }
+ if payload.get("outcome_counts") != expected_outcome_counts:
+ errors.append(
+ validation_error(
+ "OUTPUT_STRUCTURED_OUTCOME_MISMATCH",
+ f"structured outcome计数应为 {expected_outcome_counts}",
+ )
+ )
price_map = core.load_price_map(Path(core.PRICE_REFERENCE).resolve())
pricing_errors = core.apply_prices_classified(
retained,
@@ -766,6 +766,7 @@ def validate_review_structured_result_contract(
expected_outcome_counts = {
"candidate": candidate_rows,
"excluded_cancelled": sum(record.get("_OUTCOME") == "excluded_cancelled" for record in all_records),
+ "excluded_pm": sum(record.get("_OUTCOME") == "excluded_pm" for record in all_records),
"duplicate": removed_duplicates,
"excluded_rate_code": removed_by_rate,
"price_unmatched": review_rows,
@@ -982,6 +983,7 @@ def validate_legacy_direct_success_contracts(
replay_outcomes = dict(outcomes)
replay_outcomes["candidate"] = 0
replay_outcomes["excluded_cancelled"] = 0
+ replay_outcomes["excluded_pm"] = 0
replay_structured["outcome_counts"] = replay_outcomes
replay_artifacts = dict(artifacts)
replay_artifacts["manual_override_json"] = None
diff --git a/arr_ingestion/direct_postgres.py b/arr_ingestion/direct_postgres.py
index 3e39091..5adee94 100644
--- a/arr_ingestion/direct_postgres.py
+++ b/arr_ingestion/direct_postgres.py
@@ -757,12 +757,13 @@ class PostgresDirectIngestionRepository(PostgresIngestionRepository):
validation_failed_rows,
price_unmatched_rows,
excluded_cancelled_rows,
+ excluded_pm_rows,
validated_at,
result_delivery_mode
)
VALUES (
%s, %s, %s, %s, 'validated', %s, %s, %s, %s,
- %s, %s, %s, %s, %s, %s, %s, now(), 'direct_mcp'
+ %s, %s, %s, %s, %s, %s, %s, %s, now(), 'direct_mcp'
)
RETURNING id
""",
diff --git a/arr_ingestion/postgres.py b/arr_ingestion/postgres.py
index d72bd19..7eda27f 100644
--- a/arr_ingestion/postgres.py
+++ b/arr_ingestion/postgres.py
@@ -158,16 +158,20 @@ def _is_transient_database_error(error: Exception) -> bool:
return isinstance(sqlstate, str) and sqlstate in TRANSIENT_SQLSTATES
-def _outcome_counts(payload: Mapping[str, Any]) -> Tuple[int, int, int, int, int, int]:
+def _outcome_counts(payload: Mapping[str, Any]) -> Tuple[int, int, int, int, int, int, int]:
values = payload.get("outcome_counts")
if not isinstance(values, Mapping):
raise IngestionError(
"RESULT_CONTRACT_INVALID",
"structured outcome counts are invalid",
)
+ if payload.get("processor_version") == "4.4.0" and not {
+ "excluded_cancelled", "excluded_pm"
+ }.issubset(values):
+ raise IngestionError("RESULT_CONTRACT_INVALID", "structured outcome counts are invalid")
try:
counts = tuple(
- int(values.get(name, 0) if name == "excluded_cancelled" else values[name])
+ int(values.get(name, 0) if name in {"excluded_cancelled", "excluded_pm"} else values[name])
for name in (
"retained",
"excluded_rate_code",
@@ -175,6 +179,7 @@ def _outcome_counts(payload: Mapping[str, Any]) -> Tuple[int, int, int, int, int
"validation_failed",
"price_unmatched",
"excluded_cancelled",
+ "excluded_pm",
)
)
except (KeyError, TypeError, ValueError):
@@ -339,6 +344,7 @@ class PostgresIngestionRepository(IngestionRepository):
raise IngestionError("DATABASE_MIGRATION_MISSING", "direct data entry requires migration 019")
self._run_transaction(check, "direct data schema could not be checked")
self.assert_cancelled_scope_schema()
+ self.assert_pm_scope_schema()
def assert_cancelled_scope_schema(self) -> None:
"""Require the count and outcome constraints before either input is enabled."""
@@ -365,6 +371,36 @@ class PostgresIngestionRepository(IngestionRepository):
"ARR processing requires migration 020")
self._run_transaction(check, "ARR cancellation schema could not be checked")
+ def assert_pm_scope_schema(self) -> None:
+ """Require PM auditing before enabling the processor 4.4 entrypoints."""
+ def check(cursor):
+ cursor.execute("""SELECT EXISTS (
+ SELECT 1 FROM information_schema.columns
+ WHERE table_schema = 'finance' AND table_name = 'daily_versions'
+ AND column_name = 'excluded_pm_rows'
+ AND data_type = 'integer' AND is_nullable = 'NO'
+ ), (
+ SELECT pg_get_constraintdef(oid) FROM pg_constraint
+ WHERE conrelid = 'finance.daily_versions'::regclass
+ AND conname = 'daily_versions_counts_reconcile'
+ ), (
+ SELECT pg_get_constraintdef(oid) FROM pg_constraint
+ WHERE conrelid = 'finance.daily_records'::regclass
+ AND conname = 'daily_records_outcome_check'
+ ), (
+ SELECT pg_get_constraintdef(oid) FROM pg_constraint
+ WHERE conrelid = 'finance.daily_records'::regclass
+ AND conname = 'daily_records_pm_exclusion_shape'
+ )""")
+ row = cursor.fetchone()
+ if (not row or row[0] is not True
+ or "excluded_pm_rows" not in str(row[1] or "")
+ or "excluded_pm" not in str(row[2] or "")
+ or "ROOM_TYPE_PM_EXCLUDED" not in str(row[3] or "")):
+ raise IngestionError("DATABASE_MIGRATION_MISSING",
+ "ARR processing requires migration 021")
+ self._run_transaction(check, "ARR PM schema could not be checked")
+
def register_job(self, registration: JobRegistration) -> None:
if (
registration.source.role not in {"source_xml", "source_data"}
@@ -1837,13 +1873,14 @@ class PostgresIngestionRepository(IngestionRepository):
validation_failed_rows,
price_unmatched_rows,
excluded_cancelled_rows,
+ excluded_pm_rows,
failure_code,
failure_message,
validated_at
)
VALUES (
NULL, NULL, %s, %s, %s, %s, %s, 'rejected',
- %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s,
+ %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s,
%s, 'deterministic processing failed', now()
)
RETURNING id
@@ -2076,11 +2113,12 @@ class PostgresIngestionRepository(IngestionRepository):
validation_failed_rows,
price_unmatched_rows,
excluded_cancelled_rows,
+ excluded_pm_rows,
validated_at
)
VALUES (
%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, 'validated',
- %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now()
+ %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now()
)
RETURNING id
""",
diff --git a/arr_ingestion/validation.py b/arr_ingestion/validation.py
index 6f454e9..bcd06fe 100644
--- a/arr_ingestion/validation.py
+++ b/arr_ingestion/validation.py
@@ -483,16 +483,20 @@ V4_OUTCOME_FIELDS = {
"retained",
"validation_failed",
}
-CANCELLED_OUTCOME_PROCESSOR_VERSIONS = {"4.3.0"}
+CANCELLED_OUTCOME_PROCESSOR_VERSIONS = {"4.3.0", "4.4.0"}
+PM_OUTCOME_PROCESSOR_VERSIONS = {"4.4.0"}
def _v4_outcome_fields(processor_version: str) -> set[str]:
# Schema 4/5 remain readable for older immutable deliveries. The new
- # outcome is part of the explicitly approved 4.3 rule identity only; the
+ # outcomes are part of the explicitly approved rule identities only; the
# delivery validator also binds that identity to the approved rule hash.
+ outcomes = V4_OUTCOME_FIELDS
if processor_version in CANCELLED_OUTCOME_PROCESSOR_VERSIONS:
- return V4_OUTCOME_FIELDS | {"excluded_cancelled"}
- return V4_OUTCOME_FIELDS
+ outcomes = outcomes | {"excluded_cancelled"}
+ if processor_version in PM_OUTCOME_PROCESSOR_VERSIONS:
+ outcomes = outcomes | {"excluded_pm"}
+ return outcomes
V4_RESULT_METRIC_FIELDS = {
@@ -563,6 +567,15 @@ def _validate_v4_records(payload: Mapping[str, Any]) -> Tuple[Dict[str, int], in
or len(decisions) != len(set(decisions))
):
raise IngestionError("RESULT_CONTRACT_INVALID", "structured decisions are invalid")
+ if outcome == "excluded_pm" and (
+ not isinstance(values.get("room_category_label"), str)
+ or values["room_category_label"].strip().upper() != "PM"
+ or decisions != ["ROOM_TYPE_PM_EXCLUDED"]
+ or any(values.get(field) is not None for field in (
+ "real_price", "total_price", "kb_amount", "channel_key", "pricing_method"
+ ))
+ ):
+ raise IngestionError("RESULT_CONTRACT_INVALID", "structured PM exclusion is invalid")
group_code = values.get("group_code_key")
booking_status = values.get("booking_source_match_status")
if group_code is None:
@@ -709,6 +722,7 @@ def _validate_structured_payload_v4(
or counts["source_rows"]
!= counts["removed_by_rate_code"] + counts["removed_as_duplicates"] + counts["output_rows"]
+ outcomes.get("excluded_cancelled", 0)
+ + outcomes.get("excluded_pm", 0)
or channel_rows != counts["output_rows"]
or review_keys
):
@@ -737,6 +751,7 @@ def _validate_structured_payload_v4(
+ counts["candidate_rows"]
+ counts["review_required_rows"]
+ outcomes.get("excluded_cancelled", 0)
+ + outcomes.get("excluded_pm", 0)
or channels
or not review_keys
or manual_ref is not None
diff --git a/arr_web/arr_data_review.py b/arr_web/arr_data_review.py
index 35016ad..0aa52f5 100644
--- a/arr_web/arr_data_review.py
+++ b/arr_web/arr_data_review.py
@@ -277,16 +277,19 @@ class DataFieldReviews:
data["status"] = "collected_with_gaps" if unresolved else "collected"
return data
- def _issues(self, directory, data):
+ def _normalized(self, directory, data):
# Use the exact frozen processor's normalizer and field validation; the
# acquisition layer never maintains a second whitelist or numeric rule.
with tempfile.NamedTemporaryFile(suffix=".json", dir=directory) as stream:
stream.write(json_bytes(data))
stream.flush()
- business_date, records = self.rules.read_data_source(Path(stream.name))
+ return self.rules.read_data_source(Path(stream.name))
+
+ def _issues(self, directory, data):
+ business_date, records = self._normalized(directory, data)
issues = {}
for row, record in zip(data["records"], records):
- if self.rules.is_cancelled_record(record):
+ if self.rules.is_cancelled_record(record) or self.rules.is_pm_record(record):
continue
rate = record["_NORMALIZED_RATE_CODE"]
fields = {}
@@ -312,13 +315,24 @@ class DataFieldReviews:
value = observation.get("value")
return value if observation.get("state") == "available" and isinstance(value, (str, int)) else ""
+ def _excluded_sequences(self, directory, data):
+ cancelled, pm = set(), set()
+ _date, records = self._normalized(directory, data)
+ for row, record in zip(data["records"], records):
+ sequence = row["source_sequence"]
+ if self.rules.is_cancelled_record(record):
+ cancelled.add(sequence)
+ elif self.rules.is_pm_record(record):
+ pm.add(sequence)
+ return cancelled, pm
+
def _public(self, directory, meta, original, state):
data = self._derive(directory, original, state)
issues = self._issues(directory, data)
- cancelled_sequences = {row["source_sequence"] for row in data["records"]
- if self.rules.is_cancelled_record({"_RESERVATION_STATUS": row.get("reservation_status", "")})}
+ cancelled_sequences, pm_sequences = self._excluded_sequences(directory, data)
keys = set(issues) | {key for key in state["decisions"]
- if int(key.split(":")[0]) not in cancelled_sequences}
+ if int(key.split(":")[0]) not in cancelled_sequences
+ and (int(key.split(":")[0]) not in pm_sequences or key.endswith(":ROOM_CATEGORY_LABEL"))}
items = []
for item_id in sorted(keys, key=lambda key: (int(key.split(":")[0]), key.split(":")[1])):
sequence, field = item_id.split(":")
@@ -335,6 +349,7 @@ class DataFieldReviews:
"revision": state["revision"], "status": state["status"], "items": items,
"pending_count": len(issues), "total_count": len(items),
"excluded_cancelled_count": len(cancelled_sequences),
+ "excluded_pm_count": len(pm_sequences),
"can_finalize": not issues and state["status"] == "editing"}
def prepare(self, request_id, payload, manifest_sha256, report_date):
@@ -523,9 +538,11 @@ class DataFieldReviews:
actor = self._actor(actor)
data = self._derive(directory, original, state)
sequence = int(item_id.split(":")[0]) if isinstance(item_id, str) and re.fullmatch(r"[1-9][0-9]*:[A-Z_]+", item_id) else 0
- if 1 <= sequence <= len(data["records"]) and self.rules.is_cancelled_record(
- {"_RESERVATION_STATUS": data["records"][sequence - 1].get("reservation_status", "")}):
+ cancelled_sequences, pm_sequences = self._excluded_sequences(directory, data)
+ if sequence in cancelled_sequences:
raise _error(message="已取消预订不参与日报,无需完善字段")
+ if sequence in pm_sequences and not (item_id.endswith(":ROOM_CATEGORY_LABEL") and item_id in state["decisions"]):
+ raise _error(message="PM 房型不参与日报,无需完善字段")
if item_id not in self._issues(directory, data) and item_id not in state["decisions"]:
raise _error(message="只能完善当前任务列出的异常字段")
sequence, field = item_id.split(":")
diff --git a/arr_web/local_ohip.py b/arr_web/local_ohip.py
index 21f5633..b7276bb 100644
--- a/arr_web/local_ohip.py
+++ b/arr_web/local_ohip.py
@@ -215,7 +215,7 @@ def open_portal(root, port):
load_key(credential_file)
access = AccessState(root)
with job_lock(root):
- database = ReplayDatabase(root, schema_version=20)
+ database = ReplayDatabase(root, schema_version=21)
try:
database.start()
def factory(*, root, snapshot, **dependencies):
diff --git a/arr_web/local_replay_database.py b/arr_web/local_replay_database.py
index 788a6b6..68e262e 100644
--- a/arr_web/local_replay_database.py
+++ b/arr_web/local_replay_database.py
@@ -14,8 +14,8 @@ from integrations.ohip.capture_job import atomic_json, private_directory
class ReplayDatabase:
- def __init__(self, root: Path, *, schema_version: int = 20):
- if schema_version not in (18, 19, 20):
+ def __init__(self, root: Path, *, schema_version: int = 21):
+ if schema_version not in (18, 19, 20, 21):
raise ValueError("unsupported_local_schema")
self.schema_version = schema_version
self.root = root
diff --git a/arr_web/processing_runtime.py b/arr_web/processing_runtime.py
index e5daffc..b8f4928 100644
--- a/arr_web/processing_runtime.py
+++ b/arr_web/processing_runtime.py
@@ -80,6 +80,7 @@ def compose_programmatic_processing(
connect=connect,
)
repository.assert_cancelled_scope_schema()
+ repository.assert_pm_scope_schema()
policy = load_processor_policy(project_root)
ingestion = IngestionService(DeliveryValidator(object_store, policy), repository)
processor = LocalDailyProcessor(policy)
diff --git a/arr_web/static/app.js b/arr_web/static/app.js
index 3583ce4..8334a45 100644
--- a/arr_web/static/app.js
+++ b/arr_web/static/app.js
@@ -698,6 +698,9 @@
if (Number.isInteger(review?.excluded_cancelled_count) && review.excluded_cancelled_count > 0) {
$("#arr-data-review-progress").textContent += ` · ${I18N.t("data_review.excluded_cancelled", { count: review.excluded_cancelled_count })}`;
}
+ if (Number.isInteger(review?.excluded_pm_count) && review.excluded_pm_count > 0) {
+ $("#arr-data-review-progress").textContent += ` · ${I18N.t("data_review.excluded_pm", { count: review.excluded_pm_count })}`;
+ }
$("#arr-data-review-refresh").textContent = I18N.t("data_review.refresh");
$("#arr-data-review-refresh").disabled = busy;
const finalize = $("#arr-data-review-finalize");
diff --git a/arr_web/static/i18n.js b/arr_web/static/i18n.js
index cf85127..0c815c0 100644
--- a/arr_web/static/i18n.js
+++ b/arr_web/static/i18n.js
@@ -172,6 +172,7 @@
"data_review.loading": ["正在读取待完善字段…", "Loading fields for review…", "กำลังโหลดข้อมูลที่ต้องตรวจสอบ…"],
"data_review.no_items": ["没有待完善字段。", "No fields need review.", "ไม่มีข้อมูลที่ต้องตรวจสอบ"],
"data_review.excluded_cancelled": ["已排除 {count} 笔取消预订,原始数据保留", "{count} cancelled reservations excluded; original data retained", "ไม่รวมการจองที่ยกเลิก {count} รายการ โดยเก็บข้อมูลต้นฉบับไว้"],
+ "data_review.excluded_pm": ["已排除 {count} 笔 PM 预订,原始数据保留", "{count} PM reservations excluded; original data retained", "ไม่รวมการจองประเภทห้อง PM {count} รายการ โดยเก็บข้อมูลต้นฉบับไว้"],
"data_review.region": ["待完善报表字段", "Report fields to review", "ข้อมูลรายงานที่ต้องตรวจสอบ"],
"data_review.disconnected": ["字段读取中断,请刷新字段后继续。", "Field connection lost. Refresh the fields to continue.", "การเชื่อมต่อข้อมูลขัดข้อง โปรดรีเฟรชข้อมูลเพื่อดำเนินการต่อ"],
"data_review.saving": ["正在保存并确认字段…", "Saving and confirming field…", "กำลังบันทึกและยืนยันข้อมูล…"],
diff --git a/database/021_daily_pm_exclusion.down.sql b/database/021_daily_pm_exclusion.down.sql
new file mode 100644
index 0000000..f7d0d67
--- /dev/null
+++ b/database/021_daily_pm_exclusion.down.sql
@@ -0,0 +1,28 @@
+-- Never roll back after PM exclusion facts have been committed.
+BEGIN;
+DO $$
+BEGIN
+ IF current_database() <> 'booking_test' THEN
+ RAISE EXCEPTION 'ARR migration is allowed only in booking_test';
+ END IF;
+ IF EXISTS (SELECT 1 FROM finance.daily_versions WHERE excluded_pm_rows <> 0)
+ OR EXISTS (SELECT 1 FROM finance.daily_records WHERE outcome = 'excluded_pm') THEN
+ RAISE EXCEPTION 'refusing rollback: PM-exclusion facts exist';
+ END IF;
+END $$;
+
+ALTER TABLE finance.daily_records DROP CONSTRAINT daily_records_pm_exclusion_shape;
+ALTER TABLE finance.daily_records DROP CONSTRAINT daily_records_outcome_check;
+ALTER TABLE finance.daily_records
+ ADD CONSTRAINT daily_records_outcome_check CHECK (outcome IN (
+ 'retained', 'excluded_rate_code', 'duplicate', 'validation_failed',
+ 'price_unmatched', 'excluded_cancelled'
+ ));
+ALTER TABLE finance.daily_versions DROP CONSTRAINT daily_versions_counts_reconcile;
+ALTER TABLE finance.daily_versions
+ ADD CONSTRAINT daily_versions_counts_reconcile CHECK (
+ source_rows = retained_rows + excluded_rate_code_rows + duplicate_rows
+ + validation_failed_rows + price_unmatched_rows + excluded_cancelled_rows
+ );
+ALTER TABLE finance.daily_versions DROP COLUMN excluded_pm_rows;
+COMMIT;
diff --git a/database/021_daily_pm_exclusion.sql b/database/021_daily_pm_exclusion.sql
new file mode 100644
index 0000000..2279be1
--- /dev/null
+++ b/database/021_daily_pm_exclusion.sql
@@ -0,0 +1,52 @@
+-- Add PM exclusion auditing without rewriting any released daily facts.
+-- PostgreSQL 15+. Migrations 008 through 020 remain immutable.
+BEGIN;
+DO $$
+BEGIN
+ IF current_database() <> 'booking_test' THEN
+ RAISE EXCEPTION 'ARR migration is allowed only in booking_test';
+ END IF;
+ IF NOT EXISTS (
+ SELECT 1 FROM information_schema.columns
+ WHERE table_schema = 'finance' AND table_name = 'daily_versions'
+ AND column_name = 'excluded_cancelled_rows'
+ AND data_type = 'integer' AND is_nullable = 'NO'
+ ) OR NOT EXISTS (
+ SELECT 1 FROM pg_constraint
+ WHERE conrelid = 'finance.daily_records'::regclass
+ AND conname = 'daily_records_outcome_check'
+ AND pg_get_constraintdef(oid) LIKE '%excluded_cancelled%'
+ ) THEN
+ RAISE EXCEPTION 'ARR migration 020 must be applied first';
+ END IF;
+END $$;
+
+ALTER TABLE finance.daily_versions
+ ADD COLUMN excluded_pm_rows integer NOT NULL DEFAULT 0 CHECK (excluded_pm_rows >= 0);
+
+ALTER TABLE finance.daily_versions DROP CONSTRAINT daily_versions_counts_reconcile;
+ALTER TABLE finance.daily_versions
+ ADD CONSTRAINT daily_versions_counts_reconcile CHECK (
+ source_rows = retained_rows + excluded_rate_code_rows + duplicate_rows
+ + validation_failed_rows + price_unmatched_rows + excluded_cancelled_rows + excluded_pm_rows
+ );
+
+ALTER TABLE finance.daily_records DROP CONSTRAINT daily_records_outcome_check;
+ALTER TABLE finance.daily_records
+ ADD CONSTRAINT daily_records_outcome_check CHECK (outcome IN (
+ 'retained', 'excluded_rate_code', 'duplicate', 'validation_failed',
+ 'price_unmatched', 'excluded_cancelled', 'excluded_pm'
+ ));
+ALTER TABLE finance.daily_records
+ ADD CONSTRAINT daily_records_pm_exclusion_shape CHECK (
+ outcome <> 'excluded_pm' OR (
+ room_category_label IS NOT NULL AND upper(btrim(room_category_label)) = 'PM'
+ AND decision_codes = ARRAY['ROOM_TYPE_PM_EXCLUDED']::text[]
+ AND real_price IS NULL AND total_price IS NULL AND kb_amount IS NULL
+ AND channel_key IS NULL AND pricing_method IS NULL
+ )
+ );
+
+COMMENT ON COLUMN finance.daily_versions.excluded_pm_rows IS
+ 'Source rows explicitly excluded for PM room type under the pinned processor rules; older versions remain zero.';
+COMMIT;
diff --git a/tests/javascript/arr_data_review.cjs b/tests/javascript/arr_data_review.cjs
index 6a0674e..fc01822 100644
--- a/tests/javascript/arr_data_review.cjs
+++ b/tests/javascript/arr_data_review.cjs
@@ -19,6 +19,17 @@ test('cancelled source rows need no fields and the exclusion stays visible befor
assert.equal(h.calls.every(call=>call.method==='GET'),true);
});
+test('PM exclusions stay visible without requiring field completion',async()=>{
+ const h=harness(()=>response(200,{...review([]),excluded_cancelled_count:6,excluded_pm_count:19}));
+ await h.acceptARRDownloadTask(task,{sync:false});
+ const progress=h.element('#arr-data-review-progress').textContent;
+ assert.match(progress,/data_review\.excluded_cancelled/);
+ assert.match(progress,/data_review\.excluded_pm/);
+ assert.equal(h.element('#arr-data-review-body').innerHTML,'');
+ assert.equal(h.element('#arr-data-review-finalize').disabled,false);
+ assert.equal(h.calls.every(call=>call.method==='GET'),true);
+});
+
test('removed cancelled items clear obsolete drafts but retain edits still in the review',async()=>{
let current=review([item(),item({item_id:'2:DISP_ROOM_NO',source_sequence:2,field:'DISP_ROOM_NO',can_be_empty:false})]);
const h=harness(()=>response(200,current));
diff --git a/tests/local_postgres.py b/tests/local_postgres.py
index f066806..5d66697 100644
--- a/tests/local_postgres.py
+++ b/tests/local_postgres.py
@@ -76,8 +76,8 @@ class TemporaryPostgres:
connection.close()
raise
- def reset_database(self, *, schema_version=20):
- if schema_version not in (19, 20):
+ def reset_database(self, *, schema_version=21):
+ if schema_version not in (19, 20, 21):
raise ValueError("unsupported fixture schema version")
# Only this owned cluster has passed data_directory/socket checks above.
with self.connect(database="postgres", autocommit=True) as connection:
diff --git a/tests/test_arr_cancelled_ingestion.py b/tests/test_arr_cancelled_ingestion.py
index a2356ff..7ed8886 100644
--- a/tests/test_arr_cancelled_ingestion.py
+++ b/tests/test_arr_cancelled_ingestion.py
@@ -59,6 +59,7 @@ class CancelledIngestionContractTests(unittest.TestCase):
old_envelope = replace(envelope, processor_version="4.2.0")
old_payload = copy.deepcopy(payload)
old_payload["processor_version"] = "4.2.0"
+ old_payload["outcome_counts"].pop("excluded_pm", None)
with self.assertRaisesRegex(IngestionError, "outcome counts contract"):
_validate_structured_payload_v4(old_payload, old_envelope)
del old_payload["outcome_counts"]["excluded_cancelled"]
@@ -69,6 +70,7 @@ class CancelledIngestionContractTests(unittest.TestCase):
raw, store, envelope, payload = self.generated(Path(temporary))
old = copy.deepcopy(payload)
old["processor_version"] = "4.2.0"
+ old["outcome_counts"].pop("excluded_pm", None)
del old["outcome_counts"]["excluded_cancelled"]
with self.assertRaisesRegex(IngestionError, "structured outcome is invalid"):
_validate_structured_payload_v4(old, replace(envelope, processor_version="4.2.0"))
@@ -81,9 +83,9 @@ class CancelledIngestionContractTests(unittest.TestCase):
def test_finance_counts_preserve_old_payload_and_add_new_count(self):
counts = dict(retained=2, excluded_rate_code=1, duplicate=3,
validation_failed=4, price_unmatched=5)
- self.assertEqual(_outcome_counts({"outcome_counts": counts}), (2, 1, 3, 4, 5, 0))
+ self.assertEqual(_outcome_counts({"outcome_counts": counts}), (2, 1, 3, 4, 5, 0, 0))
counts["excluded_cancelled"] = 6
- self.assertEqual(_outcome_counts({"outcome_counts": counts}), (2, 1, 3, 4, 5, 6))
+ self.assertEqual(_outcome_counts({"outcome_counts": counts}), (2, 1, 3, 4, 5, 6, 0))
def test_cancelled_count_balances_pending_price_review(self):
with tempfile.TemporaryDirectory() as temporary:
@@ -152,17 +154,17 @@ class CancelledFinanceMigrationTests(unittest.TestCase):
connection.execute("RESET ROLE")
def test_empty_database_allows_guarded_rollback_and_reapply(self):
- self.database.reset_database()
+ self.database.reset_database(schema_version=20)
repository = PostgresIngestionRepository(DatabaseConfig("owned-fixture"), connect=self.database.connect)
- repository.assert_data_source_schema()
+ repository.assert_cancelled_scope_schema()
with self.database.connect(autocommit=True) as connection:
connection.execute((PROJECT / "database/020_daily_cancelled_exclusion.down.sql").read_text(), prepare=False)
with self.assertRaises(IngestionError) as missing:
- repository.assert_data_source_schema()
+ repository.assert_cancelled_scope_schema()
self.assertEqual(missing.exception.code, "DATABASE_MIGRATION_MISSING")
with self.database.connect(autocommit=True) as connection:
connection.execute((PROJECT / "database/020_daily_cancelled_exclusion.sql").read_text(), prepare=False)
- repository.assert_data_source_schema()
+ repository.assert_cancelled_scope_schema()
if __name__ == "__main__":
diff --git a/tests/test_arr_cancelled_scope.py b/tests/test_arr_cancelled_scope.py
index e2e9cc3..ce98d8d 100644
--- a/tests/test_arr_cancelled_scope.py
+++ b/tests/test_arr_cancelled_scope.py
@@ -157,7 +157,7 @@ class CancelledScopeTests(unittest.TestCase):
self.assertEqual((payload["source_rows"], payload["removed_by_rate_code"], payload["output_rows"]), (2, 0, 1))
self.assertEqual(payload["outcome_counts"]["excluded_cancelled"], 1)
self.assertEqual(set(payload["outcome_counts"]), core.FINAL_OUTCOMES)
- self.assertEqual(payload["processor_version"], "4.3.0")
+ self.assertEqual(payload["processor_version"], "4.4.0")
self.assertEqual([row["outcome"] for row in payload["records"]], ["excluded_cancelled", "retained"])
self.assertEqual(payload["records"][0]["decision_codes"], ["RESERVATION_CANCELLED_EXCLUDED"])
self.assertEqual(payload["records"][0]["confirmation_no"], "SYNTHETIC-CONF-1")
diff --git a/tests/test_arr_download_runtime.py b/tests/test_arr_download_runtime.py
index 04c5913..80ae040 100644
--- a/tests/test_arr_download_runtime.py
+++ b/tests/test_arr_download_runtime.py
@@ -41,6 +41,7 @@ class RuntimeTests(unittest.TestCase):
self.store = ManagedObjectStore(FilesystemObjectBackend(self.root / 'objects', create=True))
self.repository = InMemoryIngestionRepository()
self.repository.assert_cancelled_scope_schema = Mock()
+ self.repository.assert_pm_scope_schema = Mock()
self.processor = CountingProcessor(self.policy)
self.client = Mock()
with patch.object(processing_runtime, 'compose_object_store', return_value=
@@ -59,6 +60,21 @@ class RuntimeTests(unittest.TestCase):
def test_processing_composition_checks_finance_cancellation_schema(self):
self.repository.assert_cancelled_scope_schema.assert_called_once_with()
+ self.repository.assert_pm_scope_schema.assert_called_once_with()
+
+ def test_missing_pm_migration_closes_storage_and_disables_composition(self):
+ from arr_ingestion.contracts import IngestionError
+ self.repository.assert_pm_scope_schema.side_effect = IngestionError(
+ "DATABASE_MIGRATION_MISSING", "ARR processing requires migration 021")
+ client = Mock()
+ with patch.object(processing_runtime, 'compose_object_store', return_value=
+ processing_runtime.ObjectStoreRuntime(client, self.store)), \
+ patch.object(processing_runtime, 'PostgresIngestionRepository', return_value=self.repository), \
+ self.assertRaises(IngestionError) as error:
+ processing_runtime.compose_programmatic_processing(
+ project_root=Path(__file__).resolve().parents[1], connect=Mock())
+ self.assertEqual(error.exception.code, "DATABASE_MIGRATION_MISSING")
+ client.close.assert_called_once_with()
def test_missing_cancellation_migration_closes_storage_and_disables_composition(self):
from arr_ingestion.contracts import IngestionError
diff --git a/tests/test_arr_ingestion_postgres.py b/tests/test_arr_ingestion_postgres.py
index 04aa1af..f4fcdb8 100644
--- a/tests/test_arr_ingestion_postgres.py
+++ b/tests/test_arr_ingestion_postgres.py
@@ -533,7 +533,7 @@ class PostgresIngestionTests(unittest.TestCase):
and "INSERT INTO finance.daily_versions" in value
)
- self.assertEqual(version_insert.count("%s"), 21)
+ self.assertEqual(version_insert.count("%s"), 22)
class Migration008ContractTests(unittest.TestCase):
diff --git a/tests/test_arr_opera_daily_ingest.py b/tests/test_arr_opera_daily_ingest.py
index b9fd2cd..b6965ac 100644
--- a/tests/test_arr_opera_daily_ingest.py
+++ b/tests/test_arr_opera_daily_ingest.py
@@ -236,7 +236,7 @@ class ArrOperaDailyIngestTests(unittest.TestCase):
self.assertNotIn(forbidden, scripts)
self.assertEqual(core.RESULT_VERSION, "4.0")
self.assertEqual(core.STRUCTURED_RESULT_SCHEMA_VERSION, "4.0")
- self.assertEqual(core.PROCESSOR_VERSION, "4.3.0")
+ self.assertEqual(core.PROCESSOR_VERSION, "4.4.0")
self.assertEqual(len(core.DAILY_HEADERS), 19)
self.assertEqual(len(core.RATE_WHITELIST), 20)
@@ -302,7 +302,7 @@ class ArrOperaDailyIngestTests(unittest.TestCase):
)
self.assertEqual(payload["result_schema_version"], "4.0")
- self.assertEqual(payload["processor_version"], "4.3.0")
+ self.assertEqual(payload["processor_version"], "4.4.0")
self.assertEqual(payload["source_rows"], 5)
self.assertEqual(
payload["outcome_counts"],
@@ -310,6 +310,7 @@ class ArrOperaDailyIngestTests(unittest.TestCase):
"candidate": 0,
"duplicate": 1,
"excluded_cancelled": 0,
+ "excluded_pm": 0,
"excluded_rate_code": 1,
"price_unmatched": 0,
"retained": 3,
@@ -802,7 +803,7 @@ class ArrOperaDailyIngestTests(unittest.TestCase):
},
}
)
- self.assertEqual(delivery.processor_version, "4.3.0")
+ self.assertEqual(delivery.processor_version, "4.4.0")
self.assertEqual(delivery.result_schema_version, "4.0")
self.assertEqual(delivery.business_date, date(2026, 7, 27))
diff --git a/tests/test_arr_pm_ingestion.py b/tests/test_arr_pm_ingestion.py
new file mode 100644
index 0000000..23742aa
--- /dev/null
+++ b/tests/test_arr_pm_ingestion.py
@@ -0,0 +1,315 @@
+"""Approved PM exclusion contracts and additive Finance upgrade; all inputs synthetic."""
+from __future__ import annotations
+
+import copy
+from dataclasses import replace
+import json
+import os
+from pathlib import Path
+import tempfile
+import unittest
+
+from openpyxl import load_workbook
+
+from arr_ingestion.contracts import DeliveryEnvelope, IngestionError
+from arr_ingestion.postgres import DatabaseConfig, PostgresIngestionRepository, _outcome_counts
+from arr_ingestion.validation import DeliveryValidator, _validate_structured_payload_v4
+from arr_processing.policy import load_processor_policy
+from tests.local_postgres import TemporaryPostgres
+from tests import test_arr_direct_data as direct_helpers
+from tests.test_arr_ingestion_validation import build_delivery, policy
+from tests.test_arr_opera_daily_ingest import reservation, xml_document
+
+
+PROJECT = Path(__file__).resolve().parents[1]
+
+
+def pm_row(*, status=None):
+ value = reservation(1, rate_code="", rate_amount="0", company="").replace(
+ "SYNTHETIC ROOM TYPE",
+ " pm ",
+ ).replace("SYNTHETIC-ROOM-1", "")
+ if status:
+ value = value.replace("", f"{status}")
+ return value
+
+
+class PMIngestionContractTests(unittest.TestCase):
+ def generated(self, root, *rows):
+ raw, store, _ = build_delivery(xml_document(*rows), root)
+ envelope = DeliveryEnvelope.from_dict(json.loads(raw))
+ payload = json.loads(store.objects[envelope.artifacts["structured_result_json"].object_key])
+ return raw, store, envelope, payload
+
+ def test_pm_missing_required_fields_is_audited_before_validation_and_pricing(self):
+ with tempfile.TemporaryDirectory() as temporary:
+ raw, store, _, payload = self.generated(Path(temporary), pm_row(), reservation(2))
+ verified = DeliveryValidator(store, policy()).validate(raw)
+ self.assertEqual(verified.envelope.status, "success")
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ self.assertEqual(payload["output_rows"], 1)
+ row = payload["records"][0]
+ self.assertEqual(row["outcome"], "excluded_pm")
+ self.assertEqual(row["decision_codes"], ["ROOM_TYPE_PM_EXCLUDED"])
+ for field in ("real_price", "total_price", "kb_amount", "channel_key", "pricing_method"):
+ self.assertIsNone(row[field])
+
+ def test_pm_count_is_required_only_for_44_and_older_versions_stay_readable(self):
+ with tempfile.TemporaryDirectory() as temporary:
+ _, _, envelope, payload = self.generated(Path(temporary), reservation(1))
+ self.assertEqual(payload["processor_version"], "4.4.0")
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 0)
+ for field in ("excluded_pm", "excluded_cancelled"):
+ missing = copy.deepcopy(payload)
+ del missing["outcome_counts"][field]
+ with self.assertRaisesRegex(IngestionError, "outcome counts contract"):
+ _validate_structured_payload_v4(missing, envelope)
+ for version in ("4.3.0", "4.2.0"):
+ older = copy.deepcopy(payload)
+ older["processor_version"] = version
+ with self.assertRaisesRegex(IngestionError, "outcome counts contract"):
+ _validate_structured_payload_v4(older, replace(envelope, processor_version=version))
+ del older["outcome_counts"]["excluded_pm"]
+ if version == "4.2.0":
+ del older["outcome_counts"]["excluded_cancelled"]
+ _validate_structured_payload_v4(older, replace(envelope, processor_version=version))
+
+ def test_pm_record_cannot_be_relabelled_or_given_price_facts(self):
+ with tempfile.TemporaryDirectory() as temporary:
+ _, _, envelope, payload = self.generated(Path(temporary), pm_row(), reservation(2))
+ changes = {
+ "room_category_label": "PM-SUITE",
+ "decision_codes": ["RATE_CODE_EXCLUDED"],
+ "real_price": 0,
+ "total_price": 0,
+ "kb_amount": 0,
+ "channel_key": "Group",
+ "pricing_method": "zero_price_exception",
+ }
+ for field, value in changes.items():
+ changed = copy.deepcopy(payload)
+ changed["records"][0][field] = value
+ with self.subTest(field=field), self.assertRaisesRegex(IngestionError, "PM exclusion"):
+ _validate_structured_payload_v4(changed, envelope)
+ older = copy.deepcopy(payload)
+ older["processor_version"] = "4.3.0"
+ del older["outcome_counts"]["excluded_pm"]
+ with self.assertRaisesRegex(IngestionError, "structured outcome is invalid"):
+ _validate_structured_payload_v4(older, replace(envelope, processor_version="4.3.0"))
+ changed = copy.deepcopy(payload)
+ changed["outcome_counts"]["excluded_pm"] = 0
+ with self.assertRaisesRegex(IngestionError, "outcomes do not reconcile"):
+ _validate_structured_payload_v4(changed, envelope)
+
+ def test_pm_balances_price_review_and_failure_and_cancelled_takes_precedence(self):
+ cases = (
+ (reservation(2, rate_amount="1800"), "review_required"),
+ (reservation(2, departure="2026-07-26"), "failed"),
+ )
+ for row, status in cases:
+ with self.subTest(status=status), tempfile.TemporaryDirectory() as temporary:
+ raw, store, _, payload = self.generated(Path(temporary), pm_row(), row)
+ self.assertEqual(DeliveryValidator(store, policy()).validate(raw).envelope.status, status)
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ with tempfile.TemporaryDirectory() as temporary:
+ raw, store, _, payload = self.generated(Path(temporary), pm_row(status="CXL"), reservation(2))
+ DeliveryValidator(store, policy()).validate(raw)
+ self.assertEqual(payload["outcome_counts"]["excluded_cancelled"], 1)
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 0)
+
+ def test_finance_defaults_only_old_pm_counts_and_rejects_44_missing_keys(self):
+ counts = dict(retained=2, excluded_rate_code=1, duplicate=3,
+ validation_failed=4, price_unmatched=5)
+ self.assertEqual(_outcome_counts({"outcome_counts": counts}), (2, 1, 3, 4, 5, 0, 0))
+ counts["excluded_cancelled"] = 6
+ self.assertEqual(_outcome_counts({"processor_version": "4.3.0", "outcome_counts": counts}),
+ (2, 1, 3, 4, 5, 6, 0))
+ with self.assertRaises(IngestionError):
+ _outcome_counts({"processor_version": "4.4.0", "outcome_counts": counts})
+ counts["excluded_pm"] = 7
+ self.assertEqual(_outcome_counts({"processor_version": "4.4.0", "outcome_counts": counts}),
+ (2, 1, 3, 4, 5, 6, 7))
+ del counts["excluded_cancelled"]
+ with self.assertRaises(IngestionError):
+ _outcome_counts({"processor_version": "4.4.0", "outcome_counts": counts})
+
+
+@unittest.skipUnless(os.environ.get("ARR_TEST_LOCAL_POSTGRES") == "1", "owned PostgreSQL opt-in required")
+class PMFinanceMigrationTests(unittest.TestCase):
+ @classmethod
+ def setUpClass(cls):
+ cls.database = TemporaryPostgres().__enter__()
+ cls.addClassCleanup(cls.database.__exit__, None, None, None)
+
+ def test_additive_upgrade_preserves_cancelled_history_and_role_privileges(self):
+ import psycopg
+ self.database.reset_database(schema_version=20)
+ repository = PostgresIngestionRepository(DatabaseConfig("owned-fixture"), connect=self.database.connect)
+ repository.assert_cancelled_scope_schema()
+ with self.assertRaises(IngestionError) as missing:
+ repository.assert_pm_scope_schema()
+ self.assertEqual(missing.exception.code, "DATABASE_MIGRATION_MISSING")
+ with self.database.connect(autocommit=True) as connection:
+ artifact = connection.execute("""INSERT INTO ingestion.artifacts
+ (artifact_kind,storage_provider,bucket_alias,object_key,original_filename,sha256,byte_size)
+ VALUES ('opera_xml','local_fixture','fixture','old/source.xml','source.xml',%s,0) RETURNING id""",
+ ("a" * 64,)).fetchone()[0]
+ output_ids = {}
+ for kind in ("daily_xlsx", "result_json", "structured_result_json"):
+ output_ids[kind] = connection.execute("""INSERT INTO ingestion.artifacts
+ (artifact_kind,storage_provider,bucket_alias,object_key,original_filename,sha256,byte_size)
+ VALUES (%s,'local_fixture','fixture',%s,%s,%s,0) RETURNING id""",
+ (kind, "old/" + kind, kind, "d" * 64)).fetchone()[0]
+ run = connection.execute("""INSERT INTO ingestion.processing_runs
+ (run_key,pipeline_type,source_artifact_id,result_artifact_id,run_status,result_delivery_mode,business_date,
+ delivered_processor_version,delivered_rule_set_sha256,result_schema_version,
+ delivery_sha256,validated_at,finished_at)
+ VALUES ('old-cancelled','opera_daily',%s,%s,'accepted','artifact_callback','2026-07-27',
+ '4.3.0',%s,'4.0',%s,now(),now()) RETURNING id""",
+ (artifact, output_ids["result_json"], "b" * 64, "c" * 64)).fetchone()[0]
+ version = connection.execute("""INSERT INTO finance.daily_versions
+ (business_date,version_no,processing_run_id,source_artifact_id,daily_report_artifact_id,
+ result_json_artifact_id,structured_result_artifact_id,version_status,
+ processor_version,rule_set_sha256,result_schema_version,result_sha256,source_rows,
+ retained_rows,excluded_rate_code_rows,duplicate_rows,validation_failed_rows,
+ price_unmatched_rows,excluded_cancelled_rows,validated_at,result_delivery_mode)
+ VALUES ('2026-07-27',1,%s,%s,%s,%s,%s,'validated','4.3.0',%s,'4.0',%s,
+ 2,1,0,0,0,0,1,now(),'artifact_callback') RETURNING id""",
+ (run, artifact, output_ids["daily_xlsx"], output_ids["result_json"],
+ output_ids["structured_result_json"], "b" * 64, "c" * 64)).fetchone()[0]
+ connection.execute("""INSERT INTO finance.daily_records
+ (daily_version_id,source_sequence,source_location,outcome,decision_codes,
+ booking_source_match_status)
+ VALUES (%s,1,'reservation[1]','excluded_cancelled',
+ ARRAY['RESERVATION_CANCELLED_EXCLUDED'],'missing_group_code')""", (version,))
+ # Processor 4.3 included this PM row. A schema upgrade must preserve it.
+ connection.execute("""INSERT INTO finance.daily_records
+ (daily_version_id,source_sequence,source_location,outcome,decision_codes,block_code,
+ adults,children,company_name,company_key,confirmation_no,disp_room_no,
+ effective_rate_amount,full_name,no_of_rooms,products,rate_code,room_category_label,
+ arrival,departure,nights,real_price,total_price,channel_key,pricing_method,
+ booking_source_match_status)
+ VALUES (%s,2,'reservation[2]','retained',ARRAY['PRICE_REFERENCE_EXACT'],'',2,0,
+ 'SYNTHETIC COMPANY','SYNTHETIC COMPANY','SYNTHETIC CONF','SYNTHETIC ROOM',
+ 900,'SYNTHETIC GUEST',1,'SYNTHETIC','GRPA1','PM','2026-07-27','2026-07-28',
+ 1,900,900,'Group','price_reference_exact','missing_group_code')""", (version,))
+ before = connection.execute("SELECT to_jsonb(v) FROM finance.daily_versions v WHERE id=%s", (version,)).fetchone()[0]
+ records_before = connection.execute("SELECT to_jsonb(r) FROM finance.daily_records r ORDER BY source_sequence").fetchall()
+ guards = connection.execute("SELECT tgname,pg_get_triggerdef(oid) FROM pg_trigger WHERE NOT tgisinternal ORDER BY tgname").fetchall()
+ connection.execute("CREATE ROLE arr_pm_fixture")
+ connection.execute("CREATE ROLE arr_pm_column_fixture")
+ connection.execute("GRANT USAGE ON SCHEMA finance TO arr_pm_fixture,arr_pm_column_fixture")
+ connection.execute("GRANT SELECT,INSERT ON finance.daily_versions TO arr_pm_fixture")
+ connection.execute("GRANT INSERT (excluded_cancelled_rows) ON finance.daily_versions TO arr_pm_column_fixture")
+ connection.execute((PROJECT / "database/021_daily_pm_exclusion.sql").read_text(), prepare=False)
+ repository.assert_data_source_schema()
+ after = connection.execute("SELECT to_jsonb(v) FROM finance.daily_versions v WHERE id=%s", (version,)).fetchone()[0]
+ self.assertEqual(after.pop("excluded_pm_rows"), 0)
+ self.assertEqual(after, before)
+ self.assertEqual(connection.execute("SELECT to_jsonb(r) FROM finance.daily_records r ORDER BY source_sequence").fetchall(), records_before)
+ self.assertEqual(connection.execute("SELECT tgname,pg_get_triggerdef(oid) FROM pg_trigger WHERE NOT tgisinternal ORDER BY tgname").fetchall(), guards)
+ self.assertEqual(connection.execute("SELECT has_column_privilege('arr_pm_fixture','finance.daily_versions','excluded_pm_rows','INSERT'),has_table_privilege('arr_pm_fixture','finance.daily_versions','UPDATE'),has_column_privilege('arr_pm_column_fixture','finance.daily_versions','excluded_pm_rows','INSERT')").fetchone(), (True, False, False))
+ for statement in (
+ "UPDATE finance.daily_versions SET excluded_pm_rows=-1 WHERE id=%s",
+ "UPDATE finance.daily_versions SET excluded_pm_rows=1 WHERE id=%s",
+ ):
+ with self.assertRaises(psycopg.errors.CheckViolation):
+ connection.execute(statement, (version,))
+ connection.execute("UPDATE finance.daily_versions SET excluded_cancelled_rows=0,excluded_pm_rows=1 WHERE id=%s", (version,))
+ with self.assertRaisesRegex(psycopg.errors.RaiseException, "PM-exclusion facts exist"):
+ connection.execute((PROJECT / "database/021_daily_pm_exclusion.down.sql").read_text(), prepare=False)
+ connection.execute("ROLLBACK")
+ connection.execute("SET ROLE arr_pm_fixture")
+ with self.assertRaises(psycopg.errors.InsufficientPrivilege):
+ connection.execute("UPDATE finance.daily_versions SET excluded_pm_rows=0 WHERE id=%s", (version,))
+ connection.execute("RESET ROLE")
+
+ def test_empty_rollback_preserves_020_and_allows_reapply(self):
+ self.database.reset_database()
+ repository = PostgresIngestionRepository(DatabaseConfig("owned-fixture"), connect=self.database.connect)
+ repository.assert_data_source_schema()
+ with self.database.connect(autocommit=True) as connection:
+ connection.execute((PROJECT / "database/021_daily_pm_exclusion.down.sql").read_text(), prepare=False)
+ repository.assert_cancelled_scope_schema()
+ with self.assertRaises(IngestionError) as missing:
+ repository.assert_pm_scope_schema()
+ self.assertEqual(missing.exception.code, "DATABASE_MIGRATION_MISSING")
+ with self.database.connect(autocommit=True) as connection:
+ connection.execute((PROJECT / "database/021_daily_pm_exclusion.sql").read_text(), prepare=False)
+ repository.assert_data_source_schema()
+
+
+@unittest.skipUnless(os.environ.get("ARR_TEST_LOCAL_POSTGRES") == "1", "owned PostgreSQL opt-in required")
+class PMFinanceProcessingTests(unittest.TestCase):
+ @classmethod
+ def setUpClass(cls):
+ cls.database = TemporaryPostgres().__enter__()
+ cls.addClassCleanup(cls.database.__exit__, None, None, None)
+ cls.policy = load_processor_policy(PROJECT)
+
+ setUp = direct_helpers.DirectDataPostgresTests.setUp
+ configure = direct_helpers.DirectDataTests.configure
+ execute = direct_helpers.DirectDataTests.execute
+ sql = direct_helpers.DirectDataPostgresTests.sql
+ count = direct_helpers.DirectDataPostgresTests.count
+ monthly_worker = direct_helpers.DirectDataPostgresTests.monthly_worker
+
+ def make_first_pm(self):
+ row = self.transport.rows[0]
+ row["roomStay"]["roomType"] = "PM"
+ row["roomStay"]["currentRoomInfo"]["roomType"] = "PM"
+ for rate in row["roomStay"]["roomRates"]:
+ rate["roomType"] = "PM"
+
+ def test_direct_pm_is_audited_and_daily_monthly_exclude_it_without_refetch(self):
+ self.make_first_pm()
+ self.assertEqual(self.execute().status, "succeeded")
+ self.assertEqual(self.sql("SELECT result_schema_version,source_rows,retained_rows,excluded_pm_rows FROM finance.daily_versions"), [("5.0", 6, 5, 1)])
+ self.assertEqual(self.sql("SELECT outcome,decision_codes,real_price,channel_key FROM finance.daily_records WHERE source_sequence=1"), [("excluded_pm", ["ROOM_TYPE_PM_EXCLUDED"], None, None)])
+ self.assertEqual(self.count("finance.v_active_daily_facts"), 5)
+ self.assertEqual(self.monthly_worker().process_next().status, "published")
+ key, size = self.sql("SELECT a.object_key,a.byte_size FROM reporting.monthly_runs r JOIN ingestion.artifacts a ON a.id=r.workbook_artifact_id")[0]
+ output = self.files / "monthly.xlsx"
+ self.store.materialize(key, output, size)
+ workbook = load_workbook(output, data_only=False)
+ try:
+ self.assertEqual(sum(cell.data_type == "f" for sheet in workbook for row in sheet for cell in row), 5)
+ finally:
+ workbook.close()
+ calls = len(self.transport.calls)
+ self.assertEqual(self.execute().status, "succeeded")
+ self.assertEqual(len(self.transport.calls), calls)
+ self.assertEqual(self.count("finance.daily_versions"), 1)
+
+ def test_xml_pm_is_audited_without_required_field_or_price_review(self):
+ outcome = self.coordinator.submit("source.xml", xml_document(pm_row(), reservation(2)).encode())
+ self.assertEqual(outcome["status"], "succeeded")
+ self.assertEqual(self.sql("SELECT result_schema_version,source_rows,retained_rows,excluded_pm_rows FROM finance.daily_versions"), [("4.0", 2, 1, 1)])
+ self.assertEqual(self.sql("SELECT outcome FROM finance.daily_records ORDER BY source_sequence"), [("excluded_pm",), ("retained",)])
+ self.assertEqual(self.count("finance.v_active_daily_facts"), 1)
+ self.assertEqual(self.monthly_worker().process_next().status, "published")
+
+ def test_xml_failed_batch_keeps_pm_audit_without_monthly_commit(self):
+ outcome = self.coordinator.submit("source.xml", xml_document(pm_row(), reservation(2, departure="2026-07-26")).encode())
+ self.assertEqual(outcome["status"], "failed")
+ self.assertEqual(self.sql("SELECT version_status,source_rows,excluded_pm_rows,validation_failed_rows FROM finance.daily_versions"), [("rejected", 2, 1, 1)])
+ self.assertEqual(self.count("finance.current_daily_versions"), 0)
+ self.assertEqual(self.monthly_worker().process_next().status, "idle")
+
+ def test_database_pm_record_classification_guard_rejects_wrong_room_or_price(self):
+ import psycopg
+ self.make_first_pm()
+ self.assertEqual(self.execute().status, "succeeded")
+ with self.database.connect(autocommit=True) as connection:
+ for statement in (
+ "UPDATE finance.daily_records SET room_category_label='PM-SUITE' WHERE source_sequence=1",
+ "UPDATE finance.daily_records SET real_price=0 WHERE source_sequence=1",
+ "UPDATE finance.daily_records SET channel_key='Group' WHERE source_sequence=1",
+ "UPDATE finance.daily_records SET decision_codes=ARRAY['RATE_CODE_EXCLUDED'] WHERE source_sequence=1",
+ ):
+ with self.subTest(statement=statement), self.assertRaises(psycopg.errors.CheckViolation):
+ connection.execute(statement)
+
+
+if __name__ == "__main__":
+ unittest.main()
diff --git a/tests/test_arr_pm_review.py b/tests/test_arr_pm_review.py
new file mode 100644
index 0000000..3e9b5d5
--- /dev/null
+++ b/tests/test_arr_pm_review.py
@@ -0,0 +1,98 @@
+"""PM eligibility is shared with processing, before source-field completion."""
+import json
+from pathlib import Path
+import tempfile
+import unittest
+
+from arr_processing.policy import load_processor_policy
+from arr_web.arr_data_review import DataFieldReviews
+from arr_web.contracts import PortalError
+from tests.test_arr_data_review import source_document, raw, gap, CONTEXT, DAY, REQUEST, ACTOR, SOURCE_MANIFEST
+
+
+class PMReviewTests(unittest.TestCase):
+ def setUp(self):
+ temporary = tempfile.TemporaryDirectory(prefix="arr-pm-review-")
+ self.addCleanup(temporary.cleanup)
+ self.root = Path(temporary.name) / "reviews"
+ self.policy = load_processor_policy(Path(__file__).resolve().parents[1])
+ self.service = DataFieldReviews(root=self.root, policy=self.policy, context=CONTEXT)
+ self.document = source_document(2)
+
+ def prepare(self):
+ return self.service.prepare(REQUEST, raw(self.document), SOURCE_MANIFEST, DAY)
+
+ def test_explicit_pm_needs_no_fields_and_does_not_hide_other_rows(self):
+ self.document["records"][0]["fields"]["ROOM_CATEGORY_LABEL"] = {"state": "available", "value": " pm "}
+ for field in ("DISP_ROOM_NO", "RATE_CODE", "BLOCK_CODE", "ADULTS"):
+ gap(self.document, field)
+ gap(self.document, "DISP_ROOM_NO", sequence=2)
+ review = self.prepare()
+ self.assertEqual([item["item_id"] for item in review["items"]], ["2:DISP_ROOM_NO"])
+ self.assertEqual(review["excluded_pm_count"], 1)
+ self.assertEqual(review["excluded_cancelled_count"], 0)
+ with self.assertRaisesRegex(PortalError, "PM"):
+ self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "IGNORED", ACTOR)
+
+ def test_zero_price_or_nine_thousand_room_number_does_not_identify_pm(self):
+ for row in self.document["records"]:
+ row["fields"]["EFFECTIVE_RATE_AMOUNT"] = {"state": "available", "value": "0"}
+ row["fields"]["DISP_ROOM_NO"] = {"state": "available", "value": "9002"}
+ row["fields"]["ROOM_CATEGORY_LABEL"] = {"state": "available", "value": "PMX"}
+ gap(self.document, "ADULTS")
+ review = self.prepare()
+ self.assertEqual(review["excluded_pm_count"], 0)
+ self.assertIn("1:ADULTS", [item["item_id"] for item in review["items"]])
+
+ def test_cancelled_pm_is_counted_once_as_cancelled(self):
+ row = self.document["records"][0]
+ row["reservation_status"] = "Cancelled"
+ row["fields"]["ROOM_CATEGORY_LABEL"] = {"state": "available", "value": "PM"}
+ gap(self.document, "BLOCK_CODE", sequence=2)
+ review = self.prepare()
+ self.assertEqual((review["excluded_cancelled_count"], review["excluded_pm_count"]), (1, 0))
+
+ def test_prior_decisions_remain_audited_when_completed_room_type_is_pm(self):
+ gap(self.document, "BLOCK_CODE")
+ gap(self.document, "ROOM_CATEGORY_LABEL")
+ review = self.prepare()
+ original = self.service.original(REQUEST)
+ review = self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "VERIFIED", ACTOR)
+ review = self.service.update(REQUEST, "1:ROOM_CATEGORY_LABEL", review["revision"], "PM", ACTOR)
+ self.assertEqual([item["item_id"] for item in review["items"]], ["1:ROOM_CATEGORY_LABEL"])
+ self.assertTrue(review["items"][0]["confirmed"])
+ self.assertEqual(review["excluded_pm_count"], 1)
+ self.assertTrue(review["can_finalize"])
+ with self.assertRaisesRegex(PortalError, "PM"):
+ self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "CHANGED", ACTOR)
+ self.service.finalize(REQUEST, review["revision"], ACTOR)
+ derived = json.loads(self.service.payload(REQUEST)[0])
+ self.assertEqual(self.service.original(REQUEST), original)
+ changes = derived["manual_data_review"]["manifest"]["changes"]
+ self.assertEqual({item["field"]: item["value"] for item in changes},
+ {"BLOCK_CODE": "VERIFIED", "ROOM_CATEGORY_LABEL": "PM"})
+
+ def test_existing_room_type_normalization_is_shared_by_exclusion_count_and_edit_guard(self):
+ self.document["records"][0]["fields"]["ROOM_CATEGORY_LABEL"] = {"state": "available", "value": "PM😀"}
+ gap(self.document, "BLOCK_CODE")
+ gap(self.document, "DISP_ROOM_NO", sequence=2)
+ review = self.prepare()
+ self.assertEqual(review["excluded_pm_count"], 1)
+ self.assertEqual([item["item_id"] for item in review["items"]], ["2:DISP_ROOM_NO"])
+ with self.assertRaisesRegex(PortalError, "PM"):
+ self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "IGNORED", ACTOR)
+
+ def test_manually_completed_room_type_can_be_corrected_before_finalization(self):
+ gap(self.document, "ROOM_CATEGORY_LABEL")
+ gap(self.document, "BLOCK_CODE")
+ review = self.prepare()
+ review = self.service.update(REQUEST, "1:ROOM_CATEGORY_LABEL", review["revision"], "PM", ACTOR)
+ self.assertEqual(review["pending_count"], 0)
+ review = self.service.update(REQUEST, "1:ROOM_CATEGORY_LABEL", review["revision"], "SU6", ACTOR)
+ self.assertEqual(review["excluded_pm_count"], 0)
+ self.assertEqual(review["pending_count"], 1)
+ self.assertIn("1:BLOCK_CODE", [item["item_id"] for item in review["items"]])
+
+
+if __name__ == "__main__":
+ unittest.main()
diff --git a/tests/test_arr_pm_scope.py b/tests/test_arr_pm_scope.py
new file mode 100644
index 0000000..00e471d
--- /dev/null
+++ b/tests/test_arr_pm_scope.py
@@ -0,0 +1,228 @@
+"""Offline PM scope checks for both daily inputs and the frozen legacy replay."""
+from __future__ import annotations
+
+import argparse
+import contextlib
+import copy
+import io
+import json
+from pathlib import Path
+import tempfile
+import unittest
+import xml.etree.ElementTree as ET
+
+from tests.test_arr_opera_daily_ingest import (
+ core, reservation, run_processor, validator, xml_document,
+)
+from tests.test_arr_data_review import source_document
+
+
+def pm_xml_row(sequence, *, label="PM", room=None, status=None, incomplete=False,
+ rate_amount="900"):
+ node = ET.fromstring(reservation(sequence, room=room, rate_amount=rate_amount))
+ node.find("ROOM_CATEGORY_LABEL").text = label
+ if status is not None:
+ ET.SubElement(node, "SHORT_RESV_STATUS").text = status
+ if incomplete:
+ for field, value in (("RATE_CODE", ""), ("DISP_ROOM_NO", ""),
+ ("COMPANY_NAME", ""), ("FULL_NAME", ""),
+ ("ADULTS", "-1"), ("NO_OF_ROOMS", "0"),
+ ("EFFECTIVE_RATE_AMOUNT", "INVALID"),
+ ("TRUNC_END", "2026-07-26")):
+ node.find(field).text = value
+ return ET.tostring(node, encoding="unicode")
+
+
+class PMScopeTests(unittest.TestCase):
+ def setUp(self):
+ temporary = tempfile.TemporaryDirectory(prefix="arr-pm-scope-")
+ self.addCleanup(temporary.cleanup)
+ self.root = Path(temporary.name)
+
+ def classify_xml(self, *rows, apply_scope=True):
+ path = self.root / "source.xml"
+ path.write_text(xml_document(*rows))
+ business_date, nodes = core.read_xml(path)
+ return core.classify_input_records(nodes, business_date, apply_scope=apply_scope)
+
+ def run_data(self, document, name="output"):
+ path = self.root / (name + ".json")
+ path.write_text(json.dumps(document))
+ output = self.root / name
+ args = argparse.Namespace(xml=None, data_json=str(path), output_dir=str(output),
+ result_json=str(output / "result.json"),
+ structured_result_json=str(output / "structured-result.json"))
+ before = path.read_bytes()
+ with contextlib.redirect_stdout(io.StringIO()):
+ exit_code = core.process(args)
+ self.assertEqual(path.read_bytes(), before)
+ return exit_code, json.loads((output / "structured-result.json").read_text())
+
+ def test_xml_pm_precedes_rate_and_required_fields_but_cancelled_takes_precedence(self):
+ rows, eligible, removed_rate, duplicates, errors = self.classify_xml(
+ pm_xml_row(1, label=" pm ", incomplete=True),
+ pm_xml_row(2, status="CXL", incomplete=True), reservation(3))
+ self.assertEqual(errors, [])
+ self.assertEqual((removed_rate, duplicates), (0, 0))
+ self.assertEqual([row["_SOURCE_INDEX"] for row in eligible], [3])
+ self.assertEqual([row["_OUTCOME"] for row in rows],
+ ["excluded_pm", "excluded_cancelled", "pending"])
+ self.assertEqual(rows[0]["_DECISION_CODES"], ["ROOM_TYPE_PM_EXCLUDED"])
+ self.assertEqual(rows[1]["_DECISION_CODES"], ["RESERVATION_CANCELLED_EXCLUDED"])
+ self.assertEqual(rows[0]["ROOM_CATEGORY_LABEL"], "pm")
+ self.assertEqual(rows[0]["CONFIRMATION_NO"], "SYNTHETIC-CONF-1")
+ self.assertEqual(rows[0]["_SOURCE_LOCATION"], "reservation[1]")
+
+ def test_only_exact_explicit_pm_room_type_is_excluded(self):
+ for label in ("PM", "pm", " pM ", "\tPm\n"):
+ with self.subTest(label=label):
+ rows, _, _, _, errors = self.classify_xml(
+ pm_xml_row(1, label=label, incomplete=True), reservation(2))
+ self.assertEqual(errors, [])
+ self.assertTrue(core.is_pm_record(rows[0]))
+ for label in (None, "", "PM1", "PM ROOM", "P M", "SU6"):
+ with self.subTest(label=label):
+ rows, _, _, _, errors = self.classify_xml(
+ pm_xml_row(1, label=label, room="PM", rate_amount="0"), reservation(2))
+ self.assertEqual(errors, [])
+ self.assertFalse(core.is_pm_record(rows[0]))
+ self.assertEqual(rows[0]["_OUTCOME"], "pending")
+
+ def test_pm_does_not_consume_a_real_room_duplicate_key_or_create_price_issue(self):
+ rows, eligible, _, duplicates, errors = self.classify_xml(
+ pm_xml_row(1, room="SAME-ROOM", rate_amount="0"),
+ reservation(2, room="SAME-ROOM"), reservation(3, room="SAME-ROOM"))
+ self.assertEqual(errors, [])
+ self.assertEqual([row["_SOURCE_INDEX"] for row in eligible], [2])
+ self.assertEqual(duplicates, 1)
+ self.assertEqual([row["_OUTCOME"] for row in rows],
+ ["excluded_pm", "pending", "duplicate"])
+ self.assertEqual(rows[2]["_DUPLICATE_OF_SOURCE_SEQUENCE"], 2)
+ prices = core.load_price_map(core.PRICE_REFERENCE)
+ self.assertEqual(core.apply_prices_classified(eligible, prices), [])
+ self.assertEqual(core.review_issues(eligible, prices), [])
+
+ def test_xml_success_and_independent_replay_count_pm_and_cancelled_separately(self):
+ code, source, output, _, payload = run_processor(xml_document(
+ pm_xml_row(1, label=" pm ", incomplete=True),
+ pm_xml_row(2, status="CXL", incomplete=True), reservation(3)), self.root)
+ self.assertEqual(code, 0, payload["errors"])
+ self.assertEqual(payload["processor_version"], "4.4.0")
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ self.assertEqual(payload["outcome_counts"]["excluded_cancelled"], 1)
+ self.assertEqual(sum(payload["outcome_counts"].values()), 3)
+ self.assertEqual((payload["source_rows"], payload["output_rows"]), (3, 1))
+ self.assertEqual(payload["artifacts"]["source_xml"]["sha256"], core.sha256_file(source))
+ self.assertEqual(payload["records"][0]["room_category_label"], "pm")
+ for field in ("real_price", "total_price", "kb_amount", "channel_key", "pricing_method"):
+ self.assertIsNone(payload["records"][0][field])
+
+ # A structurally valid cancellation reclassification cannot replace PM evidence.
+ forged = copy.deepcopy(payload)
+ forged["records"][0]["outcome"] = "excluded_cancelled"
+ forged["records"][0]["decision_codes"] = ["RESERVATION_CANCELLED_EXCLUDED"]
+ forged["outcome_counts"]["excluded_cancelled"] += 1
+ forged["outcome_counts"]["excluded_pm"] -= 1
+ core.validate_structured_completeness(forged)
+ structured = output / "structured-result.json"
+ structured.write_text(json.dumps(forged))
+ errors = validator.validate(argparse.Namespace(xml=str(source),
+ daily=str(next(output.glob("*.xlsx"))), result_json=str(output / "result.json"),
+ structured_result_json=str(structured), price_reference=str(core.PRICE_REFERENCE)))
+ self.assertIn("OUTPUT_STRUCTURED_OUTCOME_MISMATCH", {error.code for error in errors})
+ self.assertIn("OUTPUT_STRUCTURED_RECORD_MISMATCH", {error.code for error in errors})
+
+ def test_pm_self_validation_and_schemas_require_explicit_type_reason_and_null_facts(self):
+ code, _, _, _, payload = run_processor(xml_document(
+ pm_xml_row(1, label="pM"), reservation(2)), self.root)
+ self.assertEqual(code, 0, payload["errors"])
+ for field, value in (("room_category_label", "SU6"),
+ ("decision_codes", ["RATE_CODE_NOT_WHITELISTED"]),
+ ("real_price", 0), ("total_price", 0), ("kb_amount", 0),
+ ("channel_key", "QBD"), ("pricing_method", "manual_review")):
+ with self.subTest(field=field):
+ forged = copy.deepcopy(payload)
+ forged["records"][0][field] = value
+ with self.assertRaises(core.ProcessingFailure):
+ core.validate_structured_completeness(forged)
+ forged = copy.deepcopy(payload)
+ del forged["outcome_counts"]["excluded_pm"]
+ with self.assertRaises(core.ProcessingFailure):
+ core.validate_structured_completeness(forged)
+ for name in ("structured-result.schema.json", "data-structured-result.schema.json"):
+ schema = json.loads((core.SKILL_ROOT / "references" / name).read_text())
+ self.assertEqual(schema["properties"]["processor_version"]["const"], "4.4.0")
+ self.assertIn("excluded_pm", schema["properties"]["outcome_counts"]["required"])
+ rule = next(item for item in schema["$defs"]["record"]["allOf"]
+ if item["if"]["properties"]["outcome"].get("const") == "excluded_pm")
+ self.assertEqual(rule["then"]["properties"]["decision_codes"]["const"],
+ ["ROOM_TYPE_PM_EXCLUDED"])
+ self.assertEqual(rule["then"]["properties"]["real_price"]["type"], "null")
+
+ def test_xml_review_and_failure_keep_pm_outside_pricing_and_other_validation(self):
+ code, _, _, _, payload = run_processor(xml_document(
+ pm_xml_row(1, incomplete=True), reservation(2, rate_amount="8765")),
+ self.root / "review")
+ self.assertEqual(code, 0, payload["errors"])
+ self.assertEqual(payload["status"], "review_required")
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ self.assertEqual(payload["review_required_rows"], 1)
+ self.assertEqual(payload["review_issues"][0]["affected_records"], 1)
+ invalid = ET.fromstring(reservation(2))
+ invalid.find("DISP_ROOM_NO").text = ""
+ code, _, _, _, payload = run_processor(xml_document(
+ pm_xml_row(1, incomplete=True), ET.tostring(invalid, encoding="unicode")),
+ self.root / "failure")
+ self.assertNotEqual(code, 0)
+ self.assertEqual(payload["status"], "failed")
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ self.assertEqual(payload["outcome_counts"]["validation_failed"], 1)
+
+ def test_data_success_review_and_failure_keep_pm_outside_business_rules(self):
+ baseline = source_document(2)
+ for field in baseline["records"][0]["fields"]:
+ baseline["records"][0]["fields"][field] = {"state": "missing", "value": None}
+ baseline["records"][0]["fields"]["ROOM_CATEGORY_LABEL"] = {"state": "available", "value": " pm "}
+ baseline["input_complete"] = False
+ baseline["status"] = "collected_with_gaps"
+ code, payload = self.run_data(baseline)
+ self.assertEqual(code, 0, payload["errors"])
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ self.assertEqual(payload["output_rows"], 1)
+ self.assertEqual(payload["removed_by_rate_code"], 0)
+ self.assertEqual(payload["records"][0]["room_category_label"], "pm")
+ review = copy.deepcopy(baseline)
+ review["records"][1]["fields"]["EFFECTIVE_RATE_AMOUNT"]["value"] = "8765"
+ code, payload = self.run_data(review, "review")
+ self.assertEqual(code, 0, payload["errors"])
+ self.assertEqual(payload["status"], "review_required")
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ self.assertEqual(payload["review_required_rows"], 1)
+ failed = copy.deepcopy(baseline)
+ failed["records"][1]["fields"]["DISP_ROOM_NO"] = {"state": "missing", "value": None}
+ code, payload = self.run_data(failed, "failure")
+ self.assertNotEqual(code, 0)
+ self.assertEqual(payload["status"], "failed")
+ self.assertEqual(payload["outcome_counts"]["excluded_pm"], 1)
+ self.assertEqual(payload["outcome_counts"]["validation_failed"], 1)
+
+ def test_frozen_legacy_v3_replay_keeps_pm_and_never_projects_new_outcomes(self):
+ code, _, _, _, payload = run_processor(xml_document(
+ pm_xml_row(1), reservation(2)), self.root, legacy_v3_output=True)
+ self.assertEqual(code, 0, payload["errors"])
+ self.assertEqual(payload["result_schema_version"], "3.0")
+ self.assertEqual(payload["output_rows"], 2)
+ self.assertEqual(set(payload["outcome_counts"]), core.LEGACY_DIRECT_FINAL_OUTCOMES)
+ self.assertEqual([row["outcome"] for row in payload["records"]], ["retained", "retained"])
+ rows, _, _, _, _ = self.classify_xml(pm_xml_row(1), reservation(2))
+ for function in (core.build_legacy_direct_structured_result,
+ core.build_legacy_direct_failed_structured_result):
+ with self.subTest(function=function.__name__), self.assertRaises(core.ProcessingFailure):
+ if function is core.build_legacy_direct_structured_result:
+ function(None, rows[:1], [], self.root / "source.xml", None, None)
+ else:
+ function(None, rows[:1], self.root / "source.xml", None, None, [])
+
+
+if __name__ == "__main__":
+ unittest.main()