Files
Holzleitner---Backend--aktu…/crates/infrastructure/src/persistence/tour_repository.rs
Dennis Nemec 954c5f52b2 feat(zahlung): Zahlungsabwicklung als eigenes Protokoll (POST /deliveries/{id}/payment)
- Neue Tabelle delivery_payments (append-only, idempotent ueber
  client_event_id): Methode + Code-Snapshot, server-seitig berechneter
  offener Betrag, Fahrer, Fahrzeug, Zeitpunkt
- Endpoint prueft unter Zeilen-Lock: Lieferung aktiv, Methode aktiv,
  offener Betrag > 0 und identisch mit dem vom Fahrer bestaetigten Betrag
- Offener Betrag als gemeinsamer Helper (open_amount_cents) fuer
  Zahlungsprotokoll und Abschluss
- Abschluss-Gate: gueltige protokollierte Zahlung erfuellt die
  Inkasso-Pflicht und liefert die Methode; delivery_completions.payment_id
  verknuepft den Abschluss mit der Zahlung. Altes payment_collected-Flag
  bleibt fuer aeltere App-Versionen gueltig
- Tour-Aggregat und Admin-Belegdetails liefern die juengste Zahlung
- Einzel-Reset loescht auch das Zahlungsprotokoll
- Integrationstest (ignored, braucht Wegwerf-DB) fuer Protokoll + Gate

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-25 14:04:35 +02:00

1348 lines
42 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

//! Postgres-Implementierung des `TourRepository`-Ports.
//!
//! Drei Operationen, getrennt umgesetzt:
//! * `find_today_for_driver` — eine Query (Tour + Lieferzahl pro Tour).
//! * `find_details_by_id` — pro Aggregat ein paar gezielte Queries,
//! anschließend in-memory zusammenbauen. Bewusst keine eine-Big-Join-
//! Query: das vervielfacht Daten über die Leitung, die wir clientseitig
//! wieder deduplizieren müssten — und Postgres handhabt 5–7 kleine
//! Prepared Statements im selben Pool effizient genug.
//! * `upsert_from_sync` — eine Transaktion, idempotent per UPSERT auf
//! den fachlichen Keys. Bestehende `scan_state`-Werte bleiben
//! unangetastet: das ERP weiß nichts davon und darf sie nicht
//! überschreiben.
use std::collections::HashMap;
use async_trait::async_trait;
use chrono::{DateTime, NaiveDate, Utc};
use sqlx::{PgPool, Postgres, Transaction};
use uuid::Uuid;
use holzleitner_application::dto::{
DeliveryOrderEntry, DeliveryWithItems, SyncContactSource, SyncDelivery, SyncDeliveryItem,
SyncTourRequest, TourDetails, TourSummary,
};
use holzleitner_application::error::ApplicationError;
use holzleitner_application::ports::TourRepository;
use holzleitner_domain::{
Address, Article, ContactChannel, ContactKind, ContactRole, ContactSource, Customer,
CustomerContact, Delivery, DeliveryCredit, DeliveryItem, DeliveryNote, DeliveryPayment,
DeliveryServiceValue, DeliveryState, ScanState, ScanStatus, Service, ServiceKind, Tour, Warehouse,
};
use super::delivery_payment_repository::{PAYMENT_COLUMNS, PaymentRow};
pub struct PgTourRepository {
pool: PgPool,
}
impl PgTourRepository {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
}
// ===== Row-Typen =========================================================
//
// Eigene FromRow-Strukturen in der Infrastructure halten die Domain-Typen
// frei von `sqlx`-Traits. Mapping in private Helfer.
#[derive(sqlx::FromRow)]
struct TourRow {
id: Uuid,
account_id: i64,
tour_date: NaiveDate,
synced_at: DateTime<Utc>,
}
#[derive(sqlx::FromRow)]
struct TourSummaryRow {
id: Uuid,
tour_date: NaiveDate,
delivery_count: i64,
}
#[derive(sqlx::FromRow)]
struct DeliveryRow {
id: Uuid,
tour_id: Uuid,
erp_belegart_id: i64,
erp_belegnummer: String,
customer_id: Uuid,
snap_street: String,
snap_house_number: String,
snap_postal_code: String,
snap_city: String,
snap_country: String,
assigned_car_id: Option<Uuid>,
desired_time: Option<String>,
special_agreements: Option<String>,
state: String,
state_reason: Option<String>,
sort_order: i32,
prepaid_amount: f64,
payment_method_id: Uuid,
}
#[derive(sqlx::FromRow)]
struct DeliveryItemRow {
id: Uuid,
delivery_id: Uuid,
article_id: Uuid,
required_quantity: i32,
warehouse_id: Uuid,
unit_price: f64,
belegzeilen_nr: i32,
komponenten_artikel_nr: Option<String>,
parent_artikel_nr: Option<String>,
scanned_quantity: i32,
credited_quantity: i32,
scan_status: String,
held_reason: Option<String>,
scan_last_updated_at: DateTime<Utc>,
}
#[derive(sqlx::FromRow)]
struct CustomerRow {
id: Uuid,
erp_customer_id: i64,
name: String,
street: String,
house_number: String,
postal_code: String,
city: String,
country: String,
}
#[derive(sqlx::FromRow)]
struct CustomerContactRow {
id: Uuid,
customer_id: Uuid,
name: String,
phone: Option<String>,
email: Option<String>,
}
#[derive(sqlx::FromRow)]
struct ArticleRow {
id: Uuid,
article_number: String,
name: String,
scannable: bool,
default_warehouse_id: Option<Uuid>,
}
#[derive(sqlx::FromRow)]
struct WarehouseRow {
id: Uuid,
code: String,
name: String,
is_standard: bool,
}
#[derive(sqlx::FromRow)]
struct ContactLinkRow {
delivery_id: Uuid,
customer_contact_id: Uuid,
}
#[derive(sqlx::FromRow)]
struct ContactSourceRow {
id: Uuid,
delivery_id: Uuid,
role: String,
anrede: Option<String>,
titel: Option<String>,
name1: Option<String>,
name2: Option<String>,
name3: Option<String>,
abteilung: Option<String>,
funktion: Option<String>,
}
#[derive(sqlx::FromRow)]
struct ContactChannelRow {
id: Uuid,
source_id: Uuid,
kind: String,
position: i16,
value: String,
}
#[derive(sqlx::FromRow)]
struct DeliveryNoteRow {
id: Uuid,
delivery_id: Uuid,
text: Option<String>,
image_attachment: Option<String>,
author_personalnummer: i64,
author_car_id: Option<Uuid>,
credit_delivery_item_id: Option<Uuid>,
is_amount_credit_note: bool,
image_attachment_deleted: bool,
created_at: DateTime<Utc>,
}
// ===== Mapping Row -> Domain =============================================
fn map_tour(row: TourRow) -> Tour {
Tour {
id: row.id,
account_id: row.account_id,
date: row.tour_date,
synced_at: row.synced_at,
}
}
fn parse_delivery_state(value: &str) -> Result<DeliveryState, ApplicationError> {
match value {
"active" => Ok(DeliveryState::Active),
"held" => Ok(DeliveryState::Held),
"canceled" => Ok(DeliveryState::Canceled),
"completed" => Ok(DeliveryState::Completed),
other => Err(ApplicationError::Repository(format!(
"unknown delivery state '{other}'"
))),
}
}
fn parse_scan_status(value: &str) -> Result<ScanStatus, ApplicationError> {
match value {
"in_progress" => Ok(ScanStatus::InProgress),
"done" => Ok(ScanStatus::Done),
"held" => Ok(ScanStatus::Held),
"removed" => Ok(ScanStatus::Removed),
other => Err(ApplicationError::Repository(format!(
"unknown scan status '{other}'"
))),
}
}
fn map_item(row: DeliveryItemRow) -> Result<DeliveryItem, ApplicationError> {
Ok(DeliveryItem {
id: row.id,
delivery_id: row.delivery_id,
article_id: row.article_id,
required_quantity: row.required_quantity,
warehouse_id: row.warehouse_id,
unit_price: row.unit_price,
belegzeilen_nr: row.belegzeilen_nr,
komponenten_artikel_nr: row.komponenten_artikel_nr,
parent_artikel_nr: row.parent_artikel_nr,
scan_state: ScanState {
scanned_quantity: row.scanned_quantity,
credited_quantity: row.credited_quantity,
status: parse_scan_status(&row.scan_status)?,
held_reason: row.held_reason,
last_updated_at: row.scan_last_updated_at,
},
})
}
fn map_customer(row: CustomerRow) -> Customer {
Customer {
id: row.id,
erp_customer_id: row.erp_customer_id,
name: row.name,
address: Address {
street: row.street,
house_number: row.house_number,
postal_code: row.postal_code,
city: row.city,
country: row.country,
},
}
}
fn map_contact(row: CustomerContactRow) -> CustomerContact {
CustomerContact {
id: row.id,
customer_id: row.customer_id,
name: row.name,
phone: row.phone,
email: row.email,
}
}
fn map_contact_source(row: ContactSourceRow) -> Result<ContactSource, ApplicationError> {
Ok(ContactSource {
id: row.id,
delivery_id: row.delivery_id,
role: role_from_db(&row.role)?,
anrede: row.anrede,
titel: row.titel,
name1: row.name1,
name2: row.name2,
name3: row.name3,
abteilung: row.abteilung,
funktion: row.funktion,
})
}
fn map_contact_channel(row: ContactChannelRow) -> Result<ContactChannel, ApplicationError> {
Ok(ContactChannel {
id: row.id,
source_id: row.source_id,
kind: kind_from_db(&row.kind)?,
position: row.position,
value: row.value,
})
}
fn map_article(row: ArticleRow) -> Article {
Article {
id: row.id,
article_number: row.article_number,
name: row.name,
scannable: row.scannable,
default_warehouse_id: row.default_warehouse_id,
}
}
fn map_note(row: DeliveryNoteRow) -> DeliveryNote {
DeliveryNote {
id: row.id,
delivery_id: row.delivery_id,
text: row.text,
image_attachment: row.image_attachment,
author_personalnummer: row.author_personalnummer,
author_car_id: row.author_car_id,
credit_delivery_item_id: row.credit_delivery_item_id,
is_amount_credit_note: row.is_amount_credit_note,
image_attachment_deleted: row.image_attachment_deleted,
created_at: row.created_at,
}
}
#[derive(sqlx::FromRow)]
struct CreditRow {
delivery_id: Uuid,
action: String,
amount_cents: i64,
reason: Option<String>,
}
/// Aktuelles Gutschrift-Ereignis → Domänenobjekt. `remove` (oder unbekannte
/// Action) liefert `None`, sodass entfernte Gutschriften nicht erscheinen.
fn map_credit(row: CreditRow) -> Option<DeliveryCredit> {
if row.action != "set" {
return None;
}
Some(DeliveryCredit {
delivery_id: row.delivery_id,
amount_cents: row.amount_cents,
reason: row.reason.unwrap_or_default(),
})
}
#[derive(sqlx::FromRow)]
struct ServiceRow {
id: Uuid,
key: String,
name: String,
kind: String,
min_value: Option<i32>,
max_value: Option<i32>,
active: bool,
sort_order: i32,
}
fn map_service(row: ServiceRow) -> Result<Service, ApplicationError> {
let kind = match row.kind.as_str() {
"boolean" => ServiceKind::Boolean,
"numeric" => ServiceKind::Numeric,
other => {
return Err(ApplicationError::Repository(format!(
"unknown service kind '{other}'"
)));
}
};
Ok(Service {
id: row.id,
key: row.key,
name: row.name,
kind,
min_value: row.min_value,
max_value: row.max_value,
active: row.active,
sort_order: row.sort_order,
})
}
#[derive(sqlx::FromRow)]
struct DeliveryServiceRow {
delivery_id: Uuid,
service_id: Uuid,
bool_value: Option<bool>,
numeric_value: Option<i32>,
}
fn map_delivery_service(row: DeliveryServiceRow) -> DeliveryServiceValue {
DeliveryServiceValue {
delivery_id: row.delivery_id,
service_id: row.service_id,
bool_value: row.bool_value,
numeric_value: row.numeric_value,
}
}
fn map_warehouse(row: WarehouseRow) -> Warehouse {
Warehouse {
id: row.id,
code: row.code,
name: row.name,
is_standard: row.is_standard,
}
}
fn map_delivery(
row: DeliveryRow,
contact_person_ids: Vec<Uuid>,
) -> Result<(Delivery, i32), ApplicationError> {
let state = parse_delivery_state(&row.state)?;
let delivery = Delivery {
id: row.id,
tour_id: row.tour_id,
erp_belegart_id: row.erp_belegart_id,
erp_belegnummer: row.erp_belegnummer,
customer_id: row.customer_id,
delivery_address_snapshot: Address {
street: row.snap_street,
house_number: row.snap_house_number,
postal_code: row.snap_postal_code,
city: row.snap_city,
country: row.snap_country,
},
assigned_car_id: row.assigned_car_id,
contact_person_ids,
desired_time: row.desired_time,
special_agreements: row.special_agreements,
state,
state_reason: row.state_reason,
prepaid_amount: row.prepaid_amount,
payment_method_id: row.payment_method_id,
};
Ok((delivery, row.sort_order))
}
// ===== Helfer: Error-Mapping =============================================
fn db<E: std::fmt::Display>(e: E) -> ApplicationError {
ApplicationError::Repository(e.to_string())
}
// ===== Trait-Implementierung =============================================
#[async_trait]
impl TourRepository for PgTourRepository {
async fn find_today_for_driver(
&self,
personalnummer: i64,
today: NaiveDate,
) -> Result<Vec<TourSummary>, ApplicationError> {
let rows = sqlx::query_as::<_, TourSummaryRow>(
r#"
SELECT t.id, t.tour_date, COUNT(d.id) AS delivery_count
FROM tours t
LEFT JOIN deliveries d ON d.tour_id = t.id
WHERE t.account_id = $1 AND t.tour_date = $2
GROUP BY t.id, t.tour_date
ORDER BY t.tour_date
"#,
)
.bind(personalnummer)
.bind(today)
.fetch_all(&self.pool)
.await
.map_err(db)?;
Ok(rows
.into_iter()
.map(|r| TourSummary {
tour_id: r.id,
tour_date: r.tour_date,
delivery_count: r.delivery_count,
})
.collect())
}
async fn find_details_by_id(
&self,
tour_id: Uuid,
) -> Result<Option<TourDetails>, ApplicationError> {
// 1. Tour selbst
let Some(tour_row) = sqlx::query_as::<_, TourRow>(
"SELECT id, account_id, tour_date, synced_at FROM tours WHERE id = $1",
)
.bind(tour_id)
.fetch_optional(&self.pool)
.await
.map_err(db)?
else {
return Ok(None);
};
let tour = map_tour(tour_row);
// 2. Lieferungen
let delivery_rows = sqlx::query_as::<_, DeliveryRow>(
r#"
SELECT
id, tour_id, erp_belegart_id, erp_belegnummer, customer_id,
snap_street, snap_house_number, snap_postal_code, snap_city, snap_country,
assigned_car_id, desired_time, special_agreements,
state, state_reason, sort_order,
prepaid_amount, payment_method_id
FROM deliveries
WHERE tour_id = $1
ORDER BY sort_order, erp_belegnummer
"#,
)
.bind(tour_id)
.fetch_all(&self.pool)
.await
.map_err(db)?;
let delivery_ids: Vec<Uuid> = delivery_rows.iter().map(|d| d.id).collect();
// 3. Kontakt-Person-Verknüpfungen
let contact_links = sqlx::query_as::<_, ContactLinkRow>(
r#"
SELECT delivery_id, customer_contact_id
FROM delivery_contact_persons
WHERE delivery_id = ANY($1)
"#,
)
.bind(&delivery_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?;
let mut contacts_per_delivery: HashMap<Uuid, Vec<Uuid>> = HashMap::new();
for link in contact_links {
contacts_per_delivery
.entry(link.delivery_id)
.or_default()
.push(link.customer_contact_id);
}
// 4. Positionen
let item_rows = sqlx::query_as::<_, DeliveryItemRow>(
r#"
SELECT
id, delivery_id, article_id, required_quantity, warehouse_id,
unit_price, belegzeilen_nr, komponenten_artikel_nr, parent_artikel_nr,
scanned_quantity, credited_quantity, scan_status, held_reason,
scan_last_updated_at
FROM delivery_items
WHERE delivery_id = ANY($1)
ORDER BY delivery_id, belegzeilen_nr, komponenten_artikel_nr NULLS FIRST
"#,
)
.bind(&delivery_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?;
let mut items_per_delivery: HashMap<Uuid, Vec<DeliveryItem>> = HashMap::new();
let mut article_ids = std::collections::BTreeSet::new();
let mut warehouse_ids = std::collections::BTreeSet::new();
for row in item_rows {
article_ids.insert(row.article_id);
warehouse_ids.insert(row.warehouse_id);
let delivery_id = row.delivery_id;
let item = map_item(row)?;
items_per_delivery
.entry(delivery_id)
.or_default()
.push(item);
}
// 5. Lieferungen + Items kombinieren
let mut customer_ids = std::collections::BTreeSet::new();
let mut deliveries = Vec::with_capacity(delivery_rows.len());
for row in delivery_rows {
customer_ids.insert(row.customer_id);
let delivery_id = row.id;
let contact_ids = contacts_per_delivery.remove(&delivery_id).unwrap_or_default();
let items = items_per_delivery.remove(&delivery_id).unwrap_or_default();
let (delivery, sort_order) = map_delivery(row, contact_ids)?;
deliveries.push(DeliveryWithItems {
delivery,
sort_order,
items,
});
}
// 6. Lookup-Stammdaten
let customer_ids_vec: Vec<Uuid> = customer_ids.into_iter().collect();
let customers = sqlx::query_as::<_, CustomerRow>(
r#"
SELECT id, erp_customer_id, name, street, house_number, postal_code, city, country
FROM customers
WHERE id = ANY($1)
ORDER BY name
"#,
)
.bind(&customer_ids_vec)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(map_customer)
.collect::<Vec<_>>();
let customer_contacts = sqlx::query_as::<_, CustomerContactRow>(
r#"
SELECT id, customer_id, name, phone, email
FROM customer_contacts
WHERE customer_id = ANY($1)
ORDER BY name
"#,
)
.bind(&customer_ids_vec)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(map_contact)
.collect::<Vec<_>>();
let article_ids_vec: Vec<Uuid> = article_ids.into_iter().collect();
let articles = sqlx::query_as::<_, ArticleRow>(
r#"
SELECT id, article_number, name, scannable, default_warehouse_id
FROM articles
WHERE id = ANY($1)
ORDER BY article_number
"#,
)
.bind(&article_ids_vec)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(map_article)
.collect::<Vec<_>>();
let warehouse_ids_vec: Vec<Uuid> = warehouse_ids.into_iter().collect();
let warehouses = sqlx::query_as::<_, WarehouseRow>(
r#"
SELECT id, code, name, is_standard
FROM warehouses
WHERE id = ANY($1)
ORDER BY code
"#,
)
.bind(&warehouse_ids_vec)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(map_warehouse)
.collect::<Vec<_>>();
// 7. Notizen aller Lieferungen dieser Tour.
let notes = sqlx::query_as::<_, DeliveryNoteRow>(
r#"
SELECT dn.id, dn.delivery_id, dn.text, dn.image_attachment,
dn.author_personalnummer, dn.author_car_id,
dn.credit_delivery_item_id, dn.is_amount_credit_note,
(att.deleted_at IS NOT NULL) AS image_attachment_deleted,
dn.created_at
FROM delivery_notes dn
LEFT JOIN attachments att ON att.id = dn.image_attachment::uuid
WHERE dn.delivery_id = ANY($1)
ORDER BY dn.delivery_id, dn.created_at
"#,
)
.bind(&delivery_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(map_note)
.collect::<Vec<_>>();
// 8. Aktuelle Betrags-Gutschriften: jüngstes Ereignis pro Lieferung,
// nur solange der letzte Stand `set` ist.
let credits = sqlx::query_as::<_, CreditRow>(
r#"
SELECT DISTINCT ON (delivery_id)
delivery_id, action, amount_cents, reason
FROM delivery_credit_audit
WHERE delivery_id = ANY($1)
ORDER BY delivery_id, recorded_at DESC, id DESC
"#,
)
.bind(&delivery_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.filter_map(map_credit)
.collect::<Vec<_>>();
// 8b. Jüngste protokollierte Zahlungsabwicklung pro Lieferung.
let payments = sqlx::query_as::<_, PaymentRow>(&format!(
r#"
SELECT DISTINCT ON (delivery_id) {PAYMENT_COLUMNS}
FROM delivery_payments
WHERE delivery_id = ANY($1)
ORDER BY delivery_id, recorded_at DESC, id DESC
"#
))
.bind(&delivery_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(DeliveryPayment::from)
.collect::<Vec<_>>();
// 9. Aktive Service-Definitionen (Stammdaten) — die App rendert daraus
// Phase 4.
let services = sqlx::query_as::<_, ServiceRow>(
r#"
SELECT id, key, name, kind, min_value, max_value, active, sort_order
FROM services
WHERE active = TRUE
ORDER BY sort_order, name
"#,
)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(map_service)
.collect::<Result<Vec<_>, _>>()?;
// 10. Pro-Lieferung gesetzte Service-Werte.
let delivery_services = sqlx::query_as::<_, DeliveryServiceRow>(
r#"
SELECT delivery_id, service_id, bool_value, numeric_value
FROM delivery_services
WHERE delivery_id = ANY($1)
"#,
)
.bind(&delivery_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?
.into_iter()
.map(map_delivery_service)
.collect::<Vec<_>>();
// 11. Kontaktdaten-Snapshots aller Lieferungen + ihre Kanäle.
// Reihenfolge: Quellen pro Lieferung nach Rolle, Kanäle pro
// Quelle nach Art und ERP-Position — so kommt „Telefon"
// vor „Telefon2", die App muss nicht extra sortieren.
let source_rows = sqlx::query_as::<_, ContactSourceRow>(
r#"
SELECT id, delivery_id, role,
anrede, titel, name1, name2, name3, abteilung, funktion
FROM delivery_contact_sources
WHERE delivery_id = ANY($1)
ORDER BY delivery_id, role
"#,
)
.bind(&delivery_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?;
let source_ids: Vec<Uuid> = source_rows.iter().map(|r| r.id).collect();
let contact_sources = source_rows
.into_iter()
.map(map_contact_source)
.collect::<Result<Vec<_>, _>>()?;
let channel_rows = sqlx::query_as::<_, ContactChannelRow>(
r#"
SELECT id, source_id, kind, position, value
FROM delivery_contact_channels
WHERE source_id = ANY($1)
ORDER BY source_id, kind, position
"#,
)
.bind(&source_ids)
.fetch_all(&self.pool)
.await
.map_err(db)?;
let contact_channels = channel_rows
.into_iter()
.map(map_contact_channel)
.collect::<Result<Vec<_>, _>>()?;
Ok(Some(TourDetails {
tour,
deliveries,
customers,
customer_contacts,
articles,
warehouses,
notes,
credits,
payments,
services,
delivery_services,
contact_sources,
contact_channels,
}))
}
async fn set_delivery_order(
&self,
tour_id: Uuid,
delivery_ids: &[Uuid],
) -> Result<Vec<DeliveryOrderEntry>, ApplicationError> {
let mut tx = self.pool.begin().await.map_err(db)?;
// 1. Lock alle Lieferungen dieser Tour. Liefert leer, wenn Tour
// nicht existiert oder keine Lieferungen hat.
let existing: Vec<Uuid> = sqlx::query_scalar(
"SELECT id FROM deliveries WHERE tour_id = $1 FOR UPDATE",
)
.bind(tour_id)
.fetch_all(&mut *tx)
.await
.map_err(db)?;
if existing.is_empty() {
tx.rollback().await.map_err(db)?;
return Err(ApplicationError::NotFound);
}
// 2. Mengen-Match: Input muss exakt der Tour entsprechen.
let existing_set: std::collections::HashSet<Uuid> = existing.iter().copied().collect();
let input_set: std::collections::HashSet<Uuid> = delivery_ids.iter().copied().collect();
if existing_set != input_set {
tx.rollback().await.map_err(db)?;
let fremde: Vec<Uuid> =
input_set.difference(&existing_set).copied().collect();
let fehlende: Vec<Uuid> =
existing_set.difference(&input_set).copied().collect();
return Err(ApplicationError::Validation(format!(
"delivery_ids match nicht zur tour (fehlende: {fehlende:?}, fremde: {fremde:?})"
)));
}
// 3. Bulk-Update via UNNEST.
let positions: Vec<i32> = (1..=delivery_ids.len() as i32).collect();
sqlx::query(
r#"
UPDATE deliveries AS d
SET sort_order = data.new_order
FROM (
SELECT UNNEST($1::uuid[]) AS id,
UNNEST($2::int[]) AS new_order
) AS data
WHERE d.id = data.id
"#,
)
.bind(delivery_ids)
.bind(&positions)
.execute(&mut *tx)
.await
.map_err(db)?;
tx.commit().await.map_err(db)?;
Ok(delivery_ids
.iter()
.zip(positions.iter())
.map(|(id, pos)| DeliveryOrderEntry {
delivery_id: *id,
sort_order: *pos,
})
.collect())
}
async fn upsert_from_sync(
&self,
request: &SyncTourRequest,
) -> Result<Uuid, ApplicationError> {
let mut tx = self.pool.begin().await.map_err(db)?;
// 0. Fahrer-/Account-Konto sicherstellen — der ERP-`Vertreter` muss als
// `accounts`-Zeile existieren (FK von `tours`). Auto-Provisionierung:
// fehlende Konten werden mit Default-Namen angelegt; bestehende
// bleiben unangetastet (DO NOTHING überschreibt keinen Namen).
sqlx::query(
r#"
INSERT INTO accounts (personalnummer, name)
VALUES ($1, $2)
ON CONFLICT (personalnummer) DO NOTHING
"#,
)
.bind(request.driver_personalnummer)
.bind(format!("Fahrer {}", request.driver_personalnummer))
.execute(&mut *tx)
.await
.map_err(db)?;
// 1. Tour upserten — Identität: (account_id, tour_date)
let tour_id: Uuid = sqlx::query_scalar(
r#"
INSERT INTO tours (account_id, tour_date)
VALUES ($1, $2)
ON CONFLICT (account_id, tour_date) DO UPDATE
SET synced_at = now()
RETURNING id
"#,
)
.bind(request.driver_personalnummer)
.bind(request.tour_date)
.fetch_one(&mut *tx)
.await
.map_err(db)?;
for delivery in &request.deliveries {
// 2. Kunde upserten — Identität: erp_customer_id
let customer_id = upsert_customer(&mut tx, delivery).await?;
// 3. Lieferung upserten — Identität: (belegart_id, belegnummer).
// Bestehende Lieferung bleibt mit ihrem state/cancellation
// erhalten; nur Stammdaten + sort_order werden refresht.
let delivery_id = upsert_delivery(&mut tx, tour_id, customer_id, delivery).await?;
// 3a. Kontaktdaten-Snapshot neu schreiben. Snapshot-Semantik:
// beim Sync wird der Stand vom ERP übernommen, ältere Stände
// verworfen. Der CASCADE-DELETE räumt auch die Channels mit.
replace_contact_sources(&mut tx, delivery_id, &delivery.contact_sources).await?;
for item in &delivery.items {
let warehouse_id = upsert_warehouse(&mut tx, item).await?;
let article_id = upsert_article(&mut tx, item, warehouse_id).await?;
upsert_delivery_item(&mut tx, delivery_id, article_id, warehouse_id, item).await?;
}
}
tx.commit().await.map_err(db)?;
Ok(tour_id)
}
async fn delete_all_tours(&self) -> Result<u64, ApplicationError> {
// DELETE FROM tours cascadet per FK auf deliveries → delivery_items →
// scan_audit, delivery_notes, delivery_credit_audit, delivery_payments,
// delivery_services,
// delivery_completions, attachments, delivery_contact_persons.
let res = sqlx::query("DELETE FROM tours")
.execute(&self.pool)
.await
.map_err(db)?;
Ok(res.rows_affected())
}
async fn reset_delivery_by_belegnummer(
&self,
belegnummer: &str,
) -> Result<u64, ApplicationError> {
let mut tx = self.pool.begin().await.map_err(db)?;
// Subquery-Filter überall identisch: Lieferung(en) mit dieser Belegnr.
const BY_BELEG: &str =
"SELECT id FROM deliveries WHERE erp_belegnummer = $1";
// Abschluss + (evtl.) offener Report-Job + Scan-/Gutschrift-Audit löschen.
sqlx::query(&format!(
"DELETE FROM delivery_report_jobs WHERE delivery_id IN ({BY_BELEG})"
))
.bind(belegnummer)
.execute(&mut *tx)
.await
.map_err(db)?;
sqlx::query(&format!(
"DELETE FROM delivery_completions WHERE delivery_id IN ({BY_BELEG})"
))
.bind(belegnummer)
.execute(&mut *tx)
.await
.map_err(db)?;
sqlx::query(
"DELETE FROM scan_audit WHERE delivery_item_id IN (\
SELECT di.id FROM delivery_items di \
JOIN deliveries d ON d.id = di.delivery_id \
WHERE d.erp_belegnummer = $1)",
)
.bind(belegnummer)
.execute(&mut *tx)
.await
.map_err(db)?;
// Nach den Abschlüssen (FK payment_id), vor dem Status-Reset.
sqlx::query(&format!(
"DELETE FROM delivery_payments WHERE delivery_id IN ({BY_BELEG})"
))
.bind(belegnummer)
.execute(&mut *tx)
.await
.map_err(db)?;
sqlx::query(&format!(
"DELETE FROM delivery_credit_audit WHERE delivery_id IN ({BY_BELEG})"
))
.bind(belegnummer)
.execute(&mut *tx)
.await
.map_err(db)?;
// Positionen zurück (verladen + gutgeschrieben = 0, Status offen).
sqlx::query(&format!(
"UPDATE delivery_items \
SET scanned_quantity = 0, scan_status = 'in_progress', \
held_reason = NULL, credited_quantity = 0, \
scan_last_updated_at = now() \
WHERE delivery_id IN ({BY_BELEG})"
))
.bind(belegnummer)
.execute(&mut *tx)
.await
.map_err(db)?;
// Lieferung wieder aktiv; Status-/Auto-/Review-Felder leeren.
let res = sqlx::query(
"UPDATE deliveries \
SET state = 'active', state_reason = NULL, assigned_car_id = NULL, \
review_resolved_at = NULL, review_resolved_by = NULL, \
review_note = NULL \
WHERE erp_belegnummer = $1",
)
.bind(belegnummer)
.execute(&mut *tx)
.await
.map_err(db)?;
tx.commit().await.map_err(db)?;
Ok(res.rows_affected())
}
async fn find_tour_id_by_belegnummer(
&self,
belegnummer: &str,
) -> Result<Option<Uuid>, ApplicationError> {
let tour_id: Option<Uuid> = sqlx::query_scalar(
"SELECT tour_id FROM deliveries WHERE erp_belegnummer = $1 LIMIT 1",
)
.bind(belegnummer)
.fetch_optional(&self.pool)
.await
.map_err(db)?;
Ok(tour_id)
}
}
// ===== Upsert-Helfer =====================================================
async fn upsert_warehouse(
tx: &mut Transaction<'_, Postgres>,
item: &SyncDeliveryItem,
) -> Result<Uuid, ApplicationError> {
let id: Uuid = sqlx::query_scalar(
r#"
INSERT INTO warehouses (code, name)
VALUES ($1, $2)
ON CONFLICT (code) DO UPDATE
SET name = EXCLUDED.name
RETURNING id
"#,
)
.bind(&item.warehouse_code)
.bind(&item.warehouse_name)
.fetch_one(&mut **tx)
.await
.map_err(db)?;
Ok(id)
}
async fn upsert_article(
tx: &mut Transaction<'_, Postgres>,
item: &SyncDeliveryItem,
fallback_warehouse_id: Uuid,
) -> Result<Uuid, ApplicationError> {
// Optional: explizites Default-Lager aus dem ERP. Wenn nicht
// geliefert, nehmen wir das Lager dieser Position als Default — das
// ist eine pragmatische Wahl, die wir später korrigieren können.
let default_warehouse_id = if let Some(code) = &item.article_default_warehouse_code {
sqlx::query_scalar::<_, Uuid>("SELECT id FROM warehouses WHERE code = $1")
.bind(code)
.fetch_optional(&mut **tx)
.await
.map_err(db)?
.unwrap_or(fallback_warehouse_id)
} else {
fallback_warehouse_id
};
let id: Uuid = sqlx::query_scalar(
r#"
INSERT INTO articles (article_number, name, scannable, default_warehouse_id)
VALUES ($1, $2, $3, $4)
ON CONFLICT (article_number) DO UPDATE
SET name = EXCLUDED.name,
scannable = EXCLUDED.scannable,
default_warehouse_id = EXCLUDED.default_warehouse_id
RETURNING id
"#,
)
.bind(&item.article_number)
.bind(&item.article_name)
.bind(item.article_scannable)
.bind(default_warehouse_id)
.fetch_one(&mut **tx)
.await
.map_err(db)?;
Ok(id)
}
async fn upsert_customer(
tx: &mut Transaction<'_, Postgres>,
delivery: &SyncDelivery,
) -> Result<Uuid, ApplicationError> {
let id: Uuid = sqlx::query_scalar(
r#"
INSERT INTO customers (
erp_customer_id, name, street, house_number, postal_code, city, country
) VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (erp_customer_id) DO UPDATE SET
name = EXCLUDED.name,
street = EXCLUDED.street,
house_number = EXCLUDED.house_number,
postal_code = EXCLUDED.postal_code,
city = EXCLUDED.city,
country = EXCLUDED.country
RETURNING id
"#,
)
.bind(delivery.erp_customer_id)
.bind(&delivery.customer_name)
.bind(&delivery.customer_address.street)
.bind(&delivery.customer_address.house_number)
.bind(&delivery.customer_address.postal_code)
.bind(&delivery.customer_address.city)
.bind(&delivery.customer_address.country)
.fetch_one(&mut **tx)
.await
.map_err(db)?;
Ok(id)
}
async fn upsert_delivery(
tx: &mut Transaction<'_, Postgres>,
tour_id: Uuid,
customer_id: Uuid,
delivery: &SyncDelivery,
) -> Result<Uuid, ApplicationError> {
// Payment-Method-Code → UUID auflösen. Fallback `"cash"` falls vom
// ERP nichts gekommen ist — `"cash"` ist Default-Stamm aus
// Migration 0008 und damit garantiert vorhanden.
let payment_code = delivery
.payment_method_code
.as_deref()
.unwrap_or("cash");
let payment_method_id: Uuid = sqlx::query_scalar(
"SELECT id FROM payment_methods WHERE code = $1",
)
.bind(payment_code)
.fetch_optional(&mut **tx)
.await
.map_err(db)?
.ok_or_else(|| {
ApplicationError::Validation(format!(
"unknown payment method code '{payment_code}'"
))
})?;
let id: Uuid = sqlx::query_scalar(
r#"
INSERT INTO deliveries (
tour_id, erp_belegart_id, erp_belegart_code, erp_belegart_name,
erp_belegnummer, customer_id,
snap_street, snap_house_number, snap_postal_code, snap_city, snap_country,
sort_order, desired_time, special_agreements,
prepaid_amount, payment_method_id
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16)
ON CONFLICT (erp_belegart_id, erp_belegnummer) DO UPDATE SET
tour_id = EXCLUDED.tour_id,
erp_belegart_code = EXCLUDED.erp_belegart_code,
erp_belegart_name = EXCLUDED.erp_belegart_name,
customer_id = EXCLUDED.customer_id,
snap_street = EXCLUDED.snap_street,
snap_house_number = EXCLUDED.snap_house_number,
snap_postal_code = EXCLUDED.snap_postal_code,
snap_city = EXCLUDED.snap_city,
snap_country = EXCLUDED.snap_country,
sort_order = EXCLUDED.sort_order,
desired_time = EXCLUDED.desired_time,
special_agreements = EXCLUDED.special_agreements,
prepaid_amount = EXCLUDED.prepaid_amount,
payment_method_id = EXCLUDED.payment_method_id
RETURNING id
"#,
)
.bind(tour_id)
.bind(delivery.belegart_id)
.bind(delivery.belegart_code.as_deref())
.bind(delivery.belegart_name.as_deref())
.bind(&delivery.belegnummer)
.bind(customer_id)
.bind(&delivery.delivery_address.street)
.bind(&delivery.delivery_address.house_number)
.bind(&delivery.delivery_address.postal_code)
.bind(&delivery.delivery_address.city)
.bind(&delivery.delivery_address.country)
.bind(delivery.sort_order)
.bind(delivery.desired_time.as_deref())
.bind(delivery.special_agreements.as_deref())
.bind(delivery.prepaid_amount)
.bind(payment_method_id)
.fetch_one(&mut **tx)
.await
.map_err(db)?;
Ok(id)
}
async fn upsert_delivery_item(
tx: &mut Transaction<'_, Postgres>,
delivery_id: Uuid,
article_id: Uuid,
warehouse_id: Uuid,
item: &SyncDeliveryItem,
) -> Result<(), ApplicationError> {
// Identität: (delivery_id, belegzeilen_nr, komponenten_artikel_nr).
// scan_state-Felder bleiben beim UPDATE bewusst unberührt — das ERP
// weiß nichts über Scans.
sqlx::query(
r#"
INSERT INTO delivery_items (
delivery_id, article_id, required_quantity, warehouse_id,
unit_price, belegzeilen_nr, komponenten_artikel_nr, parent_artikel_nr
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (delivery_id, belegzeilen_nr, komponenten_artikel_nr) DO UPDATE SET
article_id = EXCLUDED.article_id,
required_quantity = EXCLUDED.required_quantity,
warehouse_id = EXCLUDED.warehouse_id,
unit_price = EXCLUDED.unit_price,
parent_artikel_nr = EXCLUDED.parent_artikel_nr
"#,
)
.bind(delivery_id)
.bind(article_id)
.bind(item.required_quantity)
.bind(warehouse_id)
.bind(item.unit_price)
.bind(item.belegzeilen_nr)
.bind(item.komponenten_artikel_nr.as_deref())
.bind(item.parent_artikel_nr.as_deref())
.execute(&mut **tx)
.await
.map_err(db)?;
Ok(())
}
// ===== Kontaktdaten ======================================================
fn role_to_db(role: ContactRole) -> &'static str {
match role {
ContactRole::Header => "header",
ContactRole::Delivery => "delivery",
ContactRole::Billing => "billing",
ContactRole::ContactPerson => "contact_person",
ContactRole::CustomerMaster => "customer_master",
}
}
fn role_from_db(value: &str) -> Result<ContactRole, ApplicationError> {
match value {
"header" => Ok(ContactRole::Header),
"delivery" => Ok(ContactRole::Delivery),
"billing" => Ok(ContactRole::Billing),
"contact_person" => Ok(ContactRole::ContactPerson),
"customer_master" => Ok(ContactRole::CustomerMaster),
other => Err(ApplicationError::Repository(format!(
"unknown contact role in DB: {other}"
))),
}
}
fn kind_to_db(kind: ContactKind) -> &'static str {
match kind {
ContactKind::Phone => "phone",
ContactKind::Mobile => "mobile",
ContactKind::Email => "email",
ContactKind::Web => "web",
}
}
fn kind_from_db(value: &str) -> Result<ContactKind, ApplicationError> {
match value {
"phone" => Ok(ContactKind::Phone),
"mobile" => Ok(ContactKind::Mobile),
"email" => Ok(ContactKind::Email),
"web" => Ok(ContactKind::Web),
other => Err(ApplicationError::Repository(format!(
"unknown contact kind in DB: {other}"
))),
}
}
/// Snapshot-Refresh: vorhandene Sources der Lieferung löschen (Channels
/// fliegen per ON DELETE CASCADE mit), neue einfügen. Idempotent: leerer
/// Input ⇒ Lieferung hat nach dem Aufruf 0 Sources.
async fn replace_contact_sources(
tx: &mut Transaction<'_, Postgres>,
delivery_id: Uuid,
sources: &[SyncContactSource],
) -> Result<(), ApplicationError> {
sqlx::query("DELETE FROM delivery_contact_sources WHERE delivery_id = $1")
.bind(delivery_id)
.execute(&mut **tx)
.await
.map_err(db)?;
for src in sources {
let source_id: Uuid = sqlx::query_scalar(
r#"
INSERT INTO delivery_contact_sources (
delivery_id, role, anrede, titel, name1, name2, name3,
abteilung, funktion
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
RETURNING id
"#,
)
.bind(delivery_id)
.bind(role_to_db(src.role))
.bind(src.anrede.as_deref())
.bind(src.titel.as_deref())
.bind(src.name1.as_deref())
.bind(src.name2.as_deref())
.bind(src.name3.as_deref())
.bind(src.abteilung.as_deref())
.bind(src.funktion.as_deref())
.fetch_one(&mut **tx)
.await
.map_err(db)?;
for ch in &src.channels {
sqlx::query(
r#"
INSERT INTO delivery_contact_channels (
source_id, kind, position, value
) VALUES ($1, $2, $3, $4)
"#,
)
.bind(source_id)
.bind(kind_to_db(ch.kind))
.bind(ch.position)
.bind(&ch.value)
.execute(&mut **tx)
.await
.map_err(db)?;
}
}
Ok(())
}