// /services/bankStatementService.ts import axios from "axios" import { centralServicesClient } from "../push-server.client" import dayjs from "dayjs" import utc from "dayjs/plugin/utc.js" import {secrets} from "../../utils/secrets" import {FastifyInstance} from "fastify" // Drizzle imports import { bankaccounts, bankstatements, } from "../../../db/schema" import { eq, and, isNull, } from "drizzle-orm" dayjs.extend(utc) interface BalanceAmount { amount: string currency: string } interface BookedTransaction { bookingDate: string valueDate: string internalTransactionId: string transactionAmount: { amount: string; currency: string } creditorAccount?: { iban?: string } creditorName?: string debtorAccount?: { iban?: string } debtorName?: string remittanceInformationUnstructured?: string remittanceInformationStructured?: string remittanceInformationStructuredArray?: string[] additionalInformation?: string } interface TransactionsResponse { transactions: { booked: BookedTransaction[] } } const normalizeDate = (val: any) => { if (!val) return null const d = new Date(val) return isNaN(d.getTime()) ? null : d } export function bankStatementService(server: FastifyInstance) { let accessToken: string | null = null const useCentralBanking = Boolean(secrets.FEDEO_CENTRAL_SERVICES_ENABLED && centralServicesClient.configured()) // ----------------------------------------------- // ✔ TOKEN LADEN // ----------------------------------------------- const getToken = async () => { if (useCentralBanking) return console.log("Fetching GoCardless token…") const response = await axios.post( `${secrets.GOCARDLESS_BASE_URL}/token/new/`, { secret_id: secrets.GOCARDLESS_SECRET_ID, secret_key: secrets.GOCARDLESS_SECRET_KEY, } ) accessToken = response.data.access } // ----------------------------------------------- // ✔ Salden laden // ----------------------------------------------- const getBalanceData = async (accountId: string): Promise => { try { if (useCentralBanking) return await centralServicesClient.getBankingBalances(accountId) const {data} = await axios.get( `${secrets.GOCARDLESS_BASE_URL}/accounts/${accountId}/balances`, { headers: { Authorization: `Bearer ${accessToken}`, Accept: "application/json", }, } ) return data } catch (err: any) { server.log.error(err.response?.data ?? err) const expired = err.response?.data?.summary?.includes("expired") || err.response?.data?.detail?.includes("expired") if (expired) { await server.db .update(bankaccounts) .set({expired: true}) .where(eq(bankaccounts.accountId, accountId)) } return false } } // ----------------------------------------------- // ✔ Transaktionen laden // ----------------------------------------------- const getTransactionData = async (accountId: string) => { try { if (useCentralBanking) { const data = await centralServicesClient.getBankingTransactions(accountId) return data.transactions.booked } const {data} = await axios.get( `${secrets.GOCARDLESS_BASE_URL}/accounts/${accountId}/transactions`, { headers: { Authorization: `Bearer ${accessToken}`, Accept: "application/json", }, } ) return data.transactions.booked } catch (err: any) { server.log.error(err.response?.data ?? err) return null } } // ----------------------------------------------- // ✔ Haupt-Sync-Prozess // ----------------------------------------------- const syncAccounts = async (tenantId:number) => { try { console.log("Starting account sync…") // 🟦 DB: Aktive Accounts const accounts = await server.db .select() .from(bankaccounts) .where(and(eq(bankaccounts.expired, false),eq(bankaccounts.tenant, tenantId))) if (!accounts.length) return const allNewTransactions: any[] = [] for (const account of accounts) { // --------------------------- // 1. BALANCE SYNC // --------------------------- const balData = await getBalanceData(account.accountId) if (balData === false) break if (balData) { const closing = balData.balances.find( (i: any) => i.balanceType === "closingBooked" ) const bookedBal = Number(closing.balanceAmount.amount) await server.db .update(bankaccounts) .set({balance: bookedBal}) .where(eq(bankaccounts.id, account.id)) } // --------------------------- // 2. TRANSACTIONS // --------------------------- let transactions = await getTransactionData(account.accountId) if (!transactions) continue //@ts-ignore transactions = transactions.map((item) => ({ account: account.id, date: normalizeDate(item.bookingDate), credIban: item.creditorAccount?.iban ?? null, credName: item.creditorName ?? null, text: ` ${item.remittanceInformationUnstructured ?? ""} ${item.remittanceInformationStructured ?? ""} ${item.additionalInformation ?? ""} ${item.remittanceInformationStructuredArray?.join("") ?? ""} `.trim(), amount: Number(item.transactionAmount.amount), tenant: account.tenant, debIban: item.debtorAccount?.iban ?? null, debName: item.debtorName ?? null, gocardlessId: item.internalTransactionId, currency: item.transactionAmount.currency, valueDate: normalizeDate(item.valueDate), })) // Existierende Statements laden const existing = await server.db .select({gocardlessId: bankstatements.gocardlessId}) .from(bankstatements) .where(eq(bankstatements.tenant, account.tenant)) const filtered = transactions.filter( //@ts-ignore (tx) => !existing.some((x) => x.gocardlessId === tx.gocardlessId) ) allNewTransactions.push(...filtered) } // --------------------------- // 3. NEW TRANSACTIONS → DB // --------------------------- if (allNewTransactions.length > 0) { await server.db.insert(bankstatements).values(allNewTransactions) const affectedAccounts = [ ...new Set(allNewTransactions.map((t) => t.account)), ] const normalizeDate = (val: any) => { if (!val) return null const d = new Date(val) return isNaN(d.getTime()) ? null : d } for (const accId of affectedAccounts) { await server.db .update(bankaccounts) //@ts-ignore .set({syncedAt: normalizeDate(dayjs())}) .where(eq(bankaccounts.id, accId)) } } console.log("Bank statement sync completed.") } catch (error) { console.error(error) } } return { run: async (tenant) => { await getToken() await syncAccounts(tenant) console.log("Service: Bankstatement sync finished") } } }