1
0
Fork 0
mirror of https://github.com/maybe-finance/maybe.git synced 2025-08-09 07:25:19 +02:00

refactor: remove InvestmentTransaction category DB generation

This commit is contained in:
Michaël De Boey 2024-01-20 22:11:26 +01:00
parent 4b007713a3
commit 8276c56f74
No known key found for this signature in database
9 changed files with 236 additions and 39 deletions

View file

@ -1,4 +1,5 @@
import type { User } from '@prisma/client' import type { User } from '@prisma/client'
import { InvestmentTransactionCategory } from '@prisma/client'
import { PrismaClient } from '@prisma/client' import { PrismaClient } from '@prisma/client'
import { createLogger, transports } from 'winston' import { createLogger, transports } from 'winston'
import { DateTime } from 'luxon' import { DateTime } from 'luxon'
@ -131,6 +132,7 @@ describe('balance sync strategies', () => {
quantity: 10, quantity: 10,
price: 10, price: 10,
plaidType: 'buy', plaidType: 'buy',
category: InvestmentTransactionCategory.buy,
}, },
{ {
date: DateTime.fromISO('2023-02-04').toJSDate(), date: DateTime.fromISO('2023-02-04').toJSDate(),
@ -140,6 +142,7 @@ describe('balance sync strategies', () => {
quantity: 5, quantity: 5,
price: 10, price: 10,
plaidType: 'sell', plaidType: 'sell',
category: InvestmentTransactionCategory.sell,
}, },
{ {
date: DateTime.fromISO('2023-02-04').toJSDate(), date: DateTime.fromISO('2023-02-04').toJSDate(),
@ -147,6 +150,7 @@ describe('balance sync strategies', () => {
amount: 50, amount: 50,
quantity: 50, quantity: 50,
price: 1, price: 1,
category: InvestmentTransactionCategory.other,
}, },
], ],
}, },

View file

@ -1,5 +1,5 @@
import type { User } from '@prisma/client' import type { User } from '@prisma/client'
import { Prisma, PrismaClient } from '@prisma/client' import { InvestmentTransactionCategory, Prisma, PrismaClient } from '@prisma/client'
import { createLogger, transports } from 'winston' import { createLogger, transports } from 'winston'
import { DateTime } from 'luxon' import { DateTime } from 'luxon'
import type { import type {
@ -307,6 +307,7 @@ describe('insight service', () => {
price: 100, price: 100,
plaidType: 'buy', plaidType: 'buy',
plaidSubtype: 'buy', plaidSubtype: 'buy',
category: InvestmentTransactionCategory.buy,
}, },
{ {
accountId: account.id, accountId: account.id,
@ -318,6 +319,7 @@ describe('insight service', () => {
price: 200, price: 200,
plaidType: 'buy', plaidType: 'buy',
plaidSubtype: 'buy', plaidSubtype: 'buy',
category: InvestmentTransactionCategory.buy,
}, },
{ {
accountId: account.id, accountId: account.id,
@ -329,6 +331,7 @@ describe('insight service', () => {
price: 0, price: 0,
plaidType: 'cash', plaidType: 'cash',
plaidSubtype: 'dividend', plaidSubtype: 'dividend',
category: InvestmentTransactionCategory.dividend,
}, },
{ {
accountId: account.id, accountId: account.id,
@ -340,6 +343,7 @@ describe('insight service', () => {
price: 0, price: 0,
plaidType: 'cash', plaidType: 'cash',
plaidSubtype: 'dividend', plaidSubtype: 'dividend',
category: InvestmentTransactionCategory.dividend,
}, },
], ],
}) })

View file

@ -1,5 +1,5 @@
import type { PrismaClient, User } from '@prisma/client' import type { PrismaClient, User } from '@prisma/client'
import { Prisma } from '@prisma/client' import { InvestmentTransactionCategory, Prisma } from '@prisma/client'
import _ from 'lodash' import _ from 'lodash'
import { DateTime } from 'luxon' import { DateTime } from 'luxon'
import { parseCsv } from './csv' import { parseCsv } from './csv'
@ -20,6 +20,14 @@ const portfolios: Record<string, Partial<Prisma.AccountUncheckedCreateInput>> =
}, },
} }
const investmentTransactionCategoryByType: Record<string, InvestmentTransactionCategory> = {
BUY: InvestmentTransactionCategory.buy,
SELL: InvestmentTransactionCategory.sell,
DIVIDEND: InvestmentTransactionCategory.dividend,
DEPOSIT: InvestmentTransactionCategory.transfer,
WITHDRAW: InvestmentTransactionCategory.transfer,
}
export async function createTestInvestmentAccount( export async function createTestInvestmentAccount(
prisma: PrismaClient, prisma: PrismaClient,
user: User, user: User,
@ -35,7 +43,7 @@ export async function createTestInvestmentAccount(
join(__dirname, `../test-data/${portfolio}/holdings.csv`) join(__dirname, `../test-data/${portfolio}/holdings.csv`)
) )
const [_deleted, ...securities] = await prisma.$transaction([ const [, ...securities] = await prisma.$transaction([
prisma.security.deleteMany({ prisma.security.deleteMany({
where: { where: {
symbol: { symbol: {
@ -72,7 +80,7 @@ export async function createTestInvestmentAccount(
.value(), .value(),
]) ])
const account = await prisma.account.create({ return prisma.account.create({
data: { data: {
...portfolios[portfolio], ...portfolios[portfolio],
userId: user.id, userId: user.id,
@ -128,12 +136,13 @@ export async function createTestInvestmentAccount(
: it.type === 'SELL' : it.type === 'SELL'
? 'sell' ? 'sell'
: undefined, : undefined,
category:
investmentTransactionCategoryByType[it.type] ??
InvestmentTransactionCategory.other,
} }
}), }),
}, },
}, },
}, },
}) })
return account
} }

View file

@ -69,9 +69,8 @@ export class InvestmentTransactionBalanceSyncStrategy extends BalanceSyncStrateg
WHERE WHERE
it.account_id = ${pAccountId} it.account_id = ${pAccountId}
AND it.date BETWEEN ${pStart} AND now() AND it.date BETWEEN ${pStart} AND now()
AND ( -- filter for transactions that modify a position -- filter for transactions that modify a position
it.plaid_type IN ('buy', 'sell', 'transfer') AND it.category IN ('buy', 'sell', 'transfer')
)
GROUP BY GROUP BY
1, 2 1, 2
) it ON it.security_id = s.id AND it.date = d.date ) it ON it.security_id = s.id AND it.date = d.date

View file

@ -242,11 +242,7 @@ export class AccountQueryService implements IAccountQueryService {
it.account_id = ANY(${pAccountIds}) it.account_id = ANY(${pAccountIds})
AND it.date BETWEEN sd.start_date AND ${pEnd} AND it.date BETWEEN sd.start_date AND ${pEnd}
-- filter for investment_transactions that represent external flows -- filter for investment_transactions that represent external flows
AND ( AND it.category = 'transfer'
(it.plaid_type = 'cash' AND it.plaid_subtype IN ('contribution', 'deposit', 'withdrawal'))
OR (it.plaid_type = 'transfer' AND it.plaid_subtype IN ('transfer'))
OR (it.plaid_type = 'buy' AND it.plaid_subtype IN ('contribution'))
)
GROUP BY GROUP BY
1, 2 1, 2
), external_flow_totals AS ( ), external_flow_totals AS (

View file

@ -312,6 +312,7 @@ export class InsightService implements IInsightService {
{ {
plaidSubtype: 'dividend', plaidSubtype: 'dividend',
}, },
{ category: 'dividend' },
], ],
}, },
}), }),
@ -737,11 +738,7 @@ export class InsightService implements IInsightService {
LEFT JOIN account a ON a.id = it.account_id LEFT JOIN account a ON a.id = it.account_id
WHERE WHERE
it.account_id = ${accountId} it.account_id = ${accountId}
AND ( AND it.category = 'transfer'
(it.plaid_type = 'cash' AND it.plaid_subtype IN ('contribution', 'deposit', 'withdrawal'))
OR (it.plaid_type = 'transfer' AND it.plaid_subtype IN ('transfer', 'send', 'request'))
OR (it.plaid_type = 'buy' AND it.plaid_subtype IN ('contribution'))
)
-- Exclude any contributions made prior to the start date since balances will be 0 -- Exclude any contributions made prior to the start date since balances will be 0
AND (a.start_date is NULL OR it.date >= a.start_date) AND (a.start_date is NULL OR it.date >= a.start_date)
GROUP BY 1 GROUP BY 1

View file

@ -14,7 +14,8 @@ import type {
LiabilitiesObject as PlaidLiabilities, LiabilitiesObject as PlaidLiabilities,
PlaidApi, PlaidApi,
} from 'plaid' } from 'plaid'
import { Prisma } from '@prisma/client' import { InvestmentTransactionSubtype, InvestmentTransactionType } from 'plaid'
import { Prisma, InvestmentTransactionCategory } from '@prisma/client'
import { DateTime } from 'luxon' import { DateTime } from 'luxon'
import _, { chunk } from 'lodash' import _, { chunk } from 'lodash'
import { ErrorUtil, PlaidUtil } from '@maybe-finance/server/shared' import { ErrorUtil, PlaidUtil } from '@maybe-finance/server/shared'
@ -548,7 +549,7 @@ export class PlaidETL implements IETL<Connection, PlaidRawData, PlaidData> {
...chunk(investmentTransactions, 1_000).map( ...chunk(investmentTransactions, 1_000).map(
(chunk) => (chunk) =>
this.prisma.$executeRaw` this.prisma.$executeRaw`
INSERT INTO investment_transaction (account_id, security_id, plaid_investment_transaction_id, date, name, amount, fees, quantity, price, currency_code, plaid_type, plaid_subtype) INSERT INTO investment_transaction (account_id, security_id, plaid_investment_transaction_id, date, name, amount, fees, quantity, price, currency_code, plaid_type, plaid_subtype, category)
VALUES VALUES
${Prisma.join( ${Prisma.join(
chunk.map( chunk.map(
@ -584,7 +585,11 @@ export class PlaidETL implements IETL<Connection, PlaidRawData, PlaidData> {
${DbUtil.toDecimal(price)}, ${DbUtil.toDecimal(price)},
${currencyCode}, ${currencyCode},
${type}, ${type},
${subtype} ${subtype},
${this.getInvestmentTransactionCategoryByPlaidType(
type,
subtype
)}
)` )`
} }
) )
@ -602,6 +607,7 @@ export class PlaidETL implements IETL<Connection, PlaidRawData, PlaidData> {
currency_code = EXCLUDED.currency_code, currency_code = EXCLUDED.currency_code,
plaid_type = EXCLUDED.plaid_type, plaid_type = EXCLUDED.plaid_type,
plaid_subtype = EXCLUDED.plaid_subtype; plaid_subtype = EXCLUDED.plaid_subtype;
category = EXCLUDED.category;
` `
), ),
@ -669,6 +675,63 @@ export class PlaidETL implements IETL<Connection, PlaidRawData, PlaidData> {
] ]
} }
private getInvestmentTransactionCategoryByPlaidType = (
type: InvestmentTransactionType,
subType: InvestmentTransactionSubtype
): InvestmentTransactionCategory => {
if (type === InvestmentTransactionType.Buy) {
return InvestmentTransactionCategory.buy
}
if (type === InvestmentTransactionType.Sell) {
return InvestmentTransactionCategory.sell
}
if (
[
InvestmentTransactionSubtype.Dividend,
InvestmentTransactionSubtype.QualifiedDividend,
InvestmentTransactionSubtype.NonQualifiedDividend,
].includes(subType)
) {
return InvestmentTransactionCategory.dividend
}
if (
[
InvestmentTransactionSubtype.NonResidentTax,
InvestmentTransactionSubtype.Tax,
InvestmentTransactionSubtype.TaxWithheld,
].includes(subType)
) {
return InvestmentTransactionCategory.tax
}
if (
type === InvestmentTransactionType.Fee ||
[
InvestmentTransactionSubtype.AccountFee,
InvestmentTransactionSubtype.LegalFee,
InvestmentTransactionSubtype.ManagementFee,
InvestmentTransactionSubtype.MarginExpense,
InvestmentTransactionSubtype.TransferFee,
InvestmentTransactionSubtype.TrustFee,
].includes(subType)
) {
return InvestmentTransactionCategory.fee
}
if (type === InvestmentTransactionType.Cash) {
return InvestmentTransactionCategory.transfer
}
if (type === InvestmentTransactionType.Cancel) {
return InvestmentTransactionCategory.cancel
}
return InvestmentTransactionCategory.other
}
private async _extractHoldings(accessToken: string) { private async _extractHoldings(accessToken: string) {
try { try {
const { data } = await this.plaid.investmentsHoldingsGet({ access_token: accessToken }) const { data } = await this.plaid.investmentsHoldingsGet({ access_token: accessToken })

View file

@ -0,0 +1,127 @@
-- AlterTable
ALTER TABLE "investment_transaction"
RENAME COLUMN "category" TO "category_old";
DROP VIEW IF EXISTS holdings_enriched;
ALTER TABLE "investment_transaction"
ADD COLUMN "category" "InvestmentTransactionCategory" NOT NULL DEFAULT 'other'::"InvestmentTransactionCategory";
CREATE OR REPLACE VIEW holdings_enriched AS (
SELECT
h.id,
h.account_id,
h.security_id,
h.quantity,
COALESCE(pricing_latest.price_close * h.quantity * COALESCE(s.shares_per_contract, 1), h.value) AS "value",
COALESCE(h.cost_basis, tcb.cost_basis * h.quantity) AS "cost_basis",
COALESCE(h.cost_basis / h.quantity / COALESCE(s.shares_per_contract, 1), tcb.cost_basis) AS "cost_basis_per_share",
pricing_latest.price_close AS "price",
pricing_prev.price_close AS "price_prev",
h.excluded
FROM
holding h
INNER JOIN security s ON s.id = h.security_id
-- latest security pricing
LEFT JOIN LATERAL (
SELECT
price_close
FROM
security_pricing
WHERE
security_id = h.security_id
ORDER BY
date DESC
LIMIT 1
) pricing_latest ON true
-- previous security pricing (for computing daily ∆)
LEFT JOIN LATERAL (
SELECT
price_close
FROM
security_pricing
WHERE
security_id = h.security_id
ORDER BY
date DESC
LIMIT 1
OFFSET 1
) pricing_prev ON true
-- calculate cost basis from transactions
LEFT JOIN (
SELECT
it.account_id,
it.security_id,
SUM(it.quantity * it.price) / SUM(it.quantity) AS cost_basis
FROM
investment_transaction it
WHERE
it.category = 'buy'
AND it.quantity > 0
GROUP BY
it.account_id,
it.security_id
) tcb ON tcb.account_id = h.account_id AND tcb.security_id = s.id
);
CREATE OR REPLACE FUNCTION calculate_return_dietz(p_account_id account.id%type, p_start date, p_end date, out percentage numeric, out amount numeric) AS $$
DECLARE
v_start date := GREATEST(p_start, (SELECT MIN(date) FROM account_balance WHERE account_id = p_account_id));
v_end date := p_end;
v_days int := v_end - v_start;
BEGIN
SELECT
ROUND((b1.balance - b0.balance - flows.net) / NULLIF(b0.balance + flows.weighted, 0), 4) AS "percentage",
b1.balance - b0.balance - flows.net AS "amount"
INTO
percentage, amount
FROM
account a
LEFT JOIN LATERAL (
SELECT
COALESCE(SUM(-fw.flow), 0) AS "net",
COALESCE(SUM(-fw.flow * fw.weight), 0) AS "weighted"
FROM (
SELECT
SUM(it.amount) AS flow,
(v_days - (it.date - v_start))::numeric / v_days AS weight
FROM
investment_transaction it
WHERE
it.account_id = a.id
AND it.date BETWEEN v_start AND v_end
-- filter for investment_transactions that represent external flows
AND it.category = 'transfer'
GROUP BY
it.date
) fw
) flows ON TRUE
LEFT JOIN LATERAL (
SELECT
ab.balance AS "balance"
FROM
account_balance ab
WHERE
ab.account_id = a.id AND ab.date <= v_start
ORDER BY
ab.date DESC
LIMIT 1
) b0 ON TRUE
LEFT JOIN LATERAL (
SELECT
COALESCE(ab.balance, a.current_balance) AS "balance"
FROM
account_balance ab
WHERE
ab.account_id = a.id AND ab.date <= v_end
ORDER BY
ab.date DESC
LIMIT 1
) b1 ON TRUE
WHERE
a.id = p_account_id;
END;
$$ LANGUAGE plpgsql STABLE;
ALTER TABLE "investment_transaction"
DROP COLUMN "category_old";

View file

@ -217,24 +217,22 @@ enum InvestmentTransactionCategory {
} }
model InvestmentTransaction { model InvestmentTransaction {
id Int @id @default(autoincrement()) id Int @id @default(autoincrement())
createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz(6) createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz(6)
updatedAt DateTime @default(now()) @updatedAt @map("updated_at") @db.Timestamptz(6) updatedAt DateTime @default(now()) @updatedAt @map("updated_at") @db.Timestamptz(6)
account Account @relation(fields: [accountId], references: [id], onDelete: Cascade) account Account @relation(fields: [accountId], references: [id], onDelete: Cascade)
accountId Int @map("account_id") accountId Int @map("account_id")
security Security? @relation(fields: [securityId], references: [id], onDelete: Cascade) security Security? @relation(fields: [securityId], references: [id], onDelete: Cascade)
securityId Int? @map("security_id") securityId Int? @map("security_id")
date DateTime @db.Date date DateTime @db.Date
name String name String
amount Decimal @db.Decimal(19, 4) amount Decimal @db.Decimal(19, 4)
fees Decimal? @db.Decimal(19, 4) fees Decimal? @db.Decimal(19, 4)
flow TransactionFlow @default(dbgenerated("\nCASE\n WHEN (amount < (0)::numeric) THEN 'INFLOW'::\"TransactionFlow\"\n ELSE 'OUTFLOW'::\"TransactionFlow\"\nEND")) flow TransactionFlow @default(dbgenerated("\nCASE\n WHEN (amount < (0)::numeric) THEN 'INFLOW'::\"TransactionFlow\"\n ELSE 'OUTFLOW'::\"TransactionFlow\"\nEND"))
quantity Decimal @db.Decimal(36, 18) quantity Decimal @db.Decimal(36, 18)
price Decimal @db.Decimal(23, 8) price Decimal @db.Decimal(23, 8)
currencyCode String @default("USD") @map("currency_code") currencyCode String @default("USD") @map("currency_code")
category InvestmentTransactionCategory @default(other)
// Derived from provider types
category InvestmentTransactionCategory @default(dbgenerated("\nCASE\n WHEN (plaid_type = 'buy'::text) THEN 'buy'::\"InvestmentTransactionCategory\"\n WHEN (plaid_type = 'sell'::text) THEN 'sell'::\"InvestmentTransactionCategory\"\n WHEN (plaid_subtype = ANY (ARRAY['dividend'::text, 'qualified dividend'::text, 'non-qualified dividend'::text])) THEN 'dividend'::\"InvestmentTransactionCategory\"\n WHEN (plaid_subtype = ANY (ARRAY['non-resident tax'::text, 'tax'::text, 'tax withheld'::text])) THEN 'tax'::\"InvestmentTransactionCategory\"\n WHEN ((plaid_type = 'fee'::text) OR (plaid_subtype = ANY (ARRAY['account fee'::text, 'legal fee'::text, 'management fee'::text, 'margin expense'::text, 'transfer fee'::text, 'trust fee'::text]))) THEN 'fee'::\"InvestmentTransactionCategory\"\n WHEN (plaid_type = 'cash'::text) THEN 'transfer'::\"InvestmentTransactionCategory\"\n WHEN (plaid_type = 'cancel'::text) THEN 'cancel'::\"InvestmentTransactionCategory\"\n ELSE 'other'::\"InvestmentTransactionCategory\"\nEND"))
// plaid data // plaid data
plaidInvestmentTransactionId String? @unique @map("plaid_investment_transaction_id") plaidInvestmentTransactionId String? @unique @map("plaid_investment_transaction_id")