from __future__ import annotations import copy import datetime from pprint import pformat from typing import TYPE_CHECKING, Any, cast import polars as pl import sqlalchemy as sql from dopt_basics.result_pattern import wrap_result from wce_crm import db from wce_crm.constants import TIMEZONE_CEST from wce_crm.data_models import ( Beratungsgespraech_Einzelgespraech, Beratungsgespraech_Vorgang, Initrec, ) from wce_crm.logging import logger_back as logger from wce_crm.types import ( CompanyInfo, CompanyProfileConsultationEntry, CompanyProfileConsultations, ConsultingType, ContactPersonInfo, InitRecType, MainPageEntry, ) if TYPE_CHECKING: from wce_crm.types import ConsId, ExtAnId, ExtMaId, RecId def _transform_for_gui_output( data: pl.DataFrame, ) -> pl.DataFrame: q = ( data.lazy() .with_columns( pl.col(pl.Datetime).dt.to_string("%d.%m.%Y"), pl.col(pl.Date).dt.to_string("%d.%m.%Y"), pl.when(pl.col(pl.Boolean)) .then(pl.lit("Ja")) .otherwise(pl.lit("Nein")) .name.keep(), ) .with_columns(pl.all().cast(pl.String)) ) return q.collect() def initrec_comp_search_choices() -> tuple[tuple[str, int], ...]: # TODO no reload functionality logger.debug("[Call backend] comp_search_choices") q = db.DF_CRM_MASTER.lazy() counter = pl.int_range(0, pl.len()).over(pl.col.ma_unternehmensname) q = q.with_columns( dedupl=pl.when(counter == 0) .then(pl.col.ma_unternehmensname) .otherwise(pl.format("{} ({})", pl.col.ma_unternehmensname, counter)) ) df = q.collect() return tuple(zip(df["dedupl"], df["ma_id"])) def initrec_comp_search_get_info( ma_id: ExtMaId, ) -> CompanyInfo: logger.debug("[Call backend] comp_search_get_info") df = db.DF_CRM_MASTER.filter(pl.col.ma_id == ma_id) if df.height > 1 or df.height == 0: raise ValueError(f"Größe des zurückgelieferten Datenpakets ungültig: {df.height}") df = _transform_for_gui_output(df) return cast(CompanyInfo, df.row(0, named=True)) def initrec_comp_contact_person_search_choices( ma_id: ExtMaId | None, use_both_names: bool, ) -> tuple[tuple[str, int], ...]: # TODO no reload functionality logger.debug("[Call backend] contact_person_search_choices") q = db.DF_CONTACT_PERSON.lazy() if ma_id is not None: q = q.filter(pl.col.ma_id == ma_id) dedupl_col = pl.col.an_nachname if use_both_names: q = q.with_columns( name_search=(pl.format("{}, {}", pl.col.an_nachname, pl.col.an_vorname)) ) dedupl_col = pl.col.name_search counter = pl.int_range(0, pl.len()).over(dedupl_col) q = q.with_columns( dedupl=pl.when(counter == 0) .then(dedupl_col) .otherwise(pl.format("{} ({})", dedupl_col, counter)) ) df = q.collect() return tuple(zip(df["dedupl"], df["an_id"])) def initrec_comp_contact_person_search_get_info( an_id: ExtAnId, ) -> ContactPersonInfo: logger.debug("[Call backend] contact_person_search_get_info") df = db.DF_CONTACT_PERSON.filter(pl.col.an_id == an_id) if df.height > 1 or df.height == 0: raise ValueError(f"Größe des zurückgelieferten Datenpakets ungültig: {df.height}") df = _transform_for_gui_output(df) return cast(ContactPersonInfo, df.row(0, named=True)) # // internals def initrec_company_insert_initial_recording( data: dict[str, Any], ) -> RecId: logger.debug("[Call backend] insert_initial_recording") stmt = db.grunderfassung_unternehmen.insert() with db.ENGINE.begin() as conn: ret = conn.execute(stmt, data) if ret.rowcount == 0: raise IOError("Entry was not inserted correctly") prim_keys = ret.inserted_primary_key assert prim_keys return prim_keys[0] def initrec_company_update_initial_recording( id_: RecId, data: dict[str, Any], ) -> None: logger.debug("[Call backend] update_initial_recording") stmt = db.grunderfassung_unternehmen.update().where( db.grunderfassung_unternehmen.c.un_id == id_ ) with db.ENGINE.begin() as conn: conn.execute(stmt, data) def initrec_company_to_db( auto_form_data: Initrec, ) -> Initrec: logger.debug("[AutoForm -- backend] Call database saving routine...") dump_data = copy.deepcopy(auto_form_data.db_data) dump_data["geloescht"] = auto_form_data.geloescht with db.ENGINE.begin() as conn: if auto_form_data.rec_id is None: logger.debug("[AutoForm -- backend] Insert...") stmt = db.grunderfassung_unternehmen.insert().returning( db.grunderfassung_unternehmen.c.un_id, db.grunderfassung_unternehmen.c.Metadaten_aktualisierung, db.grunderfassung_unternehmen.c.geloescht, ) ret = conn.execute(stmt, dump_data) if ret.rowcount == 0: raise IOError("Entry was not inserted correctly") from_db = ret.mappings().fetchall() assert len(from_db) == 1, "expected excactly one returned row" data_from_db = from_db[0] auto_form_data.rec_id = data_from_db["un_id"] auto_form_data.geloescht = data_from_db["geloescht"] auto_form_data.db_data["Metadaten_aktualisierung"] = data_from_db[ "Metadaten_aktualisierung" ] logger.debug("[AutoForm -- backend] Inserted InitRec successfully") else: logger.debug("[AutoForm -- backend] Update...") stmt = ( db.grunderfassung_unternehmen.update() .where(db.grunderfassung_unternehmen.c.un_id == auto_form_data.rec_id) .returning( db.grunderfassung_unternehmen.c.Metadaten_aktualisierung, db.grunderfassung_unternehmen.c.geloescht, ) ) ret = conn.execute(stmt, dump_data) from_db = ret.mappings().fetchall() assert len(from_db) == 1, "expected excatly one returned row" data_from_db = from_db[0] auto_form_data.db_data["Metadaten_aktualisierung"] = data_from_db[ "Metadaten_aktualisierung" ] logger.debug( "[AutoForm -- backend] Updated InitRec with ID %d successfully", auto_form_data.rec_id, ) return auto_form_data def initrec_company_from_db( id_: RecId, ) -> Initrec: logger.debug("[AutoForm -- backend] Call database loading routine...") stmt = db.grunderfassung_unternehmen.select().where( db.grunderfassung_unternehmen.c.un_id == id_ ) with db.ENGINE.connect() as conn: ret = conn.execute(stmt) from_db = ret.mappings().all() assert len(from_db) == 1, "not excactly one company initial recording obtained" data_from_db = dict(from_db[0]) geloescht = data_from_db["geloescht"] del data_from_db["geloescht"] return Initrec( rec_id=id_, geloescht=geloescht, db_data=data_from_db, ) def initrec_person_to_db( auto_form_data: Initrec, ) -> Initrec: logger.debug("[AutoForm -- backend] Call database saving routine...") dump_data = copy.deepcopy(auto_form_data.db_data) dump_data["geloescht"] = auto_form_data.geloescht with db.ENGINE.begin() as conn: if auto_form_data.rec_id is None: logger.debug("[AutoForm -- backend] Insert...") stmt = db.grunderfassung_personen.insert().returning( db.grunderfassung_personen.c.pers_id, db.grunderfassung_personen.c.Metadaten_aktualisierung, db.grunderfassung_personen.c.geloescht, ) ret = conn.execute(stmt, dump_data) if ret.rowcount == 0: raise IOError("Entry was not inserted correctly") from_db = ret.mappings().fetchall() assert len(from_db) == 1, "expected excactly one returned row" data_from_db = from_db[0] auto_form_data.rec_id = data_from_db["un_id"] auto_form_data.geloescht = data_from_db["geloescht"] auto_form_data.db_data["Metadaten_aktualisierung"] = data_from_db[ "Metadaten_aktualisierung" ] logger.debug("[AutoForm -- backend] Inserted InitRec successfully") else: logger.debug("[AutoForm -- backend] Update...") stmt = ( db.grunderfassung_personen.update() .where(db.grunderfassung_personen.c.pers_id == auto_form_data.rec_id) .returning( db.grunderfassung_personen.c.Metadaten_aktualisierung, db.grunderfassung_personen.c.geloescht, ) ) ret = conn.execute(stmt, dump_data) from_db = ret.mappings().fetchall() assert len(from_db) == 1, "expected excatly one returned row" data_from_db = from_db[0] auto_form_data.db_data["Metadaten_aktualisierung"] = data_from_db[ "Metadaten_aktualisierung" ] logger.debug( "[AutoForm -- backend] Updated InitRec with ID %d successfully", auto_form_data.rec_id, ) return auto_form_data def initrec_person_from_db( id_: RecId, ) -> Initrec: logger.debug("[AutoForm -- backend] Call database loading routine...") stmt = db.grunderfassung_personen.select().where( db.grunderfassung_personen.c.pers_id == id_ ) with db.ENGINE.connect() as conn: ret = conn.execute(stmt) from_db = ret.mappings().all() assert len(from_db) == 1, "not excactly one company initial recording obtained" data_from_db = dict(from_db[0]) geloescht = data_from_db["geloescht"] del data_from_db["geloescht"] return Initrec( rec_id=id_, geloescht=geloescht, db_data=data_from_db, ) # def initrec_person_to_db( # auto_form_data: Initrec, # ) -> Initrec_FromDb: # logger.debug("[AutoForm -- backend] Call database saving routine...") # dump_data = copy.deepcopy(auto_form_data.form_data) # dump_data["geloescht"] = auto_form_data.geloescht # return_from_db: Initrec_FromDb # with db.ENGINE.begin() as conn: # if auto_form_data.rec_id is None: # # insert new "Vorgang" # logger.debug("[AutoForm -- backend] Insert...") # stmt = db.grunderfassung_personen.insert().returning( # db.grunderfassung_personen.c.pers_id, # db.grunderfassung_personen.c.Metadaten_aktualisierung, # db.grunderfassung_personen.c.geloescht, # ) # ret = conn.execute(stmt, dump_data) # if ret.rowcount == 0: # raise IOError("Entry was not inserted correctly") # from_db = ret.mappings().fetchall() # assert len(from_db) == 1, "expected excactly one returned row" # data_from_db = from_db[0] # return_from_db = Initrec_FromDb( # rec_id=data_from_db["un_id"], # geloescht=data_from_db["geloescht"], # Metadaten_aktualisierung=data_from_db["Metadaten_aktualisierung"], # ) # logger.debug("[AutoForm -- backend] Inserted InitRec successfully") # else: # logger.debug("[AutoForm -- backend] Update...") # stmt = ( # db.beratung_vorgang.update() # .where(db.grunderfassung_personen.c.pers_id == auto_form_data.rec_id) # .returning( # db.grunderfassung_personen.c.Metadaten_aktualisierung, # db.grunderfassung_personen.c.geloescht, # ) # ) # ret = conn.execute(stmt, dump_data) # from_db = ret.mappings().fetchall() # assert len(from_db) == 1, "expected excatly one returned row" # data_from_db = from_db[0] # return_from_db = Initrec_FromDb( # rec_id=auto_form_data.rec_id, # geloescht=data_from_db["geloescht"], # Metadaten_aktualisierung=data_from_db["Metadaten_aktualisierung"], # ) # logger.debug( # "[AutoForm -- backend] Updated InitRec with ID %d successfully", # auto_form_data.rec_id, # ) # return return_from_db def initrec_company_get_initial_recording( id_: RecId, ) -> dict[str, Any]: logger.debug("[Call backend] get_initial_recording") stmt = db.grunderfassung_unternehmen.select().where( db.grunderfassung_unternehmen.c.un_id == id_ ) with db.ENGINE.connect() as conn: ret = conn.execute(stmt) results = ret.mappings().all() if not results: raise KeyError(f"Database ID {id_} not found") assert len(results) == 1, "more than one company initial recording obtained" row = results[0] assert row, "row was not obtained" return dict(row) def initrec_company_delete_initial_recording( id_: RecId, ) -> None: logger.debug("[Call backend] delete_initial_recording") stmt = db.grunderfassung_unternehmen.delete().where( db.grunderfassung_unternehmen.c.un_id == id_ ) with db.ENGINE.begin() as conn: ret = conn.execute(stmt) if ret.rowcount == 0: raise KeyError(f"Database ID {id_} not found for deletion") def initrec_person_insert_initial_recording( data: dict[str, Any], ) -> RecId: logger.debug("[Call backend] insert_initial_recording") stmt = db.grunderfassung_personen.insert() with db.ENGINE.begin() as conn: ret = conn.execute(stmt, data) if ret.rowcount == 0: raise IOError("Entry was not inserted correctly") prim_keys = ret.inserted_primary_key assert prim_keys return prim_keys[0] def initrec_person_update_initial_recording( id_: RecId, data: dict[str, Any], ) -> None: logger.debug("[Call backend] update_initial_recording") stmt = db.grunderfassung_personen.update().where( db.grunderfassung_personen.c.pers_id == id_ ) with db.ENGINE.begin() as conn: conn.execute(stmt, data) def initrec_person_get_initial_recording( id_: RecId, ) -> dict[str, Any]: logger.debug("[Call backend] get_initial_recording person") stmt = db.grunderfassung_personen.select().where( db.grunderfassung_personen.c.pers_id == id_ ) with db.ENGINE.connect() as conn: ret = conn.execute(stmt) results = ret.mappings().all() if not results: raise KeyError(f"Database ID {id_} not found") assert len(results) == 1, "more than one person initial recording obtained" row = results[0] assert row, "row was not obtained" return dict(row) def initrec_person_delete_initial_recording( id_: RecId, ) -> None: logger.debug("[Call backend] delete_initial_recording") stmt = db.grunderfassung_personen.delete().where( db.grunderfassung_personen.c.pers_id == id_ ) with db.ENGINE.begin() as conn: ret = conn.execute(stmt) if ret.rowcount == 0: raise KeyError(f"Database ID {id_} not found for deletion") @wrap_result(10) def page_consulting_to_db( consultation_data: Beratungsgespraech_Vorgang, ) -> Beratungsgespraech_Vorgang: logger.debug("[Consulting Page] Call database saving routine...") with db.ENGINE.begin() as conn: if consultation_data.vorgang_id is None: # insert new "Vorgang" insert_data_process = consultation_data.model_dump( exclude={"vorgang_id", "beratungen"} ) logger.debug( "[Consulting Page] Call insert 'Vorgang' with data:\n%s", pformat(insert_data_process), ) stmt = db.beratung_vorgang.insert() ret = conn.execute(stmt, insert_data_process) if ret.rowcount == 0: raise IOError("Entry was not inserted correctly") prim_keys = ret.inserted_primary_key assert prim_keys consultation_data.vorgang_id = cast("ConsId", prim_keys[0]) logger.debug("[Consulting Page] Inserted 'Vorgang' successfully") else: logger.debug( "[Consulting Page] VorgangID already set. ID: %d. Update...", consultation_data.vorgang_id, ) update_data_process = consultation_data.model_dump(exclude={"beratungen"}) stmt = db.beratung_vorgang.update().where( db.beratung_vorgang.c.vorgang_id == consultation_data.vorgang_id ) conn.execute(stmt, update_data_process) rows_for_db_insert: list[dict[str, Any]] = [] rows_for_db_update: list[dict[str, Any]] = [] cons_sessions_inserted: list[Beratungsgespraech_Einzelgespraech] = [] for cons_session in consultation_data.beratungen: cons_session.vorgang_id = consultation_data.vorgang_id if cons_session.beratung_id is None: row_data = cons_session.model_dump(exclude={"beratung_id"}) rows_for_db_insert.append(row_data) cons_sessions_inserted.append(cons_session) else: row_data = cons_session.model_dump() # new bind param to avoid name clashes row_data["b_beratung_id"] = row_data["beratung_id"] del row_data["beratung_id"] rows_for_db_update.append(row_data) if rows_for_db_update: # ... update logger.debug( "[Consulting Page] Call update for sessions:\n%s", pformat(rows_for_db_update) ) stmt = db.beratung_einzelberatung.update().where( db.beratung_einzelberatung.c.beratung_id == sql.bindparam("b_beratung_id") ) conn.execute(stmt, rows_for_db_update) if rows_for_db_insert: # ... insert logger.debug( "[Consulting Page] Call insert for sessions:\n%s", pformat(rows_for_db_insert) ) stmt = db.beratung_einzelberatung.insert().returning( db.beratung_einzelberatung.c.beratung_id ) res = conn.execute(stmt, rows_for_db_insert) new_cons_session_ids = cast(list[int], [row[0] for row in res.fetchall()]) assert len(cons_sessions_inserted) == len(new_cons_session_ids) for cons_session, new_id in zip(cons_sessions_inserted, new_cons_session_ids): cons_session.beratung_id = new_id return consultation_data @wrap_result(11) def page_consulting_from_db( cons_id: ConsId, ) -> Beratungsgespraech_Vorgang: # get "Vorgang" and all associated sessions to instantiate Pydantic data model logger.debug("[Consulting Page] Call database reading routine...") with db.ENGINE.connect() as conn: # get the consultation process stmt = db.beratung_vorgang.select().where(db.beratung_vorgang.c.vorgang_id == cons_id) ret = conn.execute(stmt) results = ret.mappings().all() if not results: raise KeyError(f"Database ID {cons_id} not found") assert len(results) == 1, "more than one consulting process obtained" consultation_data_db = results[0] consultation_data = Beratungsgespraech_Vorgang( vorgang_id=consultation_data_db["vorgang_id"], un_id=consultation_data_db["un_id"], pers_id=consultation_data_db["pers_id"], titel=consultation_data_db["titel"], beratungs_typ=consultation_data_db["beratungs_typ"], erstellt=consultation_data_db["erstellt"], aktualisiert=consultation_data_db["aktualisiert"], geloescht=consultation_data_db["geloescht"], beratungen=[], ) if consultation_data.geloescht is not None: logger.info( ( "[Consulting Page] Obtained entry for ID=%d. This entry was marked " "as deleted." ), cons_id, ) # get all consultation sessions of this process, only the ones # which are not marked as deleted stmt = db.beratung_einzelberatung.select().where( db.beratung_einzelberatung.c.vorgang_id == cons_id, db.beratung_einzelberatung.c.geloescht.is_(None), ) ret = conn.execute(stmt) # empty results possible if ret.rowcount == 0: logger.debug("[Consulting Page] No sessions, return directly...") return consultation_data cons_sessions = ret.mappings() for session in cons_sessions: cons_session_pydantic = Beratungsgespraech_Einzelgespraech( beratung_id=session["beratung_id"], vorgang_id=session["vorgang_id"], nutzer_id=session["nutzer_id"], nutzer_name=session["nutzer_name"], zeitstempel=session["zeitstempel"], ansprechpartner=session["ansprechpartner"], kommunikationsweg=session["kommunikationsweg"], thema_crm_matrix=session["thema_crm_matrix"], anmerkungen=session["anmerkungen"], rueckmeldung=session["rueckmeldung"], erstellt=session["erstellt"], aktualisiert=session["aktualisiert"], ) consultation_data.beratungen.append(cons_session_pydantic) logger.debug("[Consulting Page] Returning with sessions attached...") return consultation_data # // consulting page interaction def companyprofile_page_get_consultations( un_id: RecId, ) -> CompanyProfileConsultations: logger.debug("[Call backend] _companyprofile_page_get_consultations") stmt = sql.select( db.beratung_vorgang.c.vorgang_id, db.beratung_vorgang.c.aktualisiert, db.beratung_vorgang.c.titel, db.beratung_vorgang.c.beratungs_typ, ).where(db.beratung_vorgang.c.un_id == un_id) with db.ENGINE.connect() as conn: res = conn.execute(stmt) cons_entries_pauschal: list[CompanyProfileConsultationEntry] = [] cons_entries_individual: list[CompanyProfileConsultationEntry] = [] for entry in res.mappings(): cons_id = entry["vorgang_id"] assert cons_id, "no VorgangID defined" datetime_updated = cast(datetime.datetime, entry["aktualisiert"]) datetime_updated = datetime_updated.astimezone(TIMEZONE_CEST) base_title = entry["titel"] _cons_type = entry["beratungs_typ"] cons_type = ConsultingType(_cons_type) con_entry = CompanyProfileConsultationEntry( cons_id=cons_id, title=base_title, date_updated=datetime_updated, cons_type=cons_type, ) if cons_type is ConsultingType.PAUSCHAL: cons_entries_pauschal.append(con_entry) elif cons_type is ConsultingType.INDIVIDUAL: cons_entries_individual.append(con_entry) else: raise TypeError(f"Unknown consulting type: {_cons_type}") cons_entries_pauschal.sort(key=lambda x: x.date_updated, reverse=True) cons_entries_individual.sort(key=lambda x: x.date_updated, reverse=True) return CompanyProfileConsultations( pauschal=cons_entries_pauschal, individual=cons_entries_individual, ) # // main page interaction def _main_page_get_company_list() -> list[MainPageEntry]: logger.debug("[Call backend] get_company_list") stmt = sql.select( db.grunderfassung_unternehmen.c.un_id, db.grunderfassung_unternehmen.c.Partnersuche__un_suche, db.grunderfassung_unternehmen.c.Metadaten_aktualisierung, ) with db.ENGINE.connect() as conn: res = conn.execute(stmt) main_page_companies: list[MainPageEntry] = [] for entry in res.mappings(): rec_id = entry["un_id"] assert rec_id, "no RecID defined" ma_id_external = entry["Partnersuche__un_suche"] assert ma_id_external is not None, "external MA ID is NULL" datetime_akt = cast(datetime.datetime, entry["Metadaten_aktualisierung"]) datetime_akt = datetime_akt.astimezone(TIMEZONE_CEST) comp_info = initrec_comp_search_get_info(ma_id_external) display_name = comp_info["ma_unternehmensname"] main_page_companies.append( MainPageEntry( rec_id=rec_id, display_name=display_name, Metadaten_aktualisierung=datetime_akt, type=InitRecType.COMPANY, ) ) return main_page_companies def _main_page_get_person_list() -> list[MainPageEntry]: logger.debug("[Call backend] get_company_list") stmt = sql.select( db.grunderfassung_personen.c.pers_id, db.grunderfassung_personen.c.Stammdaten__vorname, db.grunderfassung_personen.c.Stammdaten__name, db.grunderfassung_personen.c.Metadaten_aktualisierung, ) with db.ENGINE.connect() as conn: res = conn.execute(stmt) main_page_persons: list[MainPageEntry] = [] for entry in res.mappings(): rec_id = entry["pers_id"] assert rec_id, "no RecID defined" datetime_akt = cast(datetime.datetime, entry["Metadaten_aktualisierung"]) datetime_akt = datetime_akt.astimezone() surname: str = entry["Stammdaten__name"] first_name: str = entry["Stammdaten__vorname"] names_to_join = [n for n in (first_name, surname) if n] display_name = " ".join(names_to_join) main_page_persons.append( MainPageEntry( rec_id=rec_id, display_name=display_name, Metadaten_aktualisierung=datetime_akt, type=InitRecType.PERSON, ) ) return main_page_persons def main_page_get_entries() -> list[MainPageEntry]: companies = _main_page_get_company_list() persons = _main_page_get_person_list() all_entries = companies + persons all_entries.sort(key=lambda x: x.Metadaten_aktualisierung, reverse=True) return all_entries