diff --git a/msp/importers/__init__.py b/msp/importers/__init__.py
new file mode 100644
index 0000000..e69de29
diff --git a/msp/importers/adn_monthly_csv.py b/msp/importers/adn_monthly_csv.py
new file mode 100644
index 0000000..8e38ddb
--- /dev/null
+++ b/msp/importers/adn_monthly_csv.py
@@ -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,
+ )
diff --git a/msp/importers/base.py b/msp/importers/base.py
new file mode 100644
index 0000000..790d292
--- /dev/null
+++ b/msp/importers/base.py
@@ -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",
+]
diff --git a/msp/importers/builder.py b/msp/importers/builder.py
new file mode 100644
index 0000000..6fd5ea5
--- /dev/null
+++ b/msp/importers/builder.py
@@ -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 "
".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)
diff --git a/msp/importers/registry.py b/msp/importers/registry.py
new file mode 100644
index 0000000..74fc971
--- /dev/null
+++ b/msp/importers/registry.py
@@ -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
diff --git a/msp/importers/resolution.py b/msp/importers/resolution.py
new file mode 100644
index 0000000..24cd630
--- /dev/null
+++ b/msp/importers/resolution.py
@@ -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"
diff --git a/msp/importers/run_orchestrator.py b/msp/importers/run_orchestrator.py
new file mode 100644
index 0000000..4e711ce
--- /dev/null
+++ b/msp/importers/run_orchestrator.py
@@ -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
diff --git a/msp/importers/subscriptions.py b/msp/importers/subscriptions.py
new file mode 100644
index 0000000..20f9708
--- /dev/null
+++ b/msp/importers/subscriptions.py
@@ -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-``) 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,
+ },
+ ))