Compare commits
33
Commits
8478fa0b38
..
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
df3b84079c | ||
|
|
d11e4fd0c4 | ||
|
|
a2ba83e6f8 | ||
|
|
0d564ff299 | ||
|
|
bc4aa8242c | ||
|
|
9d23828345 | ||
|
|
e60ab43abd | ||
|
|
6d2150d1b2 | ||
|
|
7f78d2274c | ||
|
|
e739d5ea5b | ||
|
|
09209be500 | ||
|
|
c86a79b46e | ||
|
|
b65fad16e4 | ||
|
|
a60d0effff | ||
|
|
c230aef034 | ||
|
|
ca0df97e73 | ||
|
|
5c3f734bbc | ||
|
|
e3ed72c9b0 | ||
|
|
3c3b6dc4e6 | ||
|
|
c3796720b7 | ||
|
|
35fd2b90d2 | ||
|
|
4b2f0f3bf7 | ||
|
|
29fe8bee39 | ||
|
|
49a1700706 | ||
|
|
04b58d3075 | ||
|
|
076a573711 | ||
|
|
af0588c5ef | ||
|
|
21a63a10a4 | ||
|
|
08766dc0e0 | ||
|
|
bf47520d31 | ||
|
|
fc1cea32b5 | ||
|
|
92462bd316 | ||
|
|
f0b0cb1567 |
@@ -28,6 +28,7 @@
|
|||||||
**/values.dev.yaml
|
**/values.dev.yaml
|
||||||
/bin
|
/bin
|
||||||
/target
|
/target
|
||||||
|
!/target/**/monzo-ingestion
|
||||||
/.idea
|
/.idea
|
||||||
LICENSE
|
LICENSE
|
||||||
README.md
|
README.md
|
||||||
|
|||||||
@@ -1 +1 @@
|
|||||||
DATABASE_URL=postgres://postgres@localhost/logos_neu
|
DATABASE_URL=postgres://postgres@localhost/monzo_development
|
||||||
|
|||||||
+79
-23
@@ -1,46 +1,102 @@
|
|||||||
name: Build and Publish Docker Container
|
name: Build and Publish
|
||||||
|
|
||||||
on:
|
on:
|
||||||
push:
|
push:
|
||||||
branches:
|
branches: [ main ]
|
||||||
- main
|
|
||||||
|
env:
|
||||||
|
CARGO_TERM_COLOR: always
|
||||||
|
RUST_BINARY_NAME: monzo-ingestion
|
||||||
|
|
||||||
jobs:
|
jobs:
|
||||||
build:
|
build:
|
||||||
|
name: Build and Test
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
container:
|
container: catthehacker/ubuntu:act-latest
|
||||||
image: catthehacker/ubuntu:act-latest
|
|
||||||
|
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout code
|
- uses: actions/checkout@v3
|
||||||
uses: actions/checkout@v2
|
|
||||||
|
|
||||||
- name: Set up Docker Buildx
|
- name: Install Rust
|
||||||
uses: docker/setup-buildx-action@v1
|
uses: actions-rs/toolchain@v1
|
||||||
|
with:
|
||||||
|
toolchain: stable
|
||||||
|
profile: minimal
|
||||||
|
override: true
|
||||||
|
components: rustfmt, clippy
|
||||||
|
|
||||||
- name: Cache Docker layers
|
- name: Add ARM64 target
|
||||||
uses: https://github.com/actions/cache@v3
|
run: rustup target add aarch64-unknown-linux-musl
|
||||||
|
|
||||||
|
- name: Install ARM64 toolchain
|
||||||
|
run: |
|
||||||
|
apt-get update
|
||||||
|
apt-get install -y gcc-aarch64-linux-gnu build-essential musl-tools
|
||||||
|
|
||||||
|
- name: Cache dependencies
|
||||||
|
uses: actions/cache@v3
|
||||||
with:
|
with:
|
||||||
path: |
|
path: |
|
||||||
/tmp/.buildx-cache.bookworm
|
~/.cargo
|
||||||
/tmp/.buildx-cache.latest
|
target/
|
||||||
key: ${{ runner.os }}-buildx-${{ gitea.sha }}
|
key: "${{ runner.os }}-cargo-${{ hashFiles('**/Cargo.lock') }}"
|
||||||
restore-keys: |
|
restore-keys: |
|
||||||
${{ runner.os }}-buildx-
|
${{ runner.os }}-cargo-
|
||||||
|
|
||||||
- name: Login to Docker
|
- name: Build (x86_64)
|
||||||
uses: docker/login-action@v1
|
uses: actions-rs/cargo@v1
|
||||||
|
with:
|
||||||
|
command: build
|
||||||
|
args: --release --all-features
|
||||||
|
|
||||||
|
- name: Build (ARM64)
|
||||||
|
uses: actions-rs/cargo@v1
|
||||||
|
with:
|
||||||
|
command: build
|
||||||
|
args: --release --all-features --target aarch64-unknown-linux-musl
|
||||||
|
env:
|
||||||
|
CARGO_TARGET_AARCH64_UNKNOWN_LINUX_MUSL_LINKER: aarch64-linux-gnu-gcc
|
||||||
|
|
||||||
|
- name: Set up QEMU
|
||||||
|
uses: docker/setup-qemu-action@v2
|
||||||
|
|
||||||
|
- name: Set up Docker Buildx
|
||||||
|
uses: docker/setup-buildx-action@v2
|
||||||
|
|
||||||
|
- name: Login to DockerHub
|
||||||
|
uses: docker/login-action@v2
|
||||||
with:
|
with:
|
||||||
registry: git.joshuacoles.me
|
registry: git.joshuacoles.me
|
||||||
username: ${{ secrets.DOCKER_USERNAME }}
|
username: ${{ secrets.DOCKER_USERNAME }}
|
||||||
password: ${{ secrets.DOCKER_PASSWORD }}
|
password: ${{ secrets.DOCKER_PASSWORD }}
|
||||||
|
|
||||||
- name: Build and Push Docker image
|
- name: Upload artifacts
|
||||||
uses: docker/build-push-action@master
|
uses: actions/upload-artifact@v3
|
||||||
|
with:
|
||||||
|
name: binaries
|
||||||
|
path: |
|
||||||
|
target/release/${{ env.RUST_BINARY_NAME }}
|
||||||
|
target/aarch64-unknown-linux-musl/release/${{ env.RUST_BINARY_NAME }}
|
||||||
|
|
||||||
|
# # The target directory is kept in the .dockerignore file, so allow these to be copied in we need to move them to a
|
||||||
|
# # new directory.
|
||||||
|
# - run: mv target docker-binaries
|
||||||
|
|
||||||
|
- name: Build and push multi-arch Docker image
|
||||||
|
uses: docker/build-push-action@v4
|
||||||
with:
|
with:
|
||||||
context: .
|
context: .
|
||||||
|
file: ./Dockerfile
|
||||||
|
platforms: linux/amd64,linux/arm64
|
||||||
push: true
|
push: true
|
||||||
file: ./Dockerfile.cache
|
tags: git.joshuacoles.me/${{ github.repository }}:${{ github.sha }},git.joshuacoles.me/${{ github.repository }}:latest
|
||||||
tags: git.joshuacoles.me/personal/monzo-ingestion:latest
|
build-args: |
|
||||||
cache-from: type=registry,ref=user/app:latest
|
BINARY_NAME=${{ env.RUST_BINARY_NAME }}
|
||||||
cache-to: type=inline
|
|
||||||
|
- uses: robiningelbrecht/ntfy-action@v1.0.0
|
||||||
|
name: Notify via ntfy.sh
|
||||||
|
if: always()
|
||||||
|
with:
|
||||||
|
url: ${{ secrets.NTFY_URL }}
|
||||||
|
topic: ${{ secrets.NTFY_TOPIC }}
|
||||||
|
job_status: ${{ job.status }}
|
||||||
|
|||||||
Generated
+1635
-1097
File diff suppressed because it is too large
Load Diff
+17
-16
@@ -7,30 +7,31 @@ edition = "2021"
|
|||||||
entity = { path = "entity" }
|
entity = { path = "entity" }
|
||||||
migration = { path = "migration" }
|
migration = { path = "migration" }
|
||||||
|
|
||||||
axum = { version = "0.7.5", features = ["multipart"] }
|
axum = { version = "0.8.8", features = ["multipart"] }
|
||||||
tokio = { version = "1.37.0", features = ["full"] }
|
tokio = { version = "1.48.0", features = ["full"] }
|
||||||
sea-orm = { version = "1.0.0-rc.4", features = [
|
sea-orm = { version = "1.1.0", features = [
|
||||||
"sqlx-postgres",
|
"sqlx-postgres",
|
||||||
"runtime-tokio-rustls",
|
"runtime-tokio-rustls",
|
||||||
"macros"
|
"macros"
|
||||||
] }
|
] }
|
||||||
|
|
||||||
serde = { version = "1.0.203", features = ["derive"] }
|
serde = { version = "1.0", features = ["derive"] }
|
||||||
serde_json = "1.0.117"
|
serde_json = "1.0"
|
||||||
tracing-subscriber = "0.3.18"
|
tracing-subscriber = "0.3.18"
|
||||||
tracing = "0.1.40"
|
tracing = "0.1.40"
|
||||||
anyhow = { version = "1.0.86", features = ["backtrace"] }
|
anyhow = { version = "1.0", features = ["backtrace"] }
|
||||||
thiserror = "1.0.61"
|
thiserror = "2.0"
|
||||||
http = "1.1.0"
|
http = "1.1"
|
||||||
chrono = { version = "0.4.38", features = ["serde"] }
|
chrono = { version = "0.4", features = ["serde"] }
|
||||||
num-traits = "0.2.19"
|
num-traits = "0.2"
|
||||||
csv = "1.3.0"
|
csv = "1.3.0"
|
||||||
clap = "4.5.4"
|
clap = "4.5"
|
||||||
testcontainers = "0.17.0"
|
testcontainers = "0.26"
|
||||||
testcontainers-modules = { version = "0.5.0", features = ["postgres"] }
|
testcontainers-modules = { version = "0.14", features = ["postgres"] }
|
||||||
sqlx = { version = "0.7.4", features = ["postgres"] }
|
sqlx = { version = "0.8", features = ["postgres"] }
|
||||||
tower-http = { version = "0.5.2", features = ["trace"] }
|
tower-http = { version = "0.6", features = ["trace"] }
|
||||||
bytes = "1.6.0"
|
bytes = "1.7"
|
||||||
|
once_cell = "1.19"
|
||||||
|
|
||||||
[workspace]
|
[workspace]
|
||||||
members = [".", "migration", "entity"]
|
members = [".", "migration", "entity"]
|
||||||
|
|||||||
+12
-76
@@ -1,78 +1,14 @@
|
|||||||
# syntax=docker/dockerfile:1
|
FROM busybox AS platform_determiner
|
||||||
|
ARG TARGETPLATFORM
|
||||||
|
ARG BINARY_NAME
|
||||||
|
|
||||||
# Comments are provided throughout this file to help you get started.
|
COPY /target /target
|
||||||
# If you need more help, visit the Dockerfile reference guide at
|
RUN case "$TARGETPLATFORM" in \
|
||||||
# https://docs.docker.com/engine/reference/builder/
|
"linux/amd64") BINARY_PATH="/target/release/${BINARY_NAME}" ;; \
|
||||||
|
"linux/arm64") BINARY_PATH="/target/aarch64-unknown-linux-musl/release/${BINARY_NAME}" ;; \
|
||||||
|
*) exit 1 ;; \
|
||||||
|
esac && mv "$BINARY_PATH" "/usr/bin/monzo-ingestion" && chmod +x "/usr/bin/monzo-ingestion"
|
||||||
|
|
||||||
################################################################################
|
FROM --platform=$TARGETPLATFORM debian:bookworm-slim
|
||||||
# Create a stage for building the application.
|
COPY --from=platform_determiner /usr/bin/monzo-ingestion /usr/local/bin/monzo-ingestion
|
||||||
|
ENTRYPOINT ["/usr/local/bin/monzo-ingestion"]
|
||||||
ARG RUST_VERSION=1.76.0
|
|
||||||
ARG APP_NAME=monzo-ingestion
|
|
||||||
FROM rust:${RUST_VERSION}-slim-bullseye AS build
|
|
||||||
ARG APP_NAME
|
|
||||||
WORKDIR /app
|
|
||||||
|
|
||||||
# Build the application.
|
|
||||||
# Leverage a cache mount to /usr/local/cargo/registry/
|
|
||||||
# for downloaded dependencies and a cache mount to /app/target/ for
|
|
||||||
# compiled dependencies which will speed up subsequent builds.
|
|
||||||
# Leverage a bind mount to the src directory to avoid having to copy the
|
|
||||||
# source code into the container. Once built, copy the executable to an
|
|
||||||
# output directory before the cache mounted /app/target is unmounted.
|
|
||||||
RUN --mount=type=bind,source=src,target=src \
|
|
||||||
--mount=type=bind,source=entity,target=entity \
|
|
||||||
--mount=type=bind,source=migration,target=migration \
|
|
||||||
--mount=type=bind,source=Cargo.toml,target=Cargo.toml \
|
|
||||||
--mount=type=bind,source=Cargo.lock,target=Cargo.lock \
|
|
||||||
--mount=type=cache,target=/app/target/ \
|
|
||||||
--mount=type=cache,target=/usr/local/cargo/registry/ \
|
|
||||||
<<EOF
|
|
||||||
set -e
|
|
||||||
cargo build --locked --release
|
|
||||||
cp ./target/release/$APP_NAME /bin/server
|
|
||||||
EOF
|
|
||||||
|
|
||||||
################################################################################
|
|
||||||
# Create a new stage for running the application that contains the minimal
|
|
||||||
# runtime dependencies for the application. This often uses a different base
|
|
||||||
# image from the build stage where the necessary files are copied from the build
|
|
||||||
# stage.
|
|
||||||
#
|
|
||||||
# The example below uses the debian bullseye image as the foundation for running the app.
|
|
||||||
# By specifying the "bullseye-slim" tag, it will also use whatever happens to be the
|
|
||||||
# most recent version of that tag when you build your Dockerfile. If
|
|
||||||
# reproducability is important, consider using a digest
|
|
||||||
# (e.g., debian@sha256:ac707220fbd7b67fc19b112cee8170b41a9e97f703f588b2cdbbcdcecdd8af57).
|
|
||||||
FROM debian:bullseye-slim AS final
|
|
||||||
|
|
||||||
RUN set -ex; \
|
|
||||||
apt-get update && \
|
|
||||||
apt-get -y install --no-install-recommends \
|
|
||||||
ca-certificates curl && \
|
|
||||||
rm -rf /var/lib/apt/lists/*
|
|
||||||
|
|
||||||
# Create a non-privileged user that the app will run under.
|
|
||||||
# See https://docs.docker.com/develop/develop-images/dockerfile_best-practices/#user
|
|
||||||
ARG UID=10001
|
|
||||||
RUN adduser \
|
|
||||||
--disabled-password \
|
|
||||||
--gecos "" \
|
|
||||||
--home "/nonexistent" \
|
|
||||||
--shell "/sbin/nologin" \
|
|
||||||
--no-create-home \
|
|
||||||
--uid "${UID}" \
|
|
||||||
appuser
|
|
||||||
USER appuser
|
|
||||||
|
|
||||||
# Copy the executable from the "build" stage.
|
|
||||||
COPY --from=build /bin/server /bin/
|
|
||||||
|
|
||||||
# Expose the port that the application listens on.
|
|
||||||
EXPOSE 3000
|
|
||||||
|
|
||||||
HEALTHCHECK --interval=5s --timeout=3s --retries=3 \
|
|
||||||
CMD curl -f http://localhost:3000/health || exit 1
|
|
||||||
|
|
||||||
# What the container should run when it is started.
|
|
||||||
CMD ["/bin/server", "web", "--addr", "0.0.0.0:3000"]
|
|
||||||
|
|||||||
@@ -1,58 +0,0 @@
|
|||||||
# Stage 1: Build
|
|
||||||
ARG RUST_VERSION=1.76.0
|
|
||||||
FROM lukemathwalker/cargo-chef:latest-rust-${RUST_VERSION} as chef
|
|
||||||
WORKDIR /build/
|
|
||||||
# hadolint ignore=DL3008
|
|
||||||
RUN apt-get update && \
|
|
||||||
apt-get install -y --no-install-recommends \
|
|
||||||
lld \
|
|
||||||
clang \
|
|
||||||
libclang-dev \
|
|
||||||
&& apt-get clean \
|
|
||||||
&& rm -rf /var/lib/apt/lists/*
|
|
||||||
|
|
||||||
FROM chef as planner
|
|
||||||
COPY . .
|
|
||||||
RUN cargo chef prepare --recipe-path recipe.json
|
|
||||||
|
|
||||||
FROM chef as builder
|
|
||||||
COPY --from=planner /build/recipe.json recipe.json
|
|
||||||
# Build dependencies - this is the caching Docker layer!
|
|
||||||
RUN cargo chef cook --release -p monzo-ingestion --recipe-path recipe.json
|
|
||||||
# Build application
|
|
||||||
COPY . .
|
|
||||||
RUN cargo build --release -p monzo-ingestion
|
|
||||||
|
|
||||||
# Stage 2: Run
|
|
||||||
FROM debian:bookworm-slim AS runtime
|
|
||||||
|
|
||||||
RUN set -ex; \
|
|
||||||
apt-get update && \
|
|
||||||
apt-get -y install --no-install-recommends \
|
|
||||||
ca-certificates curl && \
|
|
||||||
rm -rf /var/lib/apt/lists/*
|
|
||||||
|
|
||||||
# Create a non-privileged user that the app will run under.
|
|
||||||
# See https://docs.docker.com/develop/develop-images/dockerfile_best-practices/#user
|
|
||||||
ARG UID=10001
|
|
||||||
RUN adduser \
|
|
||||||
--disabled-password \
|
|
||||||
--gecos "" \
|
|
||||||
--home "/nonexistent" \
|
|
||||||
--shell "/sbin/nologin" \
|
|
||||||
--no-create-home \
|
|
||||||
--uid "${UID}" \
|
|
||||||
appuser
|
|
||||||
USER appuser
|
|
||||||
|
|
||||||
# Copy the executable from the "build" stage.
|
|
||||||
COPY --from=builder /build/target/release/monzo-ingestion /bin/server
|
|
||||||
|
|
||||||
# Expose the port that the application listens on.
|
|
||||||
EXPOSE 3000
|
|
||||||
|
|
||||||
HEALTHCHECK --interval=5s --timeout=3s --retries=3 \
|
|
||||||
CMD curl -f http://localhost:3000/health || exit 1
|
|
||||||
|
|
||||||
# What the container should run when it is started.
|
|
||||||
CMD ["/bin/server", "serve", "--addr", "0.0.0.0:3000"]
|
|
||||||
+1
-1
@@ -9,5 +9,5 @@ name = "entity"
|
|||||||
path = "src/lib.rs"
|
path = "src/lib.rs"
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
sea-orm = { version = "1.0.0-rc.4" }
|
sea-orm = { version = "1.1.0" }
|
||||||
serde = { version = "1.0.203", features = ["derive"] }
|
serde = { version = "1.0.203", features = ["derive"] }
|
||||||
|
|||||||
@@ -0,0 +1,26 @@
|
|||||||
|
//! `SeaORM` Entity, @generated by sea-orm-codegen 1.1.0
|
||||||
|
|
||||||
|
use sea_orm::entity::prelude::*;
|
||||||
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, PartialEq, DeriveEntityModel, Eq, Serialize, Deserialize)]
|
||||||
|
#[sea_orm(table_name = "account")]
|
||||||
|
pub struct Model {
|
||||||
|
#[sea_orm(primary_key)]
|
||||||
|
pub id: i32,
|
||||||
|
pub name: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
||||||
|
pub enum Relation {
|
||||||
|
#[sea_orm(has_many = "super::transaction::Entity")]
|
||||||
|
Transaction,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Related<super::transaction::Entity> for Entity {
|
||||||
|
fn to() -> RelationDef {
|
||||||
|
Relation::Transaction.def()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ActiveModelBehavior for ActiveModel {}
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
//! `SeaORM` Entity. Generated by sea-orm-codegen 0.12.2
|
//! `SeaORM` Entity, @generated by sea-orm-codegen 1.1.0
|
||||||
|
|
||||||
use sea_orm::entity::prelude::*;
|
use sea_orm::entity::prelude::*;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
@@ -14,6 +14,21 @@ pub struct Model {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
||||||
pub enum Relation {}
|
pub enum Relation {
|
||||||
|
#[sea_orm(
|
||||||
|
belongs_to = "super::transaction::Entity",
|
||||||
|
from = "Column::TransactionId",
|
||||||
|
to = "super::transaction::Column::Id",
|
||||||
|
on_update = "NoAction",
|
||||||
|
on_delete = "Cascade"
|
||||||
|
)]
|
||||||
|
Transaction,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Related<super::transaction::Entity> for Entity {
|
||||||
|
fn to() -> RelationDef {
|
||||||
|
Relation::Transaction.def()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl ActiveModelBehavior for ActiveModel {}
|
impl ActiveModelBehavior for ActiveModel {}
|
||||||
|
|||||||
+2
-1
@@ -1,6 +1,7 @@
|
|||||||
//! `SeaORM` Entity. Generated by sea-orm-codegen 0.12.2
|
//! `SeaORM` Entity, @generated by sea-orm-codegen 1.1.0
|
||||||
|
|
||||||
pub mod prelude;
|
pub mod prelude;
|
||||||
|
|
||||||
|
pub mod account;
|
||||||
pub mod expenditure;
|
pub mod expenditure;
|
||||||
pub mod transaction;
|
pub mod transaction;
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
//! `SeaORM` Entity. Generated by sea-orm-codegen 0.12.2
|
//! `SeaORM` Entity, @generated by sea-orm-codegen 1.1.0
|
||||||
|
|
||||||
|
pub use super::account::Entity as Account;
|
||||||
pub use super::expenditure::Entity as Expenditure;
|
pub use super::expenditure::Entity as Expenditure;
|
||||||
pub use super::transaction::Entity as Transaction;
|
pub use super::transaction::Entity as Transaction;
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
//! `SeaORM` Entity. Generated by sea-orm-codegen 0.12.2
|
//! `SeaORM` Entity, @generated by sea-orm-codegen 1.1.0
|
||||||
|
|
||||||
use sea_orm::entity::prelude::*;
|
use sea_orm::entity::prelude::*;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
@@ -18,9 +18,33 @@ pub struct Model {
|
|||||||
pub description: Option<String>,
|
pub description: Option<String>,
|
||||||
#[sea_orm(unique)]
|
#[sea_orm(unique)]
|
||||||
pub identity_hash: Option<i64>,
|
pub identity_hash: Option<i64>,
|
||||||
|
pub account_id: Option<i32>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
||||||
pub enum Relation {}
|
pub enum Relation {
|
||||||
|
#[sea_orm(
|
||||||
|
belongs_to = "super::account::Entity",
|
||||||
|
from = "Column::AccountId",
|
||||||
|
to = "super::account::Column::Id",
|
||||||
|
on_update = "NoAction",
|
||||||
|
on_delete = "Cascade"
|
||||||
|
)]
|
||||||
|
Account,
|
||||||
|
#[sea_orm(has_many = "super::expenditure::Entity")]
|
||||||
|
Expenditure,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Related<super::account::Entity> for Entity {
|
||||||
|
fn to() -> RelationDef {
|
||||||
|
Relation::Account.def()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Related<super::expenditure::Entity> for Entity {
|
||||||
|
fn to() -> RelationDef {
|
||||||
|
Relation::Expenditure.def()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl ActiveModelBehavior for ActiveModel {}
|
impl ActiveModelBehavior for ActiveModel {}
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ path = "src/lib.rs"
|
|||||||
async-std = { version = "1", features = ["attributes", "tokio1"] }
|
async-std = { version = "1", features = ["attributes", "tokio1"] }
|
||||||
|
|
||||||
[dependencies.sea-orm-migration]
|
[dependencies.sea-orm-migration]
|
||||||
version = "1.0.0-rc.4"
|
version = "1.1.0"
|
||||||
features = [
|
features = [
|
||||||
"runtime-tokio-rustls", # `ASYNC_RUNTIME` feature
|
"runtime-tokio-rustls", # `ASYNC_RUNTIME` feature
|
||||||
"sqlx-postgres", # `DATABASE_DRIVER` feature
|
"sqlx-postgres", # `DATABASE_DRIVER` feature
|
||||||
|
|||||||
@@ -3,6 +3,9 @@ pub use sea_orm_migration::prelude::*;
|
|||||||
pub mod m20230904_141851_create_monzo_tables;
|
pub mod m20230904_141851_create_monzo_tables;
|
||||||
mod m20240529_195030_add_transaction_identity_hash;
|
mod m20240529_195030_add_transaction_identity_hash;
|
||||||
mod m20240603_162500_make_title_optional;
|
mod m20240603_162500_make_title_optional;
|
||||||
|
mod m20241015_195220_add_account_to_transactions;
|
||||||
|
mod m20241015_200222_add_expenditure_transaction_fk;
|
||||||
|
mod m20241015_200652_add_accounts;
|
||||||
|
|
||||||
pub struct Migrator;
|
pub struct Migrator;
|
||||||
|
|
||||||
@@ -17,6 +20,9 @@ impl MigratorTrait for Migrator {
|
|||||||
Box::new(m20230904_141851_create_monzo_tables::Migration),
|
Box::new(m20230904_141851_create_monzo_tables::Migration),
|
||||||
Box::new(m20240529_195030_add_transaction_identity_hash::Migration),
|
Box::new(m20240529_195030_add_transaction_identity_hash::Migration),
|
||||||
Box::new(m20240603_162500_make_title_optional::Migration),
|
Box::new(m20240603_162500_make_title_optional::Migration),
|
||||||
|
Box::new(m20241015_195220_add_account_to_transactions::Migration),
|
||||||
|
Box::new(m20241015_200222_add_expenditure_transaction_fk::Migration),
|
||||||
|
Box::new(m20241015_200652_add_accounts::Migration),
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,76 @@
|
|||||||
|
use sea_orm_migration::{prelude::*, schema::*};
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.create_table(
|
||||||
|
Table::create()
|
||||||
|
.table(Account::Table)
|
||||||
|
.if_not_exists()
|
||||||
|
.col(
|
||||||
|
ColumnDef::new(Account::Id)
|
||||||
|
.integer()
|
||||||
|
.auto_increment()
|
||||||
|
.not_null()
|
||||||
|
.primary_key(),
|
||||||
|
)
|
||||||
|
.col(
|
||||||
|
ColumnDef::new(Account::Name)
|
||||||
|
.string()
|
||||||
|
.not_null(),
|
||||||
|
)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
manager.alter_table(
|
||||||
|
TableAlterStatement::new()
|
||||||
|
.table(Transaction::Table)
|
||||||
|
.add_column(ColumnDef::new(Transaction::AccountId).integer())
|
||||||
|
.to_owned(),
|
||||||
|
).await?;
|
||||||
|
|
||||||
|
manager
|
||||||
|
.create_foreign_key(
|
||||||
|
ForeignKey::create()
|
||||||
|
.name("fk_transaction_account_id")
|
||||||
|
.from(Transaction::Table, Transaction::AccountId)
|
||||||
|
.to(Account::Table, Account::Id)
|
||||||
|
.on_delete(ForeignKeyAction::Cascade)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager.alter_table(
|
||||||
|
TableAlterStatement::new()
|
||||||
|
.table(Transaction::Table)
|
||||||
|
.drop_column(Transaction::AccountId)
|
||||||
|
.to_owned(),
|
||||||
|
).await?;
|
||||||
|
|
||||||
|
manager.drop_table(Table::drop().table(Account::Table).to_owned()).await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Account {
|
||||||
|
Table,
|
||||||
|
Id,
|
||||||
|
Name,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Transaction {
|
||||||
|
Table,
|
||||||
|
AccountId,
|
||||||
|
}
|
||||||
@@ -0,0 +1,40 @@
|
|||||||
|
use sea_orm_migration::{prelude::*, schema::*};
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.create_foreign_key(
|
||||||
|
ForeignKey::create()
|
||||||
|
.name("fk_expenditure_transaction_id")
|
||||||
|
.from(Expenditure::Table, Expenditure::TransactionId)
|
||||||
|
.to(Transaction::Table, Transaction::Id)
|
||||||
|
.on_delete(ForeignKeyAction::Cascade)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager.drop_foreign_key(ForeignKey::drop().name("fk_expenditure_transaction_id").to_owned()).await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Transaction {
|
||||||
|
Table,
|
||||||
|
Id,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Expenditure {
|
||||||
|
Table,
|
||||||
|
TransactionId,
|
||||||
|
}
|
||||||
@@ -0,0 +1,53 @@
|
|||||||
|
use sea_orm_migration::{prelude::*, schema::*};
|
||||||
|
use crate::sea_orm::sqlx;
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager.exec_stmt(
|
||||||
|
InsertStatement::new()
|
||||||
|
.into_table(Account::Table)
|
||||||
|
.columns(vec![Account::Name])
|
||||||
|
.values_panic(["Monzo".into()])
|
||||||
|
.values_panic(["Flex".into()])
|
||||||
|
.to_owned()
|
||||||
|
).await?;
|
||||||
|
|
||||||
|
let id: i32 = manager.get_connection().query_one(
|
||||||
|
sea_orm::Statement::from_string(
|
||||||
|
manager.get_database_backend(),
|
||||||
|
"SELECT id FROM account WHERE name = 'Monzo'",
|
||||||
|
)
|
||||||
|
).await?.expect("Monzo account not found").try_get_by_index(0)?;
|
||||||
|
|
||||||
|
manager.exec_stmt(
|
||||||
|
UpdateStatement::new()
|
||||||
|
.table(Transaction::Table)
|
||||||
|
.values([(Transaction::AccountId, id.into())])
|
||||||
|
.to_owned()
|
||||||
|
).await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
// This is a data creation migration, so we do not roll back the migration
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Account {
|
||||||
|
Table,
|
||||||
|
Id,
|
||||||
|
Name,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Transaction {
|
||||||
|
Table,
|
||||||
|
AccountId,
|
||||||
|
}
|
||||||
@@ -29,6 +29,7 @@ impl AppError {
|
|||||||
impl IntoResponse for AppError {
|
impl IntoResponse for AppError {
|
||||||
fn into_response(self) -> Response {
|
fn into_response(self) -> Response {
|
||||||
let status_code = match self {
|
let status_code = match self {
|
||||||
|
AppError::BadRequest(_) => StatusCode::BAD_REQUEST,
|
||||||
_ => StatusCode::INTERNAL_SERVER_ERROR,
|
_ => StatusCode::INTERNAL_SERVER_ERROR,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
+129
-23
@@ -20,16 +20,17 @@ pub struct Insertion {
|
|||||||
pub async fn insert(
|
pub async fn insert(
|
||||||
db: &DatabaseConnection,
|
db: &DatabaseConnection,
|
||||||
monzo_rows: Vec<MonzoRow>,
|
monzo_rows: Vec<MonzoRow>,
|
||||||
|
account_id: i32,
|
||||||
) -> Result<Vec<String>, AppError> {
|
) -> Result<Vec<String>, AppError> {
|
||||||
let mut new_transaction_ids = Vec::new();
|
let mut new_transaction_ids = Vec::new();
|
||||||
let insertions = monzo_rows
|
let insertions = monzo_rows
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.map(MonzoRow::into_insertion)
|
.map(|row| MonzoRow::into_insertion(row, account_id))
|
||||||
.collect::<Result<Vec<_>, _>>()?;
|
.collect::<Result<Vec<_>, _>>()?;
|
||||||
|
|
||||||
for insertions in insertions.chunks(400) {
|
for insertions in insertions.chunks(200) {
|
||||||
let (new_or_updated_insertions, inserted_transaction_ids) =
|
let (new_or_updated_insertions, inserted_transaction_ids) =
|
||||||
whittle_insertions(insertions, &db).await?;
|
whittle_insertions(insertions, db).await?;
|
||||||
|
|
||||||
if new_or_updated_insertions.is_empty() {
|
if new_or_updated_insertions.is_empty() {
|
||||||
continue;
|
continue;
|
||||||
@@ -69,7 +70,7 @@ async fn update_expenditures(
|
|||||||
|
|
||||||
expenditure::Entity::insert_many(
|
expenditure::Entity::insert_many(
|
||||||
new_or_updated_insertions
|
new_or_updated_insertions
|
||||||
.into_iter()
|
.iter()
|
||||||
.flat_map(|i| &i.contained_expenditures)
|
.flat_map(|i| &i.contained_expenditures)
|
||||||
.cloned(),
|
.cloned(),
|
||||||
)
|
)
|
||||||
@@ -122,10 +123,12 @@ async fn whittle_insertions<'a>(
|
|||||||
.all(tx)
|
.all(tx)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
|
tracing::debug!("Found existing entries: {existing_hashes:?}");
|
||||||
|
|
||||||
// We will only update those where the hash is different to avoid unnecessary updates and
|
// We will only update those where the hash is different to avoid unnecessary updates and
|
||||||
// notifications.
|
// notifications.
|
||||||
let new_or_updated_insertions = insertions
|
let new_or_updated_insertions = insertions
|
||||||
.into_iter()
|
.iter()
|
||||||
.filter(|i| {
|
.filter(|i| {
|
||||||
let hash = i.identity_hash;
|
let hash = i.identity_hash;
|
||||||
!existing_hashes
|
!existing_hashes
|
||||||
@@ -155,23 +158,29 @@ async fn notify_new_transactions(
|
|||||||
}
|
}
|
||||||
|
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{notify_new_transactions, update_expenditures, update_transactions};
|
use super::{insert, notify_new_transactions, update_expenditures, update_transactions};
|
||||||
|
use crate::ingestion::ingestion_logic::from_json_row;
|
||||||
use anyhow::Error;
|
use anyhow::Error;
|
||||||
|
use entity::account;
|
||||||
use migration::MigratorTrait;
|
use migration::MigratorTrait;
|
||||||
use sea_orm::{DatabaseConnection, TransactionTrait};
|
use sea_orm::{ActiveModelTrait, DatabaseConnection, TransactionTrait};
|
||||||
|
use serde_json::Value;
|
||||||
use sqlx::postgres::PgListener;
|
use sqlx::postgres::PgListener;
|
||||||
use sqlx::PgPool;
|
use sqlx::{Executor, PgPool};
|
||||||
use testcontainers::runners::AsyncRunner;
|
use testcontainers::runners::AsyncRunner;
|
||||||
use testcontainers::ContainerAsync;
|
use testcontainers::ContainerAsync;
|
||||||
|
use tokio::sync::OnceCell;
|
||||||
|
|
||||||
async fn initialise() -> Result<
|
#[derive(Debug)]
|
||||||
(
|
struct DatabaseInstance {
|
||||||
ContainerAsync<testcontainers_modules::postgres::Postgres>,
|
container: ContainerAsync<testcontainers_modules::postgres::Postgres>,
|
||||||
DatabaseConnection,
|
db: DatabaseConnection,
|
||||||
PgPool,
|
pool: PgPool,
|
||||||
),
|
}
|
||||||
Error,
|
|
||||||
> {
|
static INSTANCE: OnceCell<DatabaseInstance> = OnceCell::const_new();
|
||||||
|
|
||||||
|
async fn initialise_db() -> Result<DatabaseInstance, Error> {
|
||||||
let container = testcontainers_modules::postgres::Postgres::default()
|
let container = testcontainers_modules::postgres::Postgres::default()
|
||||||
.start()
|
.start()
|
||||||
.await?;
|
.await?;
|
||||||
@@ -186,15 +195,35 @@ mod tests {
|
|||||||
migration::Migrator::up(&db, None).await?;
|
migration::Migrator::up(&db, None).await?;
|
||||||
|
|
||||||
let pool = PgPool::connect(connection_string).await?;
|
let pool = PgPool::connect(connection_string).await?;
|
||||||
|
let instance = DatabaseInstance {
|
||||||
|
container,
|
||||||
|
db,
|
||||||
|
pool,
|
||||||
|
};
|
||||||
|
|
||||||
Ok((container, db, pool))
|
Ok(instance)
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn get_or_initialize_db_instance() -> Result<&'static DatabaseInstance, Error> {
|
||||||
|
Ok(INSTANCE
|
||||||
|
.get_or_init(|| async { initialise_db().await.unwrap() })
|
||||||
|
.await)
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn create_test_account(db: &DatabaseConnection) -> Result<i32, Error> {
|
||||||
|
let new_account = account::ActiveModel {
|
||||||
|
id: sea_orm::ActiveValue::NotSet,
|
||||||
|
name: sea_orm::ActiveValue::Set("Test Account".to_string()),
|
||||||
|
};
|
||||||
|
let inserted = new_account.insert(db).await?;
|
||||||
|
Ok(inserted.id)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_no_new_insertions() -> Result<(), Error> {
|
async fn test_empty_insertion_list() -> Result<(), Error> {
|
||||||
let (_container, db, _pool) = initialise().await?;
|
let db = get_or_initialize_db_instance().await?;
|
||||||
let insertions = vec![];
|
let insertions = vec![];
|
||||||
let tx = db.begin().await?;
|
let tx = db.db.begin().await?;
|
||||||
update_transactions(&tx, &insertions).await?;
|
update_transactions(&tx, &insertions).await?;
|
||||||
update_expenditures(&tx, &insertions, &vec![]).await?;
|
update_expenditures(&tx, &insertions, &vec![]).await?;
|
||||||
tx.commit().await?;
|
tx.commit().await?;
|
||||||
@@ -204,8 +233,8 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_notify() -> Result<(), Error> {
|
async fn test_notify() -> Result<(), Error> {
|
||||||
let (_container, db, pool) = initialise().await?;
|
let dbi = get_or_initialize_db_instance().await?;
|
||||||
let mut listener = PgListener::connect_with(&pool).await?;
|
let mut listener = PgListener::connect_with(&dbi.pool).await?;
|
||||||
listener.listen("monzo_new_transactions").await?;
|
listener.listen("monzo_new_transactions").await?;
|
||||||
|
|
||||||
let ids = vec![
|
let ids = vec![
|
||||||
@@ -214,7 +243,7 @@ mod tests {
|
|||||||
"test3".to_string(),
|
"test3".to_string(),
|
||||||
];
|
];
|
||||||
|
|
||||||
notify_new_transactions(&db, &ids).await?;
|
notify_new_transactions(&dbi.db, &ids).await?;
|
||||||
|
|
||||||
let notification = listener.recv().await?;
|
let notification = listener.recv().await?;
|
||||||
let payload = notification.payload();
|
let payload = notification.payload();
|
||||||
@@ -227,4 +256,81 @@ mod tests {
|
|||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_notify_on_insert() -> Result<(), Error> {
|
||||||
|
let dbi = get_or_initialize_db_instance().await?;
|
||||||
|
let account_id = create_test_account(&dbi.db).await?;
|
||||||
|
let mut listener = PgListener::connect_with(&dbi.pool).await?;
|
||||||
|
listener.listen("monzo_new_transactions").await?;
|
||||||
|
|
||||||
|
let json = include_str!("../../fixtures/transactions.json");
|
||||||
|
let json: Vec<Vec<Value>> = serde_json::from_str(json).unwrap();
|
||||||
|
let data = json
|
||||||
|
.iter()
|
||||||
|
.map(|row| from_json_row(row))
|
||||||
|
.collect::<Result<Vec<_>, anyhow::Error>>()?;
|
||||||
|
|
||||||
|
insert(&dbi.db, data.clone(), account_id).await?;
|
||||||
|
let notification = listener.recv().await?;
|
||||||
|
let payload = notification.payload();
|
||||||
|
let mut payload = serde_json::from_str::<Vec<String>>(&payload)?;
|
||||||
|
payload.sort();
|
||||||
|
|
||||||
|
let mut ids = data
|
||||||
|
.iter()
|
||||||
|
.map(|row| row.transaction_id.clone())
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
ids.sort();
|
||||||
|
|
||||||
|
assert_eq!(payload, ids, "Inserted IDs do not match");
|
||||||
|
|
||||||
|
insert(&dbi.db, data.clone(), account_id).await?;
|
||||||
|
let notification = listener.recv().await?;
|
||||||
|
let payload = notification.payload();
|
||||||
|
let payload = serde_json::from_str::<Vec<String>>(&payload)?;
|
||||||
|
assert_eq!(
|
||||||
|
payload,
|
||||||
|
Vec::<String>::new(),
|
||||||
|
"Re-inserting identical rows triggered double notification"
|
||||||
|
);
|
||||||
|
|
||||||
|
let mut altered_data = data.clone();
|
||||||
|
altered_data[0].description = Some("New description".to_string());
|
||||||
|
|
||||||
|
assert_ne!(
|
||||||
|
altered_data[0].compute_hash(),
|
||||||
|
data[0].compute_hash(),
|
||||||
|
"Alterations have the same hash"
|
||||||
|
);
|
||||||
|
|
||||||
|
insert(&dbi.db, altered_data.clone(), account_id).await?;
|
||||||
|
let notification = listener.recv().await?;
|
||||||
|
let payload = notification.payload();
|
||||||
|
let payload = serde_json::from_str::<Vec<String>>(&payload)?;
|
||||||
|
assert_eq!(
|
||||||
|
payload,
|
||||||
|
vec![altered_data[0].transaction_id.clone()],
|
||||||
|
"Re-inserting altered row failed to re-trigger notification"
|
||||||
|
);
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn get_account_id(
|
||||||
|
p0: &DatabaseConnection,
|
||||||
|
p1: Option<String>,
|
||||||
|
) -> Result<i32, AppError> {
|
||||||
|
let p1 = p1.unwrap_or("Monzo".to_string());
|
||||||
|
|
||||||
|
entity::prelude::Account::find()
|
||||||
|
.filter(entity::account::Column::Name.eq(p1))
|
||||||
|
.select_only()
|
||||||
|
.column(entity::account::Column::Id)
|
||||||
|
.into_tuple::<i32>()
|
||||||
|
.one(p0)
|
||||||
|
.await?
|
||||||
|
.ok_or(AppError::BadRequest(anyhow!("Account not found")))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,9 @@
|
|||||||
|
#[allow(dead_code)]
|
||||||
|
pub mod headings {
|
||||||
|
#[allow(unused_imports)]
|
||||||
|
pub use super::super::ingestion_logic::headings::*;
|
||||||
|
|
||||||
|
// Additional Flex headings
|
||||||
|
pub const MONEY_OUT: usize = 16;
|
||||||
|
pub const MONEY_IN: usize = 17;
|
||||||
|
}
|
||||||
@@ -9,9 +9,10 @@ use sea_orm::prelude::Decimal;
|
|||||||
use sea_orm::IntoActiveModel;
|
use sea_orm::IntoActiveModel;
|
||||||
use serde_json::Value;
|
use serde_json::Value;
|
||||||
use std::hash::Hash;
|
use std::hash::Hash;
|
||||||
|
use crate::ingestion::flex;
|
||||||
|
|
||||||
#[allow(dead_code)]
|
#[allow(dead_code)]
|
||||||
mod headings {
|
pub(crate) mod headings {
|
||||||
pub const TRANSACTION_ID: usize = 0;
|
pub const TRANSACTION_ID: usize = 0;
|
||||||
pub const DATE: usize = 1;
|
pub const DATE: usize = 1;
|
||||||
pub const TIME: usize = 2;
|
pub const TIME: usize = 2;
|
||||||
@@ -30,19 +31,19 @@ mod headings {
|
|||||||
pub const CATEGORY_SPLIT: usize = 15;
|
pub const CATEGORY_SPLIT: usize = 15;
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Eq, PartialEq, Hash)]
|
#[derive(Debug, Eq, PartialEq, Hash, Clone)]
|
||||||
pub struct MonzoRow {
|
pub struct MonzoRow {
|
||||||
category_split: Option<String>,
|
pub category_split: Option<String>,
|
||||||
primary_category: String,
|
pub primary_category: String,
|
||||||
total_amount: Decimal,
|
pub total_amount: Decimal,
|
||||||
receipt: Option<String>,
|
pub receipt: Option<String>,
|
||||||
notes: Option<String>,
|
pub notes: Option<String>,
|
||||||
emoji: Option<String>,
|
pub emoji: Option<String>,
|
||||||
description: Option<String>,
|
pub description: Option<String>,
|
||||||
transaction_type: String,
|
pub transaction_type: String,
|
||||||
title: Option<String>,
|
pub title: Option<String>,
|
||||||
timestamp: NaiveDateTime,
|
pub timestamp: NaiveDateTime,
|
||||||
transaction_id: String,
|
pub transaction_id: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl MonzoRow {
|
impl MonzoRow {
|
||||||
@@ -77,7 +78,7 @@ impl MonzoRow {
|
|||||||
hasher.finish() as i64
|
hasher.finish() as i64
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn into_insertion(self) -> Result<Insertion, anyhow::Error> {
|
pub fn into_insertion(self, account_id: i32) -> Result<Insertion, anyhow::Error> {
|
||||||
let identity_hash = self.compute_hash();
|
let identity_hash = self.compute_hash();
|
||||||
|
|
||||||
let expenditures: Vec<_> = match &self.category_split {
|
let expenditures: Vec<_> = match &self.category_split {
|
||||||
@@ -106,6 +107,7 @@ impl MonzoRow {
|
|||||||
total_amount: self.total_amount,
|
total_amount: self.total_amount,
|
||||||
description: self.description,
|
description: self.description,
|
||||||
identity_hash: Some(identity_hash),
|
identity_hash: Some(identity_hash),
|
||||||
|
account_id: Some(account_id),
|
||||||
}
|
}
|
||||||
.into_active_model(),
|
.into_active_model(),
|
||||||
|
|
||||||
@@ -138,7 +140,7 @@ fn parse_timestamp(date: &str, time: &str) -> anyhow::Result<NaiveDateTime> {
|
|||||||
Ok(date.and_time(time))
|
Ok(date.and_time(time))
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn from_json_row(row: Vec<Value>) -> anyhow::Result<MonzoRow> {
|
pub fn from_json_row(row: &[Value]) -> anyhow::Result<MonzoRow> {
|
||||||
let date = DateTime::parse_from_rfc3339(row[headings::DATE].as_str().context("No date")?)
|
let date = DateTime::parse_from_rfc3339(row[headings::DATE].as_str().context("No date")?)
|
||||||
.context("Failed to parse date")?;
|
.context("Failed to parse date")?;
|
||||||
|
|
||||||
@@ -176,7 +178,7 @@ fn test_json() {
|
|||||||
|
|
||||||
let json_rows = json
|
let json_rows = json
|
||||||
.iter()
|
.iter()
|
||||||
.map(|row| from_json_row(row.clone()))
|
.map(|row| from_json_row(&row))
|
||||||
.collect::<Result<Vec<_>, anyhow::Error>>()
|
.collect::<Result<Vec<_>, anyhow::Error>>()
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
pub mod db;
|
pub mod db;
|
||||||
pub mod ingestion_logic;
|
pub mod ingestion_logic;
|
||||||
|
pub mod flex;
|
||||||
pub mod routes;
|
pub mod routes;
|
||||||
|
|||||||
+91
-21
@@ -9,46 +9,116 @@ use sea_orm::DatabaseConnection;
|
|||||||
use serde_json::Value;
|
use serde_json::Value;
|
||||||
use std::io::Cursor;
|
use std::io::Cursor;
|
||||||
|
|
||||||
|
#[derive(serde::Deserialize, Debug)]
|
||||||
|
#[serde(untagged)]
|
||||||
|
pub enum MonzoBatchedJsonInput {
|
||||||
|
Legacy(Vec<Vec<Value>>),
|
||||||
|
New {
|
||||||
|
account_id: Option<u8>,
|
||||||
|
rows: Vec<Vec<Value>>,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
impl MonzoBatchedJsonInput {
|
||||||
|
fn account_id(&self) -> Option<u8> {
|
||||||
|
match self {
|
||||||
|
MonzoBatchedJsonInput::Legacy(_) => None,
|
||||||
|
MonzoBatchedJsonInput::New { account_id, .. } => *account_id,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn rows(&self) -> &[Vec<Value>] {
|
||||||
|
match self {
|
||||||
|
MonzoBatchedJsonInput::Legacy(rows) => rows,
|
||||||
|
MonzoBatchedJsonInput::New { rows, .. } => rows,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn monzo_batched_json(
|
pub async fn monzo_batched_json(
|
||||||
Extension(db): Extension<DatabaseConnection>,
|
Extension(db): Extension<DatabaseConnection>,
|
||||||
Json(data): Json<Vec<Vec<Value>>>,
|
Json(data): Json<MonzoBatchedJsonInput>,
|
||||||
) -> Result<&'static str, AppError> {
|
) -> Result<&'static str, AppError> {
|
||||||
let data = data
|
let rows = data
|
||||||
.into_iter()
|
.rows()
|
||||||
.skip(1) // Skip the header row.
|
.iter()
|
||||||
.map(from_json_row)
|
.skip_while(|row| row[0] == Value::String("Transaction ID".to_string()))
|
||||||
|
.map(|row| from_json_row(row.as_ref()))
|
||||||
.collect::<Result<_, _>>()?;
|
.collect::<Result<_, _>>()?;
|
||||||
|
|
||||||
db::insert(&db, data).await?;
|
// We default to the main account for JSON ingestion for now.
|
||||||
|
db::insert(&db, rows, data.account_id().unwrap_or(1) as i32).await?;
|
||||||
|
|
||||||
Ok("Ok")
|
Ok("Ok")
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn extract_csv(mut multipart: Multipart) -> Result<Option<Bytes>, MultipartError> {
|
async fn extract_csv_and_account_name(
|
||||||
let csv = loop {
|
mut multipart: Multipart,
|
||||||
match multipart.next_field().await? {
|
) -> Result<(Option<Bytes>, Option<String>), MultipartError> {
|
||||||
Some(field) if field.name() == Some("csv") => {
|
let mut csv = None;
|
||||||
break Some(field.bytes().await?);
|
let mut account_name = None;
|
||||||
|
|
||||||
|
while let Some(field) = multipart.next_field().await? {
|
||||||
|
match field.name() {
|
||||||
|
Some("csv") => {
|
||||||
|
csv = Some(field.bytes().await?);
|
||||||
}
|
}
|
||||||
|
|
||||||
Some(_) => {}
|
Some("account_id") => {
|
||||||
None => break None,
|
account_name = Some(field.text().await?);
|
||||||
}
|
}
|
||||||
};
|
|
||||||
|
|
||||||
Ok(csv)
|
_ => {}
|
||||||
|
}
|
||||||
|
|
||||||
|
if csv.is_some() && account_name.is_some() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok((csv, account_name))
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(serde::Deserialize, Debug, Clone)]
|
||||||
|
pub struct ShortcutBody {
|
||||||
|
pub body: String,
|
||||||
|
pub account_name: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn shortcuts_csv(
|
||||||
|
Extension(db): Extension<DatabaseConnection>,
|
||||||
|
Json(shortcut_body): Json<ShortcutBody>,
|
||||||
|
) -> Result<&'static str, AppError> {
|
||||||
|
let account_id = db::get_account_id(&db, Some(shortcut_body.account_name)).await?;
|
||||||
|
|
||||||
|
let csv = Cursor::new(shortcut_body.body.as_bytes());
|
||||||
|
let mut csv = csv::Reader::from_reader(csv);
|
||||||
|
let data = csv.records();
|
||||||
|
let data = data
|
||||||
|
.filter_map(|f| f.ok())
|
||||||
|
.map(from_csv_row)
|
||||||
|
.collect::<Result<_, _>>()?;
|
||||||
|
|
||||||
|
db::insert(&db, data, account_id).await?;
|
||||||
|
|
||||||
|
Ok("Ok")
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn monzo_batched_csv(
|
pub async fn monzo_batched_csv(
|
||||||
Extension(db): Extension<DatabaseConnection>,
|
Extension(db): Extension<DatabaseConnection>,
|
||||||
multipart: Multipart,
|
multipart: Multipart,
|
||||||
) -> Result<&'static str, AppError> {
|
) -> Result<&'static str, AppError> {
|
||||||
static CSV_MISSING_ERR_MSG: &'static str = "No CSV file provided. Expected a multipart request with a `csv` field containing the contents of the CSV.";
|
static CSV_MISSING_ERR_MSG: &str = "No CSV file provided. Expected a multipart request with a `csv` field containing the contents of the CSV.";
|
||||||
|
|
||||||
let csv = extract_csv(multipart)
|
let (csv, account_name) = extract_csv_and_account_name(multipart)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| AppError::BadRequest(anyhow!(e)))
|
.map_err(|e| AppError::BadRequest(anyhow!(e)))?;
|
||||||
.and_then(|csv| csv.ok_or(AppError::BadRequest(anyhow!(CSV_MISSING_ERR_MSG))))?;
|
|
||||||
|
let Some(csv) = csv else {
|
||||||
|
return Err(AppError::BadRequest(anyhow!(CSV_MISSING_ERR_MSG)));
|
||||||
|
};
|
||||||
|
|
||||||
|
let account_id = db::get_account_id(&db, account_name).await?;
|
||||||
|
|
||||||
let csv = Cursor::new(csv);
|
let csv = Cursor::new(csv);
|
||||||
let mut csv = csv::Reader::from_reader(csv);
|
let mut csv = csv::Reader::from_reader(csv);
|
||||||
@@ -58,7 +128,7 @@ pub async fn monzo_batched_csv(
|
|||||||
.map(from_csv_row)
|
.map(from_csv_row)
|
||||||
.collect::<Result<_, _>>()?;
|
.collect::<Result<_, _>>()?;
|
||||||
|
|
||||||
db::insert(&db, data).await?;
|
db::insert(&db, data, account_id).await?;
|
||||||
|
|
||||||
Ok("Ok")
|
Ok("Ok")
|
||||||
}
|
}
|
||||||
|
|||||||
+9
-3
@@ -4,7 +4,7 @@ mod ingestion;
|
|||||||
use crate::error::AppError;
|
use crate::error::AppError;
|
||||||
use crate::ingestion::db;
|
use crate::ingestion::db;
|
||||||
use crate::ingestion::ingestion_logic::from_csv_row;
|
use crate::ingestion::ingestion_logic::from_csv_row;
|
||||||
use crate::ingestion::routes::{monzo_batched_csv, monzo_batched_json};
|
use crate::ingestion::routes::{monzo_batched_csv, monzo_batched_json, shortcuts_csv};
|
||||||
use axum::routing::{get, post};
|
use axum::routing::{get, post};
|
||||||
use axum::{Extension, Router};
|
use axum::{Extension, Router};
|
||||||
use clap::{Parser, Subcommand};
|
use clap::{Parser, Subcommand};
|
||||||
@@ -44,6 +44,10 @@ enum Commands {
|
|||||||
Csv {
|
Csv {
|
||||||
/// The path of the CSV file to ingest.
|
/// The path of the CSV file to ingest.
|
||||||
csv_file: PathBuf,
|
csv_file: PathBuf,
|
||||||
|
|
||||||
|
/// The name of the account to ingest the CSV for.
|
||||||
|
#[clap(long, short)]
|
||||||
|
account: String,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -94,7 +98,7 @@ async fn main() -> anyhow::Result<()> {
|
|||||||
serve_web(addr, connection).await?;
|
serve_web(addr, connection).await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
Commands::Csv { csv_file } => {
|
Commands::Csv { csv_file, account: account_name } => {
|
||||||
let mut csv = csv::Reader::from_reader(File::open(csv_file)?);
|
let mut csv = csv::Reader::from_reader(File::open(csv_file)?);
|
||||||
let data = csv.records();
|
let data = csv.records();
|
||||||
let data = data
|
let data = data
|
||||||
@@ -102,7 +106,8 @@ async fn main() -> anyhow::Result<()> {
|
|||||||
.map(from_csv_row)
|
.map(from_csv_row)
|
||||||
.collect::<Result<_, _>>()?;
|
.collect::<Result<_, _>>()?;
|
||||||
|
|
||||||
db::insert(&connection, data).await?;
|
let account_id = db::get_account_id(&connection, Some(account_name)).await?;
|
||||||
|
db::insert(&connection, data, account_id).await?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -114,6 +119,7 @@ async fn serve_web(address: SocketAddr, connection: DatabaseConnection) -> anyho
|
|||||||
.route("/health", get(health_check))
|
.route("/health", get(health_check))
|
||||||
.route("/monzo-batch-export", post(monzo_batched_json))
|
.route("/monzo-batch-export", post(monzo_batched_json))
|
||||||
.route("/monzo-csv-ingestion", post(monzo_batched_csv))
|
.route("/monzo-csv-ingestion", post(monzo_batched_csv))
|
||||||
|
.route("/shortcuts-csv-import", post(shortcuts_csv))
|
||||||
.layer(Extension(connection.clone()))
|
.layer(Extension(connection.clone()))
|
||||||
.layer(TraceLayer::new_for_http());
|
.layer(TraceLayer::new_for_http());
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user