feat(msp): Parser-Framework + DocumentBuilder + Run-Orchestrator

Implementiert das eigentliche Import-Framework im Paket msp.importers/:

- base.py: abstrakte BaseSupplierParser-Klasse + CanonicalRow-Dataclass
  (lieferantenneutrale Repräsentation einer Rechnungsposition) sowie
  Helfer für deutsche Dezimal-/Datums-Parser und Enum-Normalisierung
  (Monthly/Annual/Triennial, P1M/P1Y/P3Y, New/Renewal/Billing/Upgrade/
  Cancel).
- registry.py: Parser-Registry, über parser_key aufgelöst. Neue
  Lieferanten bringen einen eigenen Parser mit und registrieren sich
  per @register-Dekorator.
- adn_monthly_csv.py: Parser für ADNs monatliche Rechnungs-CSV.
  Extrahiert alle relevanten Felder strukturiert, inklusive
  SUBSCRIPTION_ID_EXTERNAL, BILLINGPLAN, VERTRAGSDAUER, WARTUNGSBEGINN/
  -ENDE und BUCHUNGSTYP — Daten, die der alte adnconnect-Importer
  verworfen oder nur in HTML-Descriptions vergraben hatte.
- resolution.py: Lookups für Customer + Item und die Output-Mode-
  Selektion (Run-Override → Customer.billing_mode → Profil-Default).
- subscriptions.py: Upsert-Logik für Supply Subscription + Event,
  inklusive Deduplizierung per (subscription_external_id, period).
- builder.py: DocumentBuilder, der sowohl Sales Invoices als auch
  Delivery Notes aus kanonischen Zeilen erzeugt. Rechnungspreise werden
  per unit_price × qty gesetzt (mit Rabatt in description) — korrekt
  gemäß ADN-POSITIONSPREIS statt adnconnects fehlerhafter Listenpreis/
  discount_percentage-Kombination.
- run_orchestrator.py: verdrahtet das Ganze und wird aus dem
  Supplier-Import-Run-Controller über die Whitelisted-Methode do_import
  aufgerufen.
This commit is contained in:
David Malinowski
2026-04-14 12:41:37 +02:00
parent b399b6f37d
commit 687a1faac0
8 changed files with 1060 additions and 0 deletions
View File
+163
View File
@@ -0,0 +1,163 @@
"""Parser für die monatlichen CSV-Abrechnungsdateien der ADN Distribution.
Format: Semikolon-getrennt, Zeichenkodierung UTF-8 (BOM erlaubt), deutsche
Dezimalschreibweise. Jede Zeile ist **eine** Rechnungsposition; Kopfdaten
(RECHNUNG, DATUM, …) sind pro Zeile redundant wiederholt.
Registrierung über :func:`msp.importers.registry.register`.
"""
from __future__ import annotations
import csv
import io
from typing import Iterable
from msp.importers.base import (
BaseSupplierParser,
CanonicalRow,
ParseResult,
normalize_billing_plan,
normalize_booking_type,
normalize_term_duration,
parse_german_date,
parse_german_decimal,
)
from msp.importers.registry import register
EXPECTED_HEADER = [
"RECHNUNG", "DATUM", "KUNDENNR", "DEBITORKONTO", "SACHBEARBEITER",
"LIEFERSCHEIN", "LIEFERSCHEINDATUM", "USTID",
"RE_FIRMA", "RE_ADRESSE", "RE_PLZ", "RE_ORT", "RE_LAND",
"LI_FIRMA", "LI_STRASSE", "LI_PLZ", "LI_ORT", "LI_LAND",
"HA_FIRMA", "HA_STRASSE", "HA_PLZ", "HA_ORT", "HA_LAND",
"WARENWERT", "MWST", "GESAMTBETRAG", "NETTOZAHLBARBIS",
"LIEFERBEDINGUNGEN", "ZAHLUNGSBEDINGUNG",
"ENDKUNDE", "POSITION", "ARTIKEL", "HERSTELLERNUMMER", "ARTIKELBEZ",
"MENGE", "PREISME", "LISTPREIS", "RABATT", "EINZELPREIS", "POSITIONSPREIS",
"WARTUNGSBEGINN", "WARTUNGSENDE", "VERTRAG",
"MARKETPLACE_REF", "ORDER_REFERENCE", "ENDKUNDE_REFERENCE",
"SUBSCRIPTION_ID_EXTERNAL", "SUBSCRIPTION_START_DATE", "BUCHUNGSTYP",
"ORDERDATUM", "RECHNUNGSART", "MSERP", "MSERP_BILLINGPERIOD",
"BILLINGPLAN", "VERTRAGSDAUER", "ADDITIONALID", "ADDITIONALDESCRIPTION",
]
def _strip_bom(value: str) -> str:
return value.lstrip("\ufeff") if value else value
def _normalize_header(header: Iterable[str]) -> list[str]:
return [_strip_bom(h).strip() for h in header]
def _customer_ref_from(raw_ref: str | None, endkunde_field: str | None) -> str | None:
"""Extrahiert die ERP-Customer-Referenz.
In älteren ADN-Feeds steht die Referenz in Klammern hinter dem Endkundennamen
(``"Augenweide Optometrie [CUST-21385]"``); in neueren Versionen direkt in
der Spalte ``ENDKUNDE_REFERENCE``.
"""
if raw_ref:
return raw_ref.strip()
if endkunde_field and "[" in endkunde_field and "]" in endkunde_field:
inside = endkunde_field.split("[", 1)[1].split("]", 1)[0].strip()
return inside or None
return None
@register
class ADNMonthlyCSVParser(BaseSupplierParser):
parser_key = "adn_monthly_csv_v1"
description = "ADN Distribution — monatliche Rechnungs-CSV"
def parse(self, file_path: str) -> ParseResult:
result = ParseResult()
with open(file_path, "r", encoding="utf-8", newline="") as fh:
reader = csv.reader(fh, delimiter=";", quotechar='"')
try:
raw_header = next(reader)
except StopIteration:
result.warnings.append("Datei ist leer.")
return result
header = _normalize_header(raw_header)
result.format_drift = not self._header_matches(header)
for source_row, row in enumerate(reader, start=2):
if not row or not any(cell.strip() for cell in row):
continue
data = self._row_to_dict(header, row)
position = (data.get("POSITION") or "").strip()
# Reine Kopfzeile (Position = "1" ohne Artikel) wird übersprungen;
# ADN wiederholt die Rechnungskopfdaten als erste Zeile ohne Artikel.
if not (data.get("ARTIKEL") or "").strip() and not (data.get("HERSTELLERNUMMER") or "").strip():
continue
canonical = self._to_canonical(data, source_row)
if canonical is not None:
result.rows.append(canonical)
# Rechnungsdatum aus erster Zeile übernehmen (alle Zeilen gehören zu
# derselben Lieferantenrechnung und tragen dasselbe Datum)
if result.rows:
first_raw = result.rows[0].raw
result.invoice_date = parse_german_date(first_raw.get("DATUM"))
return result
# ---------------------------------------------------------------
# internals
# ---------------------------------------------------------------
def _header_matches(self, actual: list[str]) -> bool:
return [h.upper() for h in actual if h] == [h.upper() for h in EXPECTED_HEADER if h]
def _row_to_dict(self, header: list[str], row: list[str]) -> dict[str, str]:
# Robust gegen Zeilen mit abweichender Spaltenanzahl
return {header[i]: row[i] for i in range(min(len(header), len(row)))}
def _to_canonical(self, d: dict[str, str], source_row: int) -> CanonicalRow | None:
raw_ref = (d.get("ENDKUNDE_REFERENCE") or "").strip() or None
customer_ref = _customer_ref_from(raw_ref, d.get("ENDKUNDE"))
# Negative Menge oder negative Positionspreise kennzeichnen Gutschriften
qty = parse_german_decimal(d.get("MENGE")) or 0.0
list_price = parse_german_decimal(d.get("LISTPREIS"))
unit_price = parse_german_decimal(d.get("EINZELPREIS"))
amount = parse_german_decimal(d.get("POSITIONSPREIS"))
discount = parse_german_decimal(d.get("RABATT"))
rechnungsart = (d.get("RECHNUNGSART") or "").strip()
document_type = "Gutschrift" if "gutschrift" in rechnungsart.lower() else "Rechnung"
if amount is not None and amount < 0:
document_type = "Gutschrift"
return CanonicalRow(
source_row=source_row,
supplier_invoice_no=(d.get("RECHNUNG") or "").strip() or None,
customer_external_ref=customer_ref,
customer_display_name=(d.get("ENDKUNDE") or "").strip() or None,
vendor_product_id=(d.get("HERSTELLERNUMMER") or "").strip() or None,
supplier_article_no=(d.get("ARTIKEL") or "").strip() or None,
description=(d.get("ARTIKELBEZ") or "").strip() or None,
qty=qty,
list_price=list_price,
discount_pct=discount,
unit_price=unit_price,
amount=amount,
period_start=parse_german_date(d.get("WARTUNGSBEGINN")),
period_end=parse_german_date(d.get("WARTUNGSENDE")),
subscription_external_id=(d.get("SUBSCRIPTION_ID_EXTERNAL") or "").strip() or None,
subscription_start=parse_german_date(d.get("SUBSCRIPTION_START_DATE")),
billing_plan=normalize_billing_plan(d.get("BILLINGPLAN")),
term_duration=normalize_term_duration(d.get("VERTRAGSDAUER")),
booking_type=normalize_booking_type(d.get("BUCHUNGSTYP")),
document_type=document_type,
marketplace_ref=(d.get("MARKETPLACE_REF") or "").strip() or None,
order_ref=(d.get("ORDER_REFERENCE") or "").strip() or None,
raw=d,
)
+191
View File
@@ -0,0 +1,191 @@
"""Basis-Klassen und Datenstrukturen für das Supplier-Import-Framework.
Parser dürfen **keine** Frappe-/DB-Zugriffe tun. Sie lesen eine Datei und geben
eine Sequenz von :class:`CanonicalRow`-Objekten zurück. Validierung und Persistenz
übernimmt der DocumentBuilder / Orchestrator.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import date
from typing import Iterable, Iterator
# ---------------------------------------------------------------------------
# Kanonische Datenstruktur — gilt für *alle* Lieferanten nach Normalisierung
# ---------------------------------------------------------------------------
@dataclass
class CanonicalRow:
# Identifikation
source_row: int = 0 # Zeilennummer in der Quelldatei
supplier_invoice_no: str | None = None # z. B. ADN RECHNUNG
# Kunde
customer_external_ref: str | None = None # z. B. "CUST-21254"
customer_display_name: str | None = None # Langname aus der Quelldatei
# Produkt
vendor_product_id: str | None = None # z. B. "CFQ7TTC0LCHC:0002"
supplier_article_no: str | None = None # interne Artikelnummer des Lieferanten
description: str | None = None
# Mengen & Preise
qty: float = 0.0
list_price: float | None = None
discount_pct: float | None = None
unit_price: float | None = None
amount: float | None = None
# Laufzeit & Billing
period_start: date | None = None
period_end: date | None = None
subscription_external_id: str | None = None
subscription_start: date | None = None
billing_plan: str | None = None # normalisiert: "Monthly" / "Annual" / "Triennial"
term_duration: str | None = None # normalisiert: "P1M" / "P1Y" / "P3Y"
booking_type: str | None = None # normalisiert: "New" / "Renewal" / "Billing" / "Upgrade" / "Cancel"
document_type: str | None = None # "Rechnung" / "Gutschrift"
# Referenzen
marketplace_ref: str | None = None
order_ref: str | None = None
# Forensik
raw: dict = field(default_factory=dict)
# ---------------------------------------------------------------------------
# Parse-Ergebnis inkl. Kopfdaten & Warnungen
# ---------------------------------------------------------------------------
@dataclass
class ParseResult:
rows: list[CanonicalRow] = field(default_factory=list)
warnings: list[str] = field(default_factory=list)
format_drift: bool = False # Spalten wichen vom erwarteten Layout ab
invoice_date: date | None = None # falls alle Zeilen aus einer Lieferantenrechnung stammen
# ---------------------------------------------------------------------------
# Basis-Parser
# ---------------------------------------------------------------------------
class BaseSupplierParser:
"""Alle Lieferanten-Parser erben hiervon und implementieren :meth:`parse`.
Kein DB-Zugriff — reine Datei-in → Objekte-out-Logik.
"""
#: Eindeutiger Schlüssel in der Registry (z. B. ``adn_monthly_csv_v1``).
parser_key: str = ""
#: Für Menschen lesbare Beschreibung.
description: str = ""
def parse(self, file_path: str) -> ParseResult:
raise NotImplementedError
# ---------------------------------------------------------------------------
# Hilfsfunktionen zur Wert-Normalisierung
# ---------------------------------------------------------------------------
def normalize_billing_plan(raw: str | None) -> str | None:
if not raw:
return None
r = raw.strip().lower()
if r in ("month", "monthly", "monatlich"):
return "Monthly"
if r in ("year", "annual", "yearly", "jährlich", "jaehrlich"):
return "Annual"
if r in ("triennial", "3years", "3y"):
return "Triennial"
return None
def normalize_term_duration(raw: str | None) -> str | None:
if not raw:
return None
r = raw.strip().lower()
if r in ("monthly", "month", "p1m", "1m"):
return "P1M"
if r in ("annual", "yearly", "p1y", "1y"):
return "P1Y"
if r in ("triennial", "p3y", "3y"):
return "P3Y"
return None
def normalize_booking_type(raw: str | None) -> str | None:
if not raw:
return None
r = raw.strip().lower()
mapping = {
"new": "New",
"neu": "New",
"renewal": "Renewal",
"verlängerung": "Renewal",
"billing": "Billing",
"abrechnung": "Billing",
"upgrade": "Upgrade",
"cancel": "Cancel",
"cancellation": "Cancel",
"kündigung": "Cancel",
}
return mapping.get(r)
def parse_german_decimal(value: str | None) -> float | None:
if value is None:
return None
s = str(value).strip()
if not s:
return None
# deutsches Format: Komma als Dezimaltrenner, evtl. Punkt als Tausendertrenner
s = s.replace(".", "").replace(",", ".") if s.count(",") == 1 and s.count(".") <= 1 else s.replace(",", ".")
try:
return float(s)
except ValueError:
return None
def parse_german_date(value: str | None) -> date | None:
if not value:
return None
s = str(value).strip()
if not s or s in ("0", "-"):
return None
# Verbreitete ADN-Formate: TT.MM.JJJJ oder JJJJ-MM-TT
from datetime import datetime
for fmt in ("%d.%m.%Y", "%Y-%m-%d", "%d/%m/%Y"):
try:
return datetime.strptime(s, fmt).date()
except ValueError:
continue
# Häufiger ADN-Fehler: "226.02.2025" → führender Tippfehler, erste Ziffer droppen
if len(s) == 11 and s[0].isdigit() and s[1].isdigit() and s[2] == ".":
try:
return datetime.strptime(s[1:], "%d.%m.%Y").date()
except ValueError:
pass
return None
def iter_parser_output(parse_result: ParseResult) -> Iterator[CanonicalRow]:
"""Kleine Convenience: Iteration nur über Rows."""
yield from parse_result.rows
__all__ = [
"CanonicalRow",
"ParseResult",
"BaseSupplierParser",
"normalize_billing_plan",
"normalize_term_duration",
"normalize_booking_type",
"parse_german_decimal",
"parse_german_date",
"iter_parser_output",
]
+273
View File
@@ -0,0 +1,273 @@
"""DocumentBuilder — erzeugt Sales Invoices oder Delivery Notes aus kanonischen Zeilen."""
from __future__ import annotations
from dataclasses import dataclass, field
import frappe
from msp.importers.base import CanonicalRow
from msp.importers.resolution import resolve_customer, resolve_item
from msp.importers.subscriptions import (
create_event,
event_exists_for_period,
upsert_supply_subscription,
)
@dataclass
class BuildResult:
created_documents: list[tuple[str, str]] = field(default_factory=list) # (doctype, name)
line_outcomes: list[dict] = field(default_factory=list)
errors: list[str] = field(default_factory=list)
class DocumentBuilder:
def __init__(self, profile):
self.profile = profile
self.output_mode: str = "Sales Invoice" # wird pro Gruppe gesetzt
# -----------------------------------------------------------------
# Public API
# -----------------------------------------------------------------
def build(self, rows: list[CanonicalRow], output_mode: str) -> BuildResult:
"""Erzeugt Dokumente für *eine* Kunden-Gruppe (alle rows gehören demselben Kunden)."""
self.output_mode = output_mode
result = BuildResult()
if not rows:
return result
# Gutschriften separieren (negative Positionspreise oder Rechnungsart=Gutschrift)
invoice_rows = [r for r in rows if (r.amount or 0) >= 0 and r.document_type != "Gutschrift"]
credit_rows = [r for r in rows if (r.amount or 0) < 0 or r.document_type == "Gutschrift"]
if invoice_rows:
self._build_one(invoice_rows, is_return=False, result=result)
if credit_rows:
self._build_one(credit_rows, is_return=True, result=result)
return result
# -----------------------------------------------------------------
# Internals
# -----------------------------------------------------------------
def _build_one(self, rows: list[CanonicalRow], *, is_return: bool, result: BuildResult) -> None:
first = rows[0]
customer, is_fallback = resolve_customer(
first.customer_external_ref, self.profile.fallback_customer,
)
if not customer:
for r in rows:
result.line_outcomes.append({
"row": r, "status": "Error",
"error": f"Customer nicht gefunden (ref={r.customer_external_ref})",
})
return
doc = self._new_document(customer, first, is_return=is_return)
outcomes_this_group: list[dict] = []
for r in rows:
outcome = self._append_line(doc, r, is_return=is_return)
outcomes_this_group.append(outcome)
result.line_outcomes.append(outcome)
if not doc.items:
# Kein Dokument erzeugen. Unterscheiden: all skipped vs. all errored.
statuses = {o["status"] for o in outcomes_this_group}
if statuses == {"Skipped"}:
# Idempotent: alles schon da. Kein Fehler.
return
# Sonst waren echte Fehler im Spiel.
result.errors.append(f"Kunde {customer}: keine gültigen Positionen, Dokument nicht erzeugt.")
return
try:
doc.insert(ignore_permissions=True)
except Exception as e:
result.errors.append(f"Insert fehlgeschlagen ({customer}): {e}")
# Aufräumen: bereits gesetzte Outcomes als Error markieren
for o in result.line_outcomes[-len(rows):]:
if o["status"] == "Created":
o["status"] = "Error"
o["error"] = str(e)
return
result.created_documents.append((doc.doctype, doc.name))
# Subscription-Events erst nach erfolgreichem Insert, damit wir die
# echten Item-Row-IDs nutzen können.
self._post_create_subscriptions(doc, rows, result)
def _new_document(self, customer: str, first: CanonicalRow, *, is_return: bool):
"""Erzeugt das Grundgerüst (ohne Zeilen) der Sales Invoice / Delivery Note."""
profile = self.profile
period = self._format_period_label(first)
title_prefix = profile.title_prefix_credit_note if is_return else profile.title_prefix
title = f"{(title_prefix or '').strip()} {period}".strip()
if self.output_mode == "Sales Invoice":
doc = frappe.new_doc("Sales Invoice")
doc.customer = customer
doc.company = profile.default_company
doc.currency = frappe.db.get_value("Company", profile.default_company, "default_currency") or "EUR"
doc.posting_date = first.raw.get("DATUM") or frappe.utils.today()
doc.posting_date = self._coerce_german_date(doc.posting_date)
doc.set_posting_time = 1
if profile.default_cost_center:
doc.cost_center = profile.default_cost_center
if profile.default_taxes_and_charges:
doc.taxes_and_charges = profile.default_taxes_and_charges
if profile.default_payment_terms_template:
doc.payment_terms_template = profile.default_payment_terms_template
if profile.default_tc_name:
doc.tc_name = profile.default_tc_name
if profile.intro_text_invoice and not is_return:
doc.introduction_text = profile.intro_text_invoice
if profile.intro_text_credit_note and is_return:
doc.introduction_text = profile.intro_text_credit_note
if title:
doc.title = title
if is_return:
doc.is_return = 1
return doc
# Delivery Note
doc = frappe.new_doc("Delivery Note")
doc.customer = customer
doc.company = profile.default_company
doc.posting_date = self._coerce_german_date(first.raw.get("DATUM") or frappe.utils.today())
doc.set_posting_time = 1
if profile.default_cost_center:
doc.cost_center = profile.default_cost_center
if title:
doc.title = title
if is_return:
doc.is_return = 1
return doc
def _append_line(self, doc, row: CanonicalRow, *, is_return: bool) -> dict:
item_code, reason = resolve_item(row.vendor_product_id)
if not item_code:
return {"row": row, "status": "Error",
"error": f"Item-Lookup {reason}: vendor_product_id={row.vendor_product_id}"}
# Deduplizierung: bereits Event für diese Subscription+Period vorhanden?
if row.subscription_external_id and event_exists_for_period(
row.subscription_external_id, row.period_start, row.period_end):
return {"row": row, "status": "Skipped",
"error": "Bereits importiert (Subscription+Period existieren)"}
qty = row.qty or 0
if is_return and qty > 0:
qty = -qty
# ERPNext berechnet amount = qty * rate. `rate` muss daher der
# tatsächliche Netto-Einzelpreis nach Rabatt sein (EINZELPREIS aus ADN).
# Fallback-Reihenfolge: unit_price → amount/qty → list_price
if row.unit_price is not None:
rate = row.unit_price
elif row.amount is not None and qty:
rate = row.amount / qty
elif row.list_price is not None:
rate = row.list_price
else:
rate = 0
item_row = {
"item_code": item_code,
"qty": qty,
"rate": rate,
"description": self._format_description(row),
}
# Listenpreis + Rabatt zur Nachvollziehbarkeit in der UI mitführen
if row.list_price is not None:
item_row["price_list_rate"] = row.list_price
if row.discount_pct:
item_row["discount_percentage"] = row.discount_pct
item = doc.append("items", item_row)
# Warehouse für Delivery Note erforderlich
if self.output_mode == "Delivery Note" and self.profile.default_warehouse:
item.warehouse = self.profile.default_warehouse
return {"row": row, "status": "Created", "error": None}
def _post_create_subscriptions(self, doc, rows: list[CanonicalRow], result: BuildResult) -> None:
"""Erzeugt pro successful row ein Supply Subscription Event (sofern eine Subscription-ID vorhanden ist)."""
# Wir gehen davon aus, dass doc.items in gleicher Reihenfolge wie rows angelegt wurde,
# abzüglich der error/skipped-Zeilen.
created_rows = [o for o in result.line_outcomes[-len(rows):] if o["status"] == "Created"]
if len(created_rows) != len(doc.items):
# Fallback: keine 1:1 Zuordnung — überspringen, wird später nachgezogen
return
for outcome, item_row in zip(created_rows, doc.items):
r: CanonicalRow = outcome["row"]
if not r.subscription_external_id:
continue
vendor = self.profile.vendor
supplier = self.profile.supplier
sub_name = upsert_supply_subscription(
r, customer=doc.customer, supplier=supplier, vendor=vendor, item=item_row.item_code,
)
if not sub_name:
continue
try:
event_name = create_event(
r,
supply_subscription=sub_name,
target_doctype=doc.doctype,
target_name=doc.name,
target_row=item_row.name,
)
# Link auf dem Item-Row setzen
frappe.db.set_value(
f"{doc.doctype} Item", item_row.name,
"supply_subscription", sub_name,
)
outcome["event"] = event_name
outcome["supply_subscription"] = sub_name
except Exception as e:
outcome["subscription_error"] = str(e)
# -----------------------------------------------------------------
# Formatter / Helpers
# -----------------------------------------------------------------
def _format_description(self, row: CanonicalRow) -> str:
parts = []
if row.description:
parts.append(row.description)
if row.period_start and row.period_end:
parts.append(f"Zeitraum von {row.period_start.strftime('%d.%m.%Y')} bis {row.period_end.strftime('%d.%m.%Y')}")
if row.subscription_external_id:
parts.append(f"Subscription: {row.subscription_external_id}")
return "<br>".join(parts) if parts else (row.description or "")
def _format_period_label(self, row: CanonicalRow) -> str:
"""Monats-/Jahreslabel für den Dokumenttitel — z. B. '02.2026'."""
d = row.period_end or row.period_start
if d:
return d.strftime("%m.%Y")
from datetime import date as _d
from msp.importers.base import parse_german_date
invdate = parse_german_date(row.raw.get("DATUM")) if row.raw else None
if invdate:
return invdate.strftime("%m.%Y")
return _d.today().strftime("%m.%Y")
@staticmethod
def _coerce_german_date(value):
"""Akzeptiert strings oder date-Objekte; liefert date-tauglichen Wert für Frappe."""
if not value:
from datetime import date
return date.today()
if hasattr(value, "strftime"):
return value
from msp.importers.base import parse_german_date
return parse_german_date(value) or frappe.utils.getdate(value)
+46
View File
@@ -0,0 +1,46 @@
"""Registry der verfügbaren Parser.
Parser werden über ihren ``parser_key`` aus einem Supplier Import Profile referenziert.
Neue Lieferanten ⇒ neuer Parser + Eintrag hier.
"""
from __future__ import annotations
from typing import Type
from msp.importers.base import BaseSupplierParser
_REGISTRY: dict[str, Type[BaseSupplierParser]] = {}
def register(parser_cls: Type[BaseSupplierParser]) -> Type[BaseSupplierParser]:
"""Dekorator zum Registrieren einer Parser-Klasse anhand ihres ``parser_key``."""
key = getattr(parser_cls, "parser_key", "")
if not key:
raise ValueError(f"{parser_cls.__name__} ohne parser_key kann nicht registriert werden")
if key in _REGISTRY:
raise ValueError(f"Parser-Key bereits belegt: {key}")
_REGISTRY[key] = parser_cls
return parser_cls
def resolve(parser_key: str) -> BaseSupplierParser:
"""Gibt eine Instanz des Parsers zum gegebenen Schlüssel zurück."""
if parser_key not in _REGISTRY:
_autoload()
if parser_key not in _REGISTRY:
raise KeyError(f"Kein Parser für Key '{parser_key}' registriert. "
f"Verfügbar: {sorted(_REGISTRY)}")
return _REGISTRY[parser_key]()
def available_keys() -> list[str]:
_autoload()
return sorted(_REGISTRY)
def _autoload() -> None:
"""Importiert bekannte Parser-Module, damit ihre @register-Decorators laufen."""
# Noqa — Imports dienen nur der Registrierung
from msp.importers import adn_monthly_csv # noqa: F401
+52
View File
@@ -0,0 +1,52 @@
"""ERP-Lookups: Customer + Item aus kanonischen Feldern auflösen."""
from __future__ import annotations
import frappe
def resolve_customer(external_ref: str | None, fallback: str | None = None) -> tuple[str | None, bool]:
"""Gibt (customer_name, is_fallback) zurück. None wenn weder Match noch Fallback greift."""
if external_ref:
exists = frappe.db.exists("Customer", external_ref)
if exists:
return external_ref, False
# Fallback-Match auf customer_name, falls ID nicht existiert
match = frappe.db.get_value("Customer", {"customer_name": external_ref}, "name")
if match:
return match, False
if fallback and frappe.db.exists("Customer", fallback):
return fallback, True
return None, False
def resolve_item(vendor_product_id: str | None) -> tuple[str | None, str]:
"""Gibt (item_code, reason) zurück. Reason bei Nichtauflösbarkeit: 'no-match' oder 'ambiguous'."""
if not vendor_product_id:
return None, "no-vendor-product-id"
matches = frappe.db.get_all(
"Item",
filters={"hersteller_artikel_nummer": vendor_product_id},
fields=["name"],
limit=2,
)
if not matches:
return None, "no-match"
if len(matches) > 1:
return None, "ambiguous"
return matches[0]["name"], "ok"
def determine_effective_output_mode(run_mode: str | None, profile_default: str,
customer_billing_mode: str | None) -> str:
"""Reihenfolge: expliziter Run-Override → Customer.billing_mode → Profil-Default.
Customer.billing_mode-Werte (aus MSP-Custom-Field):
ASAP, Collective Bill, per Item Group, per Sales Order, remaining per Item Group
Alles außer "ASAP" ⇒ Delivery Note (wird später gesammelt verrechnet).
"""
if run_mode:
return run_mode
if customer_billing_mode and customer_billing_mode != "ASAP":
return "Delivery Note"
return profile_default or "Sales Invoice"
+233
View File
@@ -0,0 +1,233 @@
"""Orchestriert einen Supplier Import Run.
Ablauf:
1. Parser per ``profile.parser_key`` auflösen
2. Datei lesen und via Parser in CanonicalRow-Objekte zerlegen
3. Zeilen nach Kunde gruppieren
4. Output-Modus je Gruppe bestimmen (Run-Override → Customer.billing_mode → Profil)
5. DocumentBuilder aufrufen (Sales Invoice / Delivery Note)
6. Supplier Import Line Records schreiben + Statistiken/Log füllen
"""
from __future__ import annotations
import json
import os
import traceback
from collections import defaultdict
from typing import Iterable
import frappe
from frappe.utils import now_datetime
from msp.importers.base import CanonicalRow, ParseResult
from msp.importers.builder import BuildResult, DocumentBuilder
from msp.importers.registry import resolve
from msp.importers.resolution import determine_effective_output_mode
def execute_run(run) -> dict:
"""Einstiegspunkt. `run` ist ein :class:`Supplier Import Run`-Dokument."""
log_buffer: list[str] = []
def log(msg: str) -> None:
log_buffer.append(msg)
_set_status(run, status="Running", started_at=now_datetime())
try:
profile = frappe.get_doc("Supplier Import Profile", run.profile)
if not profile.enabled:
raise Exception(f"Profile '{profile.profile_name}' ist deaktiviert.")
parser = resolve(profile.parser_key)
file_path = _resolve_file_path(run.file)
log(f"Parser: {profile.parser_key} · Datei: {file_path}")
parse_result: ParseResult = parser.parse(file_path)
log(f"Parsed: {len(parse_result.rows)} Zeilen "
f"(format_drift={parse_result.format_drift})")
if parse_result.format_drift and not run.acknowledge_format_drift:
raise Exception(
"Format der Quelldatei weicht vom erwarteten Layout ab. "
"Mit 'Abweichendes Format akzeptieren' übersteuern."
)
grouped = _group_rows(parse_result.rows)
log(f"Gruppen (Kunden): {len(grouped)}")
document_count = 0
line_outcomes: list[dict] = []
errors: list[str] = []
builder = DocumentBuilder(profile)
for customer_ref, group_rows in grouped.items():
customer_billing_mode = _customer_billing_mode(customer_ref, profile.fallback_customer)
effective_mode = determine_effective_output_mode(
run.output_mode, profile.default_output_mode, customer_billing_mode,
)
log(f" {customer_ref}: {len(group_rows)} Zeilen → {effective_mode}")
try:
group_result: BuildResult = builder.build(group_rows, output_mode=effective_mode)
document_count += len(group_result.created_documents)
errors.extend(group_result.errors)
line_outcomes.extend(group_result.line_outcomes)
for dt, name in group_result.created_documents:
log(f" created {dt} {name}")
except Exception as e:
tb = traceback.format_exc(limit=3)
errors.append(f"Gruppe {customer_ref}: {e}")
log(f" ERROR: {e}\n{tb}")
# Zeilenresultate in Supplier Import Line child table schreiben
_persist_lines(run, line_outcomes)
error_count = sum(1 for o in line_outcomes if o.get("status") == "Error")
skipped_count = sum(1 for o in line_outcomes if o.get("status") == "Skipped")
if errors or error_count:
status = "Partial" if document_count else "Failed"
else:
# "Succeeded" gilt auch, wenn alle Zeilen rechtmäßig übersprungen wurden
# (Idempotenz bei erneutem Import derselben Datei)
status = "Succeeded"
log(f"Ergebnis: documents={document_count} lines={len(line_outcomes)} "
f"errors={error_count} skipped={skipped_count} → status={status}")
# Persistenter Statusschreib: direkt per DB (umgeht Hook-/Timestamp-Kollisionen
# nach dem set-lines-Save, der direkt davor passiert ist).
frappe.db.set_value("Supplier Import Run", run.name, {
"status": status,
"line_count": len(line_outcomes),
"document_count": document_count,
"error_count": error_count,
"finished_at": now_datetime(),
"log": "\n".join(log_buffer),
}, update_modified=False)
frappe.db.commit()
return {
"status": status,
"document_count": document_count,
"line_count": len(line_outcomes),
"error_count": error_count,
"skipped_count": skipped_count,
}
except Exception as e:
log(f"FATAL: {e}\n{traceback.format_exc(limit=5)}")
run.reload()
run.status = "Failed"
run.finished_at = now_datetime()
run.log = "\n".join(log_buffer)
run.save(ignore_permissions=True)
frappe.db.commit()
raise
# ---------------------------------------------------------------------------
# helpers
# ---------------------------------------------------------------------------
def _set_status(run, *, status: str, started_at=None) -> None:
run.reload()
run.status = status
if started_at is not None:
run.started_at = started_at
run.save(ignore_permissions=True)
frappe.db.commit()
def _resolve_file_path(attached_file: str) -> str:
"""Konvertiert den /private/files/... bzw. /files/... Pfad in den absoluten Dateisystempfad."""
if not attached_file:
raise Exception("Keine Datei angehängt.")
# Absoluter Pfad, der bereits existiert → direkt nutzen
if os.path.isabs(attached_file) and os.path.exists(attached_file):
return attached_file
site_path = os.path.abspath(frappe.get_site_path())
if attached_file.startswith("/private/files/"):
return os.path.join(site_path, "private", "files", attached_file.split("/private/files/", 1)[1])
if attached_file.startswith("/files/"):
return os.path.join(site_path, "public", "files", attached_file.split("/files/", 1)[1])
raise Exception(f"Unbekannter Dateipfad: {attached_file}")
def _group_rows(rows: list[CanonicalRow]) -> dict[str, list[CanonicalRow]]:
grouped: dict[str, list[CanonicalRow]] = defaultdict(list)
for r in rows:
key = r.customer_external_ref or "__NO_CUSTOMER__"
grouped[key].append(r)
return dict(grouped)
def _customer_billing_mode(customer_ref: str, fallback: str | None) -> str | None:
if not customer_ref or customer_ref == "__NO_CUSTOMER__":
customer_ref = fallback
if not customer_ref:
return None
mode = frappe.db.get_value("Customer", customer_ref, "billing_mode")
return mode
def _persist_lines(run, outcomes: Iterable[dict]) -> None:
"""Schreibt die Import-Ergebnisse als Supplier Import Line-Zeilen in den Run."""
run.reload()
run.set("lines", [])
for o in outcomes:
row: CanonicalRow = o["row"]
target_doctype = None
target_name = None
target_row = None
if o.get("status") == "Created":
# Aus Subscription-Event das Ziel rückgewinnen (best effort)
event_name = o.get("event")
if event_name:
ev = frappe.db.get_value(
"Supply Subscription Event", event_name,
["sales_invoice", "sales_invoice_item_row",
"delivery_note", "delivery_note_item_row"],
as_dict=True,
)
if ev:
if ev.get("sales_invoice"):
target_doctype = "Sales Invoice"
target_name = ev["sales_invoice"]
target_row = ev.get("sales_invoice_item_row")
elif ev.get("delivery_note"):
target_doctype = "Delivery Note"
target_name = ev["delivery_note"]
target_row = ev.get("delivery_note_item_row")
run.append("lines", {
"source_row": row.source_row,
"customer": _customer_or_null(row.customer_external_ref),
"customer_external_ref": row.customer_external_ref,
"vendor_product_id": row.vendor_product_id,
"qty": row.qty or 0,
"rate": row.unit_price or row.list_price,
"amount": row.amount,
"period_start": row.period_start,
"period_end": row.period_end,
"subscription_external_id": row.subscription_external_id,
"subscription_start": row.subscription_start,
"billing_plan": row.billing_plan,
"term_duration": row.term_duration,
"booking_type": row.booking_type,
"marketplace_ref": row.marketplace_ref,
"order_ref": row.order_ref,
"target_doc_type": target_doctype,
"target_doc_name": target_name,
"target_doc_row": target_row,
"line_status": o.get("status") or "Pending",
"error_message": o.get("error"),
"raw_payload": json.dumps(row.raw, default=str, ensure_ascii=False) if row.raw else None,
})
run.save(ignore_permissions=True)
def _customer_or_null(ref: str | None) -> str | None:
if not ref:
return None
if frappe.db.exists("Customer", ref):
return ref
return None
+102
View File
@@ -0,0 +1,102 @@
"""Upsert-Logik für Supply Subscription + Supply Subscription Event."""
from __future__ import annotations
import frappe
from msp.importers.base import CanonicalRow
def upsert_supply_subscription(row: CanonicalRow, *, customer: str,
supplier: str | None, vendor: str | None,
item: str | None) -> str | None:
"""Legt bei Bedarf eine Supply Subscription an oder aktualisiert das Golden Record.
Rückgabe: der Name der Supply Subscription (z. B. ``SUB-<external_id>``) oder
None, wenn kein ``subscription_external_id`` in der Zeile steht (dann wird
keine Subscription geführt).
"""
external_id = (row.subscription_external_id or "").strip()
if not external_id:
return None
existing = frappe.db.get_value("Supply Subscription", {"external_id": external_id}, "name")
payload = {
"doctype": "Supply Subscription",
"external_id": external_id,
"customer": customer,
"supplier": supplier,
"vendor": vendor,
"item": item,
"vendor_product_id": row.vendor_product_id,
"billing_plan": row.billing_plan,
"term_duration": row.term_duration,
"start_date": row.subscription_start,
"status": "Active",
}
if existing:
doc = frappe.get_doc("Supply Subscription", existing)
# Nur leere/stale Felder auffüllen, keine Werte überschreiben, die vom
# User manuell gesetzt wurden (Ausnahme: Stammdaten wie Produkt/Item).
doc.customer = customer or doc.customer
doc.supplier = supplier or doc.supplier
doc.vendor = vendor or doc.vendor
if item and not doc.item:
doc.item = item
doc.vendor_product_id = row.vendor_product_id or doc.vendor_product_id
doc.billing_plan = row.billing_plan or doc.billing_plan
doc.term_duration = row.term_duration or doc.term_duration
if row.subscription_start and not doc.start_date:
doc.start_date = row.subscription_start
doc.save(ignore_permissions=True)
return doc.name
new = frappe.get_doc(payload)
new.insert(ignore_permissions=True)
return new.name
def create_event(row: CanonicalRow, *, supply_subscription: str,
target_doctype: str, target_name: str, target_row: str | None = None,
supplier_import_line_row: str | None = None) -> str:
"""Erzeugt ein Supply Subscription Event für eine abgerechnete Periode."""
event = frappe.get_doc({
"doctype": "Supply Subscription Event",
"supply_subscription": supply_subscription,
"period_start": row.period_start,
"period_end": row.period_end,
"qty": int(row.qty) if row.qty else 0,
"rate": row.unit_price or row.list_price or 0,
"amount": row.amount or 0,
"booking_type": row.booking_type or "Billing",
"sales_invoice": target_name if target_doctype == "Sales Invoice" else None,
"sales_invoice_item_row": target_row if target_doctype == "Sales Invoice" else None,
"delivery_note": target_name if target_doctype == "Delivery Note" else None,
"delivery_note_item_row": target_row if target_doctype == "Delivery Note" else None,
"supplier_import_line": supplier_import_line_row,
})
event.insert(ignore_permissions=True)
return event.name
def event_exists_for_period(subscription_external_id: str, period_start,
period_end) -> bool:
"""Deduplizierungs-Check: Wurde diese Subscription+Period bereits importiert?"""
if not (subscription_external_id and period_start and period_end):
return False
subscription_name = frappe.db.get_value(
"Supply Subscription",
{"external_id": subscription_external_id}, "name",
)
if not subscription_name:
return False
return bool(frappe.db.exists(
"Supply Subscription Event",
{
"supply_subscription": subscription_name,
"period_start": period_start,
"period_end": period_end,
},
))