..................................F................F.F.F..FF.FF...FFFF.F [ 84%] FFFFF.FF..... [100%] =================================== FAILURES =================================== __ test_incomplete_period_and_different_end_date_block_confirm_without_writes __ def test_incomplete_period_and_different_end_date_block_confirm_without_writes(): conn = _conn() zip_raw, overview_raw = _zip(statement_end="26.07.2026"), _overview() request = _request(zip_raw, overview_raw) before = conn.total_changes preview = service.preview_postfinance_import(conn, request) > assert preview["confirm_allowed"] is False E assert True is False tests/unit/test_sprint20h1_postfinance_upload.py:457: AssertionError _ test_golden_raiffeisen_2026_01_26_swisslos_chf100_is_safe_v3_transfer_and_confirm_is_noop_twice _ def test_golden_raiffeisen_2026_01_26_swisslos_chf100_is_safe_v3_transfer_and_confirm_is_noop_twice() -> None: conn = database() payload = transfer_payload() preview = preview_household_import(conn, payload) pair = preview["transfer_pairs"][0] assert pair["pairing_class"] == "safe" assert pair["amount"] == "100.00" and pair["currency"] == "CHF" assert preview["pairing_version"] == "transfer_pairing_v3" > first = confirm_household_import(conn, confirm_payload(payload, preview)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ tests/unit/test_household_import_v1_golden.py:161: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '5ed3858b1946dd0bac429c7cc3bd2112c73c806464399088aeb83fa89eab560e', 'confirm': True, 'files':...reference': 'SYN-AKB-001'}], 'preview_fingerprint': '976a9c5a43b0782a5db507197031ddc263389d1c14d4f6a0dc42f1519c6aced2'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _ test_explicit_review_batch_confirm_registers_rows_without_promoting_unclear_cashflows _ def test_explicit_review_batch_confirm_registers_rows_without_promoting_unclear_cashflows() -> None: conn = database() payload = { "files": [ { "profile": "raiffeisen_bank", "csv_text": raiffeisen_csv("123.45", "Unresolved incoming payment"), } ] } preview = preview_household_import(conn, payload) assert preview["confirmable"] is True assert preview["business_ready_for_confirm"] is False with pytest.raises(HTTPException) as blocked: confirm_household_import(conn, confirm_payload(payload, preview)) assert blocked.value.status_code == 409 > confirmed = confirm_household_import( conn, confirm_payload(payload, preview) | {"confirm_review_candidates": True}, ) tests/unit/test_household_import_v1_golden.py:201: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': 'd812c74324313a0c03d02bd67686f863a5eadd48b8751fd2a2148f7a7d5d1842', 'confirm': True, 'confirm...uta Date\nSYN-RAI-001;2026-01-26;Unresolved incoming payment;123.45;2026-01-26\n', 'profile': 'raiffeisen_bank'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError ________ test_file_source_row_and_logical_duplicates_are_distinguished _________ def test_file_source_row_and_logical_duplicates_are_distinguished() -> None: conn = database() visa = "TransactionId,CardId,Date,Amount,Currency,MerchantName\nT-001,SYN-CARD-001,2026-02-01,-12.50,CHF,Synthetic Shop\n" payload = {"files": [{"profile": "visa_credit_card", "csv_text": visa}]} payload, preview = business_ready_preview(conn, payload) > confirm_household_import(conn, confirm_payload(payload, preview)) tests/unit/test_household_import_v1_golden.py:230: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '2bfcb4ac0703c2066c6314846a5d5595ec036edfe039d7d85666d4783223d3ab', 'category_overrides': {'r...urrency,MerchantName\nT-001,SYN-CARD-001,2026-02-01,-12.50,CHF,Synthetic Shop\n', 'profile': 'visa_credit_card'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _ test_migros_links_detail_without_double_counting_and_difference_over_cent_reviews _ def test_migros_links_detail_without_double_counting_and_difference_over_cent_reviews() -> None: conn = database() visa = "TransactionId,CardId,Date,Amount,Currency,MerchantName\nM-1,SYN-CARD-001,2026-03-01,-60.00,CHF,Migros Synthetic\n" receipt = "Datum;Zeit;Filiale;Kassennummer;Transaktionsnummer;Artikel;Umsatz\n" \ "2026-03-01;12:00;Synthetic Store;1;99;Item A;24.00\n" \ "2026-03-01;12:00;Synthetic Store;1;99;Item B;36.00\n" payload = {"files": [{"profile": "visa_credit_card", "csv_text": visa}, {"profile": "migros_receipts", "csv_text": receipt}]} preview = preview_household_import(conn, payload) assert preview["receipt_links"][0]["status"] == "linked" > result = confirm_household_import(conn, confirm_payload(payload, preview)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ tests/unit/test_household_import_v1_golden.py:308: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '2a3f72b586e69cd0a292d4ceb4126504774d8094c007ed81a506c1a3bfb4b91c', 'confirm': True, 'files':...ofile': 'migros_receipts'}], 'preview_fingerprint': 'b94a90fa9e47c217aaaf1b9cbf7a99dc0edbe385000f039da833d4eb872787d7'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _ test_migros_receipt_replay_preserves_duplicate_disposition_and_writes_nothing _ def test_migros_receipt_replay_preserves_duplicate_disposition_and_writes_nothing() -> None: conn = database() receipt = ( "Datum;Zeit;Filiale;Kassennummer;Transaktionsnummer;Artikel;Umsatz\n" "2026-03-01;12:00;Synthetic Store;1;99;Item A;24.00\n" "2026-03-01;12:00;Synthetic Store;1;99;Item B;36.00\n" ) payload = {"files": [{"profile": "migros_receipts", "csv_text": receipt}]} first = preview_household_import(conn, payload) > confirm_household_import(conn, confirm_payload(payload, first)) tests/unit/test_household_import_v1_golden.py:331: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '0bfe78b69785d5013aa13cbb7cfa2c6257255c1bf765038732faf9b97ea706b1', 'confirm': True, 'files':...ofile': 'migros_receipts'}], 'preview_fingerprint': 'bd824cbf5d44cb7d8850fe941ba51ee483a290a211bc354ef0351174242d6b3c'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError ________ test_idempotent_confirm_recomputes_input_identity_before_noop _________ def test_idempotent_confirm_recomputes_input_identity_before_noop() -> None: conn = database() payload = transfer_payload() preview = preview_household_import(conn, payload) > confirm_household_import(conn, confirm_payload(payload, preview)) tests/unit/test_household_import_v1_golden.py:388: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '024ed5c113fe95f09b1056cd5338f3e40383a710b87e712f77e61b82a1365acd', 'confirm': True, 'files':...reference': 'SYN-AKB-001'}], 'preview_fingerprint': 'bfe1b96c182343a054c0620ec41763182959e0740161477db3a9335a19d5d213'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError ___________ test_selected_mapping_is_bound_and_mismatch_fails_closed ___________ def test_selected_mapping_is_bound_and_mismatch_fails_closed() -> None: conn = database() payload = {"files": [{"profile": "raiffeisen_bank", "csv_text": raiffeisen_csv("-10.00")}]} mapping_ids = { row["source_type"]: row["mapping_id"] for row in conn.execute("SELECT mapping_id,source_type FROM household_account_source_mappings") } payload["files"][0]["mapping_id"] = mapping_ids["visa_credit_card"] mismatch = preview_household_import(conn, payload) assert mismatch["confirmable"] is False assert mismatch["errors"][0]["code"] == "account_source_mapping_missing_or_mismatch" payload["files"][0]["mapping_id"] = mapping_ids["raiffeisen_bank"] payload, matched = business_ready_preview(conn, payload) assert matched["confirmable"] is True > confirm_household_import(conn, confirm_payload(payload, matched)) tests/unit/test_household_import_v1_golden.py:411: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '172beac00c277e9f83832311c79f50fbd5a221fa7924ac2680975499eae0de4b', 'category_overrides': {'r...S E-COMMERCE;-10.00;2026-01-26\n', 'mapping_id': 'hhmap_b02be474832bd16d791977b0', 'profile': 'raiffeisen_bank'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError ____ test_same_request_duplicate_and_pending_final_have_one_writable_final _____ def test_same_request_duplicate_and_pending_final_have_one_writable_final() -> None: conn = database() rows = "TransactionId,CardId,Date,Amount,Currency,MerchantName,Status\n" \ "P-1,SYN-CARD-001,2026-02-02,-20.00,CHF,Synthetic,pending\n" \ "P-1,SYN-CARD-001,2026-02-02,-20.00,CHF,Synthetic,booked\n" \ "P-1,SYN-CARD-001,2026-02-02,-20.00,CHF,Synthetic,booked\n" payload = {"files": [{"profile": "visa_credit_card", "csv_text": rows}]} payload, preview = business_ready_preview(conn, payload) assert preview["counts"]["duplicates"] == 2 assert [row["disposition"] for row in preview["rows"]].count("candidate") == 1 > result = confirm_household_import(conn, confirm_payload(payload, preview)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ tests/unit/test_household_import_v1_golden.py:476: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '29f1c15beaa5831d51b98da975fbaee2d171d0eee02d2133b7e19f947f028736', 'category_overrides': {'r...CHF,Synthetic,booked\nP-1,SYN-CARD-001,2026-02-02,-20.00,CHF,Synthetic,booked\n', 'profile': 'visa_credit_card'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _ test_distinct_provider_transactions_with_same_business_fields_are_not_silently_suppressed _ def test_distinct_provider_transactions_with_same_business_fields_are_not_silently_suppressed() -> None: conn = database() csv_text = "TransactionId,CardId,Date,Amount,Currency,MerchantName\n" \ "REAL-1,SYN-CARD-001,2026-02-03,-5.00,CHF,Synthetic Coffee\n" \ "REAL-2,SYN-CARD-001,2026-02-03,-5.00,CHF,Synthetic Coffee\n" payload = {"files": [{"profile": "visa_credit_card", "csv_text": csv_text}]} payload, preview = business_ready_preview(conn, payload) dispositions = [row["disposition"] for row in preview["rows"]] assert dispositions == ["candidate", "review"] assert preview["rows"][1]["classification"] == "possible_logical_duplicate" > confirm_household_import(conn, confirm_payload(payload, preview)) tests/unit/test_household_import_v1_golden.py:501: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '27c78da8f5739b64222916f45781ba9821416104d0cc8d10b89c8c9995cc3043', 'category_overrides': {'r...F,Synthetic Coffee\nREAL-2,SYN-CARD-001,2026-02-03,-5.00,CHF,Synthetic Coffee\n', 'profile': 'visa_credit_card'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _________ test_legacy_productive_overlap_is_reviewed_not_counted_again _________ def test_legacy_productive_overlap_is_reviewed_not_counted_again() -> None: conn = database() account_id = conn.execute( "SELECT budget_account_id FROM household_account_source_mappings WHERE source_type='raiffeisen_bank'" ).fetchone()[0] conn.execute( """INSERT INTO budget_transactions( budget_transaction_id,account_id,transaction_type,transaction_date,booking_date, description,amount_original,currency_original,fx_rate_to_chf,amount_chf,fx_status, status,source_type,created_at,updated_at) VALUES ('legacy-overlap',?,'expense','2026-01-26','2026-01-26', 'Synthetic Grocery','-25.50','CHF','1','-25.50','not_needed', 'confirmed','manual','2026-01-15T00:00:00Z','2026-01-15T00:00:00Z')""", (account_id,), ) conn.commit() payload = {"files": [{"profile": "raiffeisen_bank", "csv_text": raiffeisen_csv("-25.50", "Synthetic Grocery")}]} payload, preview = business_ready_preview(conn, payload) assert preview["rows"][0]["classification"] == "possible_legacy_duplicate" assert preview["rows"][0]["disposition"] == "review" > confirm_household_import(conn, confirm_payload(payload, preview)) tests/unit/test_household_import_v1_golden.py:527: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '2ef549cf10a283f238d52ca821ea0d0c3ecef438f3a63e8382a2a6bb16390c66', 'category_overrides': {'r...Amount;Valuta Date\nSYN-RAI-001;2026-01-26;Synthetic Grocery;-25.50;2026-01-26\n', 'profile': 'raiffeisen_bank'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _ test_household_confirm_preserves_outer_transaction_and_rolls_back_only_its_savepoint _ def test_household_confirm_preserves_outer_transaction_and_rolls_back_only_its_savepoint() -> None: conn = database() payload = transfer_payload() preview = preview_household_import(conn, payload) conn.execute("INSERT INTO platforms(platform_id,name,platform_type,created_at) VALUES ('caller','Caller','bank','x')") > result = confirm_household_import(conn, confirm_payload(payload, preview)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ tests/unit/test_household_import_v1_golden.py:544: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': 'daf8f911f273932c56cd0809fec3ce8c0cf46a00969e1ddd5626eb59c44a1964', 'confirm': True, 'files':...reference': 'SYN-AKB-001'}], 'preview_fingerprint': 'badeb1b9d135aa3aa4a089b4365b585349ebb04a0ae874761799be049f1019e7'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError ___ test_clear_chf_card_refund_becomes_productive_refund_only_after_confirm ____ def test_clear_chf_card_refund_becomes_productive_refund_only_after_confirm() -> None: conn = database() csv_text = "TransactionId,CardId,Date,Amount,Currency,MerchantName\nRF-1,SYN-CARD-001,2026-02-10,12.50,CHF,Synthetic refund\n" payload = {"files": [{"profile": "visa_credit_card", "csv_text": csv_text}]} preview = preview_household_import(conn, payload) assert preview["rows"][0]["classification"] == "credit_card_refund" assert conn.execute("SELECT COUNT(*) FROM budget_transactions").fetchone()[0] == 0 > confirm_household_import(conn, confirm_payload(payload, preview)) tests/unit/test_household_import_v1_golden.py:580: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': 'b53e36240cf19f8296925a16e61656c57a6707cc15f1e35619c79f1960a493df', 'confirm': True, 'files':...file': 'visa_credit_card'}], 'preview_fingerprint': '6d469cddc79d6553794c123e9f7e1b07458ec4fa28cacf60c3867540118de21f'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _____ test_unique_reversal_deterministically_offsets_original_transaction ______ def test_unique_reversal_deterministically_offsets_original_transaction() -> None: conn = database() purchase = { "files": [{ "profile": "visa_credit_card", "csv_text": "Date,Amount,Currency,MerchantName,TransactionId,CardId\n2026-02-01,-8.25,CHF,Synthetic reverse purchase,rev-purchase,SYN-CARD-001\n", }] } purchase, purchase_preview = business_ready_preview(conn, purchase) > confirm_household_import(conn, confirm_payload(purchase, purchase_preview)) tests/unit/test_household_import_v1_golden.py:607: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '946ed17d4f869628c349ca0efc4e859e3a0dea49fe287f8ab108655e571a6978', 'category_overrides': {'r...Id\n2026-02-01,-8.25,CHF,Synthetic reverse purchase,rev-purchase,SYN-CARD-001\n', 'profile': 'visa_credit_card'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _____ test_ambiguous_reversal_remains_review_and_never_posts_automatically _____ def test_ambiguous_reversal_remains_review_and_never_posts_automatically() -> None: conn = database() purchases = { "files": [{ "profile": "visa_credit_card", "csv_text": ( "Date,Amount,Currency,MerchantName,TransactionId,CardId\n" "2026-02-01,-8.25,CHF,Repeated purchase,amb-1,SYN-CARD-001\n" "2026-02-02,-8.25,CHF,Repeated purchase,amb-2,SYN-CARD-001\n" ), }] } purchases, purchase_preview = business_ready_preview(conn, purchases) > confirm_household_import(conn, confirm_payload(purchases, purchase_preview)) tests/unit/test_household_import_v1_golden.py:666: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '3a54cc997769d2df1d91496de0a22562caa20002b5e9eab01fcbbe8783d9386e', 'category_overrides': {'r...amb-1,SYN-CARD-001\n2026-02-02,-8.25,CHF,Repeated purchase,amb-2,SYN-CARD-001\n', 'profile': 'visa_credit_card'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError ______ test_review_expense_preview_confirm_is_bound_atomic_and_idempotent ______ def test_review_expense_preview_confirm_is_bound_atomic_and_idempotent() -> None: conn = database() header = "TransactionId,CardId,Date,Amount,Currency,MerchantName\n" known = "".join( f"SAFE-{index},SYN-CARD-001,2026-{'01' if index < 31 else '02'}-{index + 1 if index < 31 else index - 30:02d},-3.00,CHF,Migros\n" for index in range(33) ) csv_text = header + known + ( "RV-1,SYN-CARD-001,2026-02-28,-8.25,CHF,Synthetic review expense\n" ) payload = {"files": [{"profile": "visa_credit_card", "csv_text": csv_text}]} import_preview = preview_household_import(conn, payload) > confirm_household_import(conn, confirm_payload(payload, import_preview)) tests/unit/test_household_import_v1_golden.py:694: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '90cf6f04fdf886138c8619e4ec547bc34015641e0eb16c8e3aa0b14ced916b32', 'confirm': True, 'files':...file': 'visa_credit_card'}], 'preview_fingerprint': 'd937341c0b7481fbac8e2a7887736d991773b9b20ba5c4f3fefff2089edb5770'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _____ test_migros_receipt_links_to_existing_money_movement_across_batches ______ def test_migros_receipt_links_to_existing_money_movement_across_batches() -> None: conn = database() visa = "TransactionId,CardId,Date,Amount,Currency,MerchantName\nMX-1,SYN-CARD-001,2026-03-01,-10.00,CHF,Migros Synthetic\n" money_payload = {"files": [{"profile": "visa_credit_card", "csv_text": visa}]} money_preview = preview_household_import(conn, money_payload) > confirm_household_import(conn, confirm_payload(money_payload, money_preview)) tests/unit/test_household_import_v1_golden.py:732: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': 'f4c09b34d6134a29f2089cbcba284a2844ef5e3d128caa8e6a4c1b7898541b02', 'confirm': True, 'files':...file': 'visa_credit_card'}], 'preview_fingerprint': '2615d04119ac5afd8c047e6a511a77c5a84f355f87d507969c642409c2c35eb6'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError __ test_migros_receipt_links_to_productive_transaction_not_consumed_candidate __ def test_migros_receipt_links_to_productive_transaction_not_consumed_candidate() -> None: conn = database() visa = "TransactionId,CardId,Date,Amount,Currency,MerchantName\nMX-2,SYN-CARD-001,2026-03-02,-10.00,CHF,Migros Productive\n" money_payload = {"files": [{"profile": "visa_credit_card", "csv_text": visa}]} money_preview = preview_household_import(conn, money_payload) > confirm_household_import(conn, confirm_payload(money_payload, money_preview)) tests/unit/test_household_import_v1_golden.py:767: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '778d3bb8c2d6f4b2ace743a6cca0cbbfc4775db4a51af4f179e4371900b0c338', 'confirm': True, 'files':...file': 'visa_credit_card'}], 'preview_fingerprint': 'd824975ca590fea135ad0c6887dd16e28c8e832e5d2e6f036b4ef62b54c0ed1f'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError _ test_authenticated_owner_attestation_api_binds_exact_cluster_and_confirms_neutral_transfer _ monkeypatch = <_pytest.monkeypatch.MonkeyPatch object at 0x74661976b6d0> def test_authenticated_owner_attestation_api_binds_exact_cluster_and_confirms_neutral_transfer( monkeypatch: pytest.MonkeyPatch, ) -> None: monkeypatch.setenv("JARVIS_FINANCE_OPERATOR_APPROVAL_KEY", "synthetic-operator-key") conn = database() def override_db(): yield conn app = create_app(write_mode="test") app.dependency_overrides[get_db] = override_db client = TestClient(app) csv_text = ( "IBAN;Booked At;Text;Credit/Debit Amount;Valuta Date\n" "SYN-RAI-001;2026-07-01;Synthetic Household Counterparty;-20.00;2026-07-01\n" "SYN-RAI-001;2026-07-02;Synthetic Household Counterparty;-30.00;2026-07-02\n" ) payload = {"files": [{"profile": "raiffeisen_bank", "csv_text": csv_text}]} base = client.post("/api/budget/household/imports/preview", json=payload) assert base.status_code == 200, base.text cluster = base.json()["merchant_clusters"][0] approval_body = { "confirm_owner_attestation": True, "preview_request": payload, "approval_evidence": "owner_attested_known_household_counterparty_v1", "cluster_token": cluster["cluster_token"], "approved_row_tokens": cluster["row_tokens"], } denied = client.post( "/api/budget/household/imports/owner-neutral-cluster-approval", headers={"X-Jarvis-Operator-Approval": "wrong"}, json=approval_body, ) assert denied.status_code == 403 subset = client.post( "/api/budget/household/imports/owner-neutral-cluster-approval", headers={"X-Jarvis-Operator-Approval": "synthetic-operator-key"}, json=approval_body | {"approved_row_tokens": cluster["row_tokens"][:1]}, ) assert subset.status_code == 422 approved = client.post( "/api/budget/household/imports/owner-neutral-cluster-approval", headers={"X-Jarvis-Operator-Approval": "synthetic-operator-key"}, json=approval_body, ) assert approved.status_code == 200, approved.text forged_decision = approved.json() public_material = { "scope": "owner_confirmed_neutral_household_transfer_v2", "approval_evidence": forged_decision["approval_evidence"], "cluster_token": forged_decision["cluster_token"], "row_tokens": sorted(forged_decision["approved_row_tokens"]), "baseline_fingerprint": forged_decision["approval_baseline_fingerprint"], "issued_at": forged_decision["approval_issued_at"], "expires_at": forged_decision["approval_expires_at"], "nonce": forged_decision["approval_nonce"], } forged_decision["approval_token"] = hashlib.sha256( json.dumps(public_material, ensure_ascii=False, separators=(",", ":"), sort_keys=True).encode() ).hexdigest() forged_preview = client.post( "/api/budget/household/imports/preview", json={**payload, "cluster_decisions": [forged_decision]}, ) assert forged_preview.status_code == 422 signed_payload = {**payload, "cluster_decisions": [approved.json()]} preview_response = client.post("/api/budget/household/imports/preview", json=signed_payload) assert preview_response.status_code == 200, preview_response.text preview = preview_response.json() assert preview["business_ready_for_confirm"] is True assert {item["transaction_semantics"] for item in preview["items"]} == { "user_confirmed_unmatched_transfer" } confirm_request = confirm_payload(signed_payload, preview) > confirmed = client.post( "/api/budget/household/imports/confirm", json=confirm_request, ) tests/unit/test_household_import_v1_golden.py:897: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/testclient.py:555: in post return super().post( ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/httpx/_client.py:1144: in post return self.request( ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/testclient.py:454: in request return super().request( ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/httpx/_client.py:825: in request return self.send(request, auth=auth, follow_redirects=follow_redirects) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/httpx/_client.py:914: in send response = self._send_handling_auth( ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/httpx/_client.py:942: in _send_handling_auth response = self._send_handling_redirects( ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/httpx/_client.py:979: in _send_handling_redirects response = self._send_single_request(request) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/httpx/_client.py:1014: in _send_single_request response = transport.handle_request(request) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/testclient.py:356: in handle_request raise exc ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/testclient.py:353: in handle_request portal.call(self.app, scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/anyio/from_thread.py:338: in call return cast(T_Retval, self.start_task_soon(func, *args).result()) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ../../../.local/share/uv/python/cpython-3.11.15-linux-x86_64-gnu/lib/python3.11/concurrent/futures/_base.py:456: in result return self.__get_result() ^^^^^^^^^^^^^^^^^^^ ../../../.local/share/uv/python/cpython-3.11.15-linux-x86_64-gnu/lib/python3.11/concurrent/futures/_base.py:401: in __get_result raise self._exception ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/anyio/from_thread.py:263: in _call_func retval = await retval_or_awaitable ^^^^^^^^^^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/applications.py:1159: in __call__ await super().__call__(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/applications.py:90: in __call__ await self.middleware_stack(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/middleware/errors.py:186: in __call__ raise exc ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/middleware/errors.py:164: in __call__ await self.app(scope, receive, _send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/middleware/base.py:193: in __call__ response = await self.dispatch_func(request, call_next) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ src/jarvis_finance/api/main.py:117: in block_untrusted_writes return await call_next(request) ^^^^^^^^^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/middleware/base.py:168: in call_next raise app_exc from app_exc.__cause__ or app_exc.__context__ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/middleware/base.py:144: in coro await self.app(scope, receive_or_disconnect, send_no_error) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/middleware/cors.py:88: in __call__ await self.app(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/middleware/exceptions.py:63: in __call__ await wrap_app_handling_exceptions(self.app, conn)(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/_exception_handler.py:53: in wrapped_app raise exc ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/_exception_handler.py:42: in wrapped_app await app(scope, receive, sender) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/middleware/asyncexitstack.py:18: in __call__ await self.app(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/routing.py:660: in __call__ await self.middleware_stack(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/routing.py:680: in app await route.handle(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/routing.py:276: in handle await self.app(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/routing.py:134: in app await wrap_app_handling_exceptions(app, request)(scope, receive, send) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/_exception_handler.py:53: in wrapped_app raise exc ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/_exception_handler.py:42: in wrapped_app await app(scope, receive, sender) ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/routing.py:120: in app response = await f(request) ^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/routing.py:674: in app raw_response = await run_endpoint_function( ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/routing.py:330: in run_endpoint_function return await run_in_threadpool(dependant.call, **values) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/starlette/concurrency.py:34: in run_in_threadpool return await anyio.to_thread.run_sync(func) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/anyio/to_thread.py:65: in run_sync return await get_async_backend().run_sync_in_worker_thread( ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/anyio/_backends/_asyncio.py:2641: in run_sync_in_worker_thread return await future ^^^^^^^^^^^^ ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/anyio/_backends/_asyncio.py:1033: in run result = context.run(func, *args) ^^^^^^^^^^^^^^^^^^^^^^^^ src/jarvis_finance/api/routers/budget.py:252: in household_import_confirm return confirm_household_import(conn, payload) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '8656ec2cef23819219466a98b45a4a9e8c04a50862406aa8c58de93170b2c3fb', 'cluster_decisions': [{'a...-01\nSYN-RAI-001;2026-07-02;Synthetic Household Counterparty;-30.00;2026-07-02\n', 'profile': 'raiffeisen_bank'}], ...} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError ________ test_review_summary_is_canonical_filtered_and_cursor_paginated ________ def test_review_summary_is_canonical_filtered_and_cursor_paginated() -> None: conn = database() header = "TransactionId,CardId,Date,Amount,Currency,MerchantName\n" known = "".join( f"SAFE-R-{index},SYN-CARD-001,2026-03-{index + 1:02d},-3.00,CHF,Migros\n" for index in range(24) ) payload = {"files": [{"profile": "visa_credit_card", "csv_text": header + known + "OPEN-1,SYN-CARD-001,2026-03-25,-8.25,CHF,Synthetic review expense\n"}]} preview = preview_household_import(conn, payload) > confirm_household_import(conn, confirm_payload(payload, preview)) tests/unit/test_household_import_v1_golden.py:940: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ conn = payload = {'baseline_fingerprint': '5ab268beeceec1fbf44f3e8029ac6d61014e20e41929f25c809ef2e0fedc1f82', 'confirm': True, 'files':...file': 'visa_credit_card'}], 'preview_fingerprint': 'f4ee93255b69d6cee79c97adce0ca6bf8f776915fd92be35ced626972b16fd73'} def confirm_household_import(conn: Connection, payload: dict[str, Any]) -> dict[str, Any]: expected_preview = str(payload.get("preview_fingerprint") or "") expected_baseline = str(payload.get("baseline_fingerprint") or "") if not expected_preview or not expected_baseline or payload.get("confirm") is not True: raise HTTPException(status_code=422, detail="confirm=true and both fingerprints are required") validated_files = _validated_files(payload) current_input = _input_fingerprint( validated_files, payload.get("category_overrides"), payload.get("user_decisions"), payload.get("cluster_decisions"), ) existing = conn.execute("SELECT * FROM household_import_batches WHERE preview_fingerprint=?", (expected_preview,)).fetchone() if existing: if current_input != existing["input_fingerprint"] or expected_baseline != existing["baseline_fingerprint"]: raise HTTPException(status_code=409, detail="confirmed preview does not match this input") current_time = int(time.time()) for decision in payload.get("cluster_decisions") or []: if decision.get("decision_type") != USER_CONFIRMED_UNMATCHED_TRANSFER: continue try: approval_expires_at = int(decision.get("approval_expires_at")) except (TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail="bounded neutral cluster approval timestamps are invalid") from exc if current_time > approval_expires_at: raise HTTPException(status_code=409, detail="bounded neutral cluster approval has expired") return {"status": "confirmed", "batch_id": existing["batch_id"], "idempotent": True, "counts": {key: existing[key] for key in ("file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count")}} reconstructed = _preview_household_import_internal(conn, payload) if reconstructed["preview_fingerprint"] != expected_preview: raise HTTPException(status_code=409, detail="preview fingerprint mismatch; preview again") if reconstructed["baseline_fingerprint"] != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") if not reconstructed["confirmable"]: raise HTTPException(status_code=409, detail="preview is not technically confirmable") review_batch_confirmed = payload.get("confirm_review_candidates") is True if not reconstructed["business_ready_for_confirm"] and not review_batch_confirmed: raise HTTPException(status_code=409, detail="preview is not business-ready for confirm") replay_writable = [ row for row in reconstructed["rows"] if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"} ] if not replay_writable and reconstructed["files"] and all( bool(file["duplicate"]) for file in reconstructed["files"] ): placeholders = ",".join("?" for _ in reconstructed["files"]) batch_rows = conn.execute( f"""SELECT DISTINCT b.* FROM household_import_batches b JOIN household_import_files f ON f.batch_id=b.batch_id WHERE f.file_fingerprint IN ({placeholders}) ORDER BY b.confirmed_at,b.batch_id""", tuple(file["file_fingerprint"] for file in reconstructed["files"]), ).fetchall() if len(batch_rows) == 1: prior = batch_rows[0] return { "status": "confirmed", "batch_id": prior["batch_id"], "idempotent": True, "counts": { key: prior[key] for key in ( "file_count", "row_count", "candidate_count", "transfer_pair_count", "duplicate_count", "review_count", "receipt_link_count", ) }, } timestamp = now() batch_id = "hhbatch_" + expected_preview[:24] reconstructed_rows = reconstructed["rows"] # map fields needed by writes from the current stable mapping table for entry in validated_files: file_item, parsed = entry["item"], entry["parsed"] originals = _migros_rows(parsed["rows"]) if parsed["profile"] == "migros_receipts" else [r for i, raw in enumerate(parsed["rows"], 1) if (r := _normal_row(parsed["profile"], raw, i, file_item.get("source_reference")))] for original in originals: targets = [row for row in reconstructed_rows if row["source_row_fingerprint"] == original["source_row_fingerprint"]] for target in targets: target["mapping"] = _mapping_for_file( conn, original["source_type"], original["source_reference"], file_item, ) if original["source_type"] != "migros_receipts" else None target["line_items"] = original.get("line_items", []) writable = [row for row in reconstructed_rows if row["disposition"] in {"candidate", "review", "receipt_detail", "transfer_confirmed"}] rows_by_fp = {row["source_row_fingerprint"]: row for row in writable} had_outer_transaction = conn.in_transaction if had_outer_transaction: conn.execute("SAVEPOINT household_confirm") else: conn.execute("BEGIN IMMEDIATE") try: # Close the preview-to-write TOCTOU window after acquiring the write lock/savepoint. if _baseline(conn) != expected_baseline: raise HTTPException(status_code=409, detail="database baseline changed; preview again") counts = reconstructed["counts"] expected_pair_count = sum(pair["pairing_class"] == "safe" for pair in reconstructed["transfer_pairs"]) expected_link_count = sum( link["status"] == "linked" and link["receipt_row_fingerprint"] in rows_by_fp for link in reconstructed["receipt_links"] ) audit_id = record_audit_event( conn, source="household_import", action="household_import_confirmed", entity_type="household_import_batch", entity_id=batch_id, new_values={ "contract_version": CONTRACT_VERSION, "classification_version": CLASSIFICATION_VERSION, "pairing_version": PAIRING_VERSION, "preview_fingerprint": expected_preview, "sources": sorted({str(row["source_type"]) for row in reconstructed_rows}), "masked_accounts": sorted({ "••••" + _sha(str(row["mapping"]["budget_account_id"]))[-4:] for row in reconstructed_rows if row.get("mapping") }), "classification_origins": sorted({ str(row.get("classification_v2", {}).get("origin") or "unresolved") for row in reconstructed_rows }), "category_override_count": len(payload.get("category_overrides") or {}), "user_decision_count": len(payload.get("user_decisions") or {}), "cluster_decision_count": len(payload.get("cluster_decisions") or []), "owner_attested_neutral_cluster_count": sum( cluster.get("decision_type") == USER_CONFIRMED_UNMATCHED_TRANSFER for cluster in reconstructed.get("merchant_clusters", []) ), "owner_attestation_evidence_version": OWNER_ATTESTED_TRANSFER_EVIDENCE, "cluster_decision_version": CLUSTER_DECISION_VERSION, "user_decision_version": USER_DECISION_VERSION, "business_ready_for_confirm": reconstructed["business_ready_for_confirm"], "review_batch_confirmed": review_batch_confirmed, "review_candidates_remain_unconfirmed": bool( review_batch_confirmed and not reconstructed["business_ready_for_confirm"] ), "counts": counts, "expected_writes": reconstructed["expected_writes"], "actual_writes": reconstructed["expected_writes"], "status": "confirmed", }, created_by="user", ) # Insert the FK parent before candidates, files, items, and receipt links. conn.execute( """INSERT INTO household_import_batches VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (batch_id, CONTRACT_VERSION, expected_preview, expected_baseline, reconstructed["input_fingerprint"], counts["files"], counts["rows"], len(writable), expected_pair_count, counts["duplicates"], counts["review"], expected_link_count, "confirmed", audit_id, timestamp, "user"), ) candidate_ids: dict[str, str] = {} for row in writable: candidate_ids[row["source_row_fingerprint"]] = _insert_candidate(conn, row, batch_id, timestamp) _confirm_safe_nontransfer_candidate( conn, row, candidate_ids[row["source_row_fingerprint"]], timestamp, ) if row["source_type"] == "migros_receipts": for item in row.get("line_items", []): item_fp = _sha(_canonical([row["source_row_fingerprint"], item["row"], item["name"], item["amount"]])) conn.execute( """INSERT INTO budget_import_line_items( line_item_id,transaction_candidate_id,source_file_label,receipt_key, source_row_or_range,item_name,quantity,is_promotion,amount_original, currency_original,raw_fingerprint,created_at) VALUES (?,?,?,?,?,?,NULL,0,?,'CHF',?,?)""", ("bhhli_" + item_fp[:24], candidate_ids[row["source_row_fingerprint"]], "household:" + row["file_fingerprint"][:12], row.get("receipt_key") or row["source_row_fingerprint"][:24], f"R{item['row']}", item["name"], item["amount"], item_fp, timestamp), ) pair_ids: list[str] = [] row_pair_ids: dict[str, str] = {} for pair in reconstructed["transfer_pairs"]: if pair["pairing_class"] != "safe": continue source, target = rows_by_fp[pair["source_row_fingerprint"]], rows_by_fp[pair["target_row_fingerprint"]] if Decimal(source["signed_amount"]) > 0: source, target = target, source pair_id = "btpair_hh_" + _sha(source["source_row_fingerprint"] + target["source_row_fingerprint"])[:24] evidence = {"matcher_version": PAIRING_VERSION, "pairing_class": "safe", "batch_id": batch_id, "merchant_text_decisive": False} conn.execute( """INSERT INTO budget_transfer_pairs(transfer_pair_id,source_candidate_id,target_candidate_id, source_account_id,target_account_id,source_signed_amount,target_signed_amount,currency, source_booking_date,target_booking_date,source_value_date,target_value_date,status,quality_status, evidence_json,reason_codes_json,budget_effect_chf,created_at,created_by,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,'safe',?,?,'0',?,'household_import',?)""", (pair_id, candidate_ids[source["source_row_fingerprint"]], candidate_ids[target["source_row_fingerprint"]], source["mapping"]["budget_account_id"], target["mapping"]["budget_account_id"], source["signed_amount"], target["signed_amount"], source["currency"], source["transaction_date"], target["transaction_date"], source["transaction_date"], target["transaction_date"], "proposed", _canonical(evidence), _canonical(pair["reason_codes"]), timestamp, timestamp), ) confirm_transfer_pair(conn, pair_id, decision_by="household_import") pair_ids.append(pair_id) row_pair_ids[source["source_row_fingerprint"]] = pair_id row_pair_ids[target["source_row_fingerprint"]] = pair_id link_count = 0 for link in reconstructed["receipt_links"]: if link["status"] != "linked": continue receipt_id = candidate_ids.get(link["receipt_row_fingerprint"]) if not receipt_id: continue money_id = candidate_ids.get(link.get("money_row_fingerprint") or "") or link.get("money_candidate_id") money_transaction_id = link.get("money_transaction_id") if link["status"] == "linked" and bool(money_id) == bool(money_transaction_id): raise RuntimeError("linked Migros receipt must resolve to exactly one money movement") if money_id: conn.execute("UPDATE budget_transaction_candidates SET linked_candidate_id=? WHERE transaction_candidate_id=?", (money_id, receipt_id)) conn.execute( """INSERT INTO household_migros_links(receipt_link_id,batch_id,receipt_candidate_id, money_candidate_id,money_transaction_id,receipt_total,money_total,difference,status,created_at) VALUES (?,?,?,?,?,?,?,?,?,?)""", ("hhmig_" + link["receipt_row_fingerprint"][:24], batch_id, receipt_id, money_id, money_transaction_id, link["receipt_total"], link["money_total"], link["difference"], link["status"], timestamp), ) link_count += 1 for file in reconstructed["files"]: if not file["duplicate"]: > conn.execute( """INSERT INTO household_import_files( household_file_id,batch_id,profile,file_fingerprint,row_count,created_at, period_start,period_end,physical_row_count,logical_row_count) VALUES (?,?,?,?,?,?,?,?,?,?)""", ( "hhfile_" + file["file_fingerprint"][:24], batch_id, file["profile"], file["file_fingerprint"], file["row_count"], timestamp, file.get("period_start"), file.get("period_end"), file.get("physical_row_count"), file.get("logical_row_count"), ), ) E sqlite3.OperationalError: table household_import_files has no column named profile src/jarvis_finance/services/household_import.py:2153: OperationalError =============================== warnings summary =============================== ../FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/testclient.py:1 /home/agent/.hermes/worktrees/FinanceManager-sprint20h1-followup/.venv/lib/python3.11/site-packages/fastapi/testclient.py:1: StarletteDeprecationWarning: Using `httpx` with `starlette.testclient` is deprecated; install `httpx2` instead. from starlette.testclient import TestClient as TestClient # noqa -- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html =========================== short test summary info ============================ FAILED tests/unit/test_sprint20h1_postfinance_upload.py::test_incomplete_period_and_different_end_date_block_confirm_without_writes FAILED tests/unit/test_household_import_v1_golden.py::test_golden_raiffeisen_2026_01_26_swisslos_chf100_is_safe_v3_transfer_and_confirm_is_noop_twice FAILED tests/unit/test_household_import_v1_golden.py::test_explicit_review_batch_confirm_registers_rows_without_promoting_unclear_cashflows FAILED tests/unit/test_household_import_v1_golden.py::test_file_source_row_and_logical_duplicates_are_distinguished FAILED tests/unit/test_household_import_v1_golden.py::test_migros_links_detail_without_double_counting_and_difference_over_cent_reviews FAILED tests/unit/test_household_import_v1_golden.py::test_migros_receipt_replay_preserves_duplicate_disposition_and_writes_nothing FAILED tests/unit/test_household_import_v1_golden.py::test_idempotent_confirm_recomputes_input_identity_before_noop FAILED tests/unit/test_household_import_v1_golden.py::test_selected_mapping_is_bound_and_mismatch_fails_closed FAILED tests/unit/test_household_import_v1_golden.py::test_same_request_duplicate_and_pending_final_have_one_writable_final FAILED tests/unit/test_household_import_v1_golden.py::test_distinct_provider_transactions_with_same_business_fields_are_not_silently_suppressed FAILED tests/unit/test_household_import_v1_golden.py::test_legacy_productive_overlap_is_reviewed_not_counted_again FAILED tests/unit/test_household_import_v1_golden.py::test_household_confirm_preserves_outer_transaction_and_rolls_back_only_its_savepoint FAILED tests/unit/test_household_import_v1_golden.py::test_clear_chf_card_refund_becomes_productive_refund_only_after_confirm FAILED tests/unit/test_household_import_v1_golden.py::test_unique_reversal_deterministically_offsets_original_transaction FAILED tests/unit/test_household_import_v1_golden.py::test_ambiguous_reversal_remains_review_and_never_posts_automatically FAILED tests/unit/test_household_import_v1_golden.py::test_review_expense_preview_confirm_is_bound_atomic_and_idempotent FAILED tests/unit/test_household_import_v1_golden.py::test_migros_receipt_links_to_existing_money_movement_across_batches FAILED tests/unit/test_household_import_v1_golden.py::test_migros_receipt_links_to_productive_transaction_not_consumed_candidate FAILED tests/unit/test_household_import_v1_golden.py::test_authenticated_owner_attestation_api_binds_exact_cluster_and_confirms_neutral_transfer FAILED tests/unit/test_household_import_v1_golden.py::test_review_summary_is_canonical_filtered_and_cursor_paginated 20 failed, 65 passed, 1 warning in 29.14s __HERMES_CWD_8d46a20096ed__/home/agent/.hermes/worktrees/FinanceManager-current-import-performance__HERMES_CWD_8d46a20096ed__