B-39: 数据库持久化与不可变导入
- 版本化 SQLite 迁移(向前/回滚)与连接助手。 - 导入管线:上传文件按 SHA-256 内容哈希不可变保存,文件/行/批次 幂等;重复上传返回 duplicate 并复用已有批次,不产生第二份事实。 - 每条规范化源行反查文件、工作表、原始行号与模板版本。 - 解析失败保留 exception/failed 批次与诊断,不产生已确认交易。 - 决策记录见 docs/decisions/002-persistence.md。
This commit is contained in:
@@ -0,0 +1,252 @@
|
||||
"""Immutable statement import pipeline.
|
||||
|
||||
Every uploaded file is hashed (SHA-256) and written once to content-addressed
|
||||
storage before parsing. A repeated upload of identical bytes never creates a
|
||||
second set of facts: it records a ``duplicate`` batch that points at the
|
||||
original batch. Parse failures keep the batch and its diagnostics as an
|
||||
``exception`` batch without producing any confirmed source rows.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import sqlite3
|
||||
|
||||
from .db import utc_now
|
||||
from .models import StatementBatch
|
||||
from .parser import StatementParseError, parse_statement
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ImportResult:
|
||||
batch_id: int
|
||||
status: str # parsed | duplicate | exception
|
||||
sha256: str
|
||||
source_file_id: int
|
||||
batches: tuple[StatementBatch, ...] = ()
|
||||
message: str | None = None
|
||||
|
||||
|
||||
def import_statement(
|
||||
connection: sqlite3.Connection,
|
||||
storage_dir: str | Path,
|
||||
original_filename: str,
|
||||
content: bytes,
|
||||
company_id: int | None = None,
|
||||
) -> ImportResult:
|
||||
sha256 = hashlib.sha256(content).hexdigest()
|
||||
existing_file = connection.execute(
|
||||
"SELECT id FROM source_files WHERE sha256 = ?", (sha256,)
|
||||
).fetchone()
|
||||
|
||||
if existing_file is not None:
|
||||
return _record_duplicate(connection, existing_file["id"], sha256, company_id)
|
||||
|
||||
stored_path = _store_immutable(Path(storage_dir), original_filename, content, sha256)
|
||||
now = utc_now()
|
||||
with connection:
|
||||
cursor = connection.execute(
|
||||
"""
|
||||
INSERT INTO source_files (sha256, original_filename, size_bytes, storage_path, created_at)
|
||||
VALUES (?, ?, ?, ?, ?)
|
||||
""",
|
||||
(sha256, original_filename, len(content), str(stored_path), now),
|
||||
)
|
||||
source_file_id = cursor.lastrowid
|
||||
cursor = connection.execute(
|
||||
"""
|
||||
INSERT INTO import_batches (source_file_id, status, company_id, created_at, updated_at)
|
||||
VALUES (?, 'parsing', ?, ?, ?)
|
||||
""",
|
||||
(source_file_id, company_id, now, now),
|
||||
)
|
||||
batch_id = cursor.lastrowid
|
||||
|
||||
try:
|
||||
batches = parse_statement(stored_path)
|
||||
except StatementParseError as exc:
|
||||
message = _clean_message(str(exc), stored_path, original_filename)
|
||||
with connection:
|
||||
_insert_exception(connection, batch_id, "parse", message, original_filename)
|
||||
_set_batch_status(connection, batch_id, "exception")
|
||||
return ImportResult(batch_id, "exception", sha256, source_file_id, message=message)
|
||||
except Exception as exc:
|
||||
message = f"文件解析失败,请检查文件是否完整。({type(exc).__name__})"
|
||||
with connection:
|
||||
_insert_exception(connection, batch_id, "internal", message, original_filename)
|
||||
_set_batch_status(connection, batch_id, "failed")
|
||||
raise
|
||||
|
||||
with connection:
|
||||
for batch in batches:
|
||||
_insert_sheet_batch(connection, batch_id, batch)
|
||||
_set_batch_status(connection, batch_id, "parsed")
|
||||
return ImportResult(batch_id, "parsed", sha256, source_file_id, batches=batches)
|
||||
|
||||
|
||||
def _record_duplicate(
|
||||
connection: sqlite3.Connection,
|
||||
source_file_id: int,
|
||||
sha256: str,
|
||||
company_id: int | None = None,
|
||||
) -> ImportResult:
|
||||
original = connection.execute(
|
||||
"""
|
||||
SELECT id FROM import_batches
|
||||
WHERE source_file_id = ? AND status = 'parsed'
|
||||
ORDER BY id LIMIT 1
|
||||
""",
|
||||
(source_file_id,),
|
||||
).fetchone()
|
||||
if original is None:
|
||||
original = connection.execute(
|
||||
"""
|
||||
SELECT id FROM import_batches
|
||||
WHERE source_file_id = ? AND status != 'duplicate'
|
||||
ORDER BY id LIMIT 1
|
||||
""",
|
||||
(source_file_id,),
|
||||
).fetchone()
|
||||
now = utc_now()
|
||||
with connection:
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO import_batches
|
||||
(source_file_id, status, duplicate_of_id, company_id, diagnostics, created_at, updated_at)
|
||||
VALUES (?, 'duplicate', ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
source_file_id,
|
||||
original["id"],
|
||||
company_id,
|
||||
json.dumps({"note": "内容哈希相同,复用已有批次,不产生第二份事实。"}, ensure_ascii=False),
|
||||
now,
|
||||
now,
|
||||
),
|
||||
)
|
||||
return ImportResult(
|
||||
original["id"],
|
||||
"duplicate",
|
||||
sha256,
|
||||
source_file_id,
|
||||
message="相同内容的文件已导入,本次按重复上传处理。",
|
||||
)
|
||||
|
||||
|
||||
def _store_immutable(
|
||||
storage_dir: Path, original_filename: str, content: bytes, sha256: str
|
||||
) -> Path:
|
||||
suffix = Path(original_filename).suffix.lower()
|
||||
target_dir = storage_dir / sha256[:2]
|
||||
target_dir.mkdir(parents=True, exist_ok=True)
|
||||
target = target_dir / f"{sha256}{suffix}"
|
||||
# Publish via a temporary file + hard link: the content-addressed target
|
||||
# either appears complete or not at all, and is never overwritten.
|
||||
temp = target_dir / f".{sha256}.tmp"
|
||||
temp.write_bytes(content)
|
||||
try:
|
||||
os.link(temp, target)
|
||||
except FileExistsError:
|
||||
# Content-addressed name means identical bytes; never overwrite.
|
||||
pass
|
||||
finally:
|
||||
temp.unlink(missing_ok=True)
|
||||
return target
|
||||
|
||||
|
||||
def _insert_sheet_batch(
|
||||
connection: sqlite3.Connection, batch_id: int, batch: StatementBatch
|
||||
) -> None:
|
||||
now = utc_now()
|
||||
cursor = connection.execute(
|
||||
"""
|
||||
INSERT INTO sheet_batches (
|
||||
import_batch_id, sheet_name, bank_name, template_id, template_version,
|
||||
header_row, own_account, own_name, period_start, period_end,
|
||||
transaction_count, warnings, created_at
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
batch_id,
|
||||
batch.sheet_name,
|
||||
batch.bank_name,
|
||||
batch.template_id,
|
||||
batch.template_version,
|
||||
batch.header_row,
|
||||
batch.own_account,
|
||||
batch.own_name,
|
||||
batch.period_start.isoformat() if batch.period_start else None,
|
||||
batch.period_end.isoformat() if batch.period_end else None,
|
||||
len(batch.transactions),
|
||||
json.dumps(list(batch.warnings), ensure_ascii=False),
|
||||
now,
|
||||
),
|
||||
)
|
||||
sheet_batch_id = cursor.lastrowid
|
||||
for transaction in batch.transactions:
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO source_rows (
|
||||
sheet_batch_id, source_row, transaction_at, income, expense, balance,
|
||||
own_account, own_name, counterparty_account, counterparty_name,
|
||||
counterparty_bank, summary, purpose, reference, currency, created_at
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
sheet_batch_id,
|
||||
transaction.source_row,
|
||||
transaction.transaction_at.isoformat(),
|
||||
str(transaction.income),
|
||||
str(transaction.expense),
|
||||
str(transaction.balance) if transaction.balance is not None else None,
|
||||
transaction.own_account,
|
||||
transaction.own_name,
|
||||
transaction.counterparty_account,
|
||||
transaction.counterparty_name,
|
||||
transaction.counterparty_bank,
|
||||
transaction.summary,
|
||||
transaction.purpose,
|
||||
transaction.reference,
|
||||
transaction.currency,
|
||||
now,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _insert_exception(
|
||||
connection: sqlite3.Connection,
|
||||
batch_id: int,
|
||||
stage: str,
|
||||
message: str,
|
||||
original_filename: str,
|
||||
) -> None:
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO import_exceptions (import_batch_id, stage, message, diagnostics, created_at)
|
||||
VALUES (?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
batch_id,
|
||||
stage,
|
||||
message,
|
||||
json.dumps({"original_filename": original_filename}, ensure_ascii=False),
|
||||
utc_now(),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _set_batch_status(connection: sqlite3.Connection, batch_id: int, status: str) -> None:
|
||||
connection.execute(
|
||||
"UPDATE import_batches SET status = ?, updated_at = ? WHERE id = ?",
|
||||
(status, utc_now(), batch_id),
|
||||
)
|
||||
|
||||
|
||||
def _clean_message(message: str, stored_path: Path, original_filename: str) -> str:
|
||||
return message.replace(str(stored_path), original_filename).replace(
|
||||
stored_path.name, original_filename
|
||||
)
|
||||
Reference in New Issue
Block a user