Compare commits
35
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e2ffd900b4 | ||
|
|
4b2f0f3bf7 | ||
|
|
29fe8bee39 | ||
|
|
49a1700706 | ||
|
|
04b58d3075 | ||
|
|
076a573711 | ||
|
|
af0588c5ef | ||
|
|
21a63a10a4 | ||
|
|
08766dc0e0 | ||
|
|
bf47520d31 | ||
|
|
fc1cea32b5 | ||
|
|
92462bd316 | ||
|
|
f0b0cb1567 | ||
|
|
8478fa0b38 | ||
|
|
3b2c1aeda0 | ||
|
|
b37273cfbe | ||
|
|
f344d69419 | ||
|
|
901aba9c7f | ||
|
|
d7d7fa9718 | ||
|
|
b8c1faced2 | ||
|
|
b563bbd02c | ||
|
|
7fd85550ea | ||
|
|
97f57803e5 | ||
|
|
046ce44d23 | ||
|
|
f70d844ff3 | ||
|
|
61de38b9bf | ||
|
|
6c5d3910dc | ||
|
|
4bb9f2813d | ||
|
|
3df05b2d9c | ||
|
|
f19f861297 | ||
|
|
89d0d12e26 | ||
|
|
62d5ad8dc2 | ||
|
|
a2a47b276a | ||
|
|
8ce60ce278 | ||
|
|
fbf473b3b4 |
+82
-13
@@ -1,33 +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: Install Rust
|
||||||
|
uses: actions-rs/toolchain@v1
|
||||||
|
with:
|
||||||
|
toolchain: stable
|
||||||
|
profile: minimal
|
||||||
|
override: true
|
||||||
|
components: rustfmt, clippy
|
||||||
|
|
||||||
|
- name: Add ARM64 target
|
||||||
|
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:
|
||||||
|
path: |
|
||||||
|
~/.cargo
|
||||||
|
target/
|
||||||
|
key: "${{ runner.os }}-cargo-${{ hashFiles('**/Cargo.lock') }}"
|
||||||
|
restore-keys: |
|
||||||
|
${{ runner.os }}-cargo-
|
||||||
|
|
||||||
|
- name: Build (x86_64)
|
||||||
|
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
|
- name: Set up Docker Buildx
|
||||||
uses: docker/setup-buildx-action@v1
|
uses: docker/setup-buildx-action@v2
|
||||||
|
|
||||||
- name: Login to Docker
|
- name: Login to DockerHub
|
||||||
uses: docker/login-action@v1
|
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
|
||||||
tags: git.joshuacoles.me/personal/monzo-ingestion:latest
|
tags: git.joshuacoles.me/${{ github.repository }}:${{ github.sha }},git.joshuacoles.me/${{ github.repository }}:latest
|
||||||
|
build-args: |
|
||||||
|
BINARY_NAME=${{ env.RUST_BINARY_NAME }}
|
||||||
|
|
||||||
|
- 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
+1385
-282
File diff suppressed because it is too large
Load Diff
+23
-9
@@ -9,23 +9,37 @@ migration = { path = "migration" }
|
|||||||
|
|
||||||
axum = { version = "0.7.5", features = ["multipart"] }
|
axum = { version = "0.7.5", features = ["multipart"] }
|
||||||
tokio = { version = "1.37.0", features = ["full"] }
|
tokio = { version = "1.37.0", features = ["full"] }
|
||||||
sea-orm = { version = "1.0.0-rc.4", features = [
|
sea-orm = { version = "1.0.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 = "1.0.86"
|
anyhow = { version = "1.0", features = ["backtrace"] }
|
||||||
thiserror = "1.0.61"
|
thiserror = "1.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.21"
|
||||||
|
testcontainers-modules = { version = "0.9", features = ["postgres"] }
|
||||||
|
sqlx = { version = "0.7", features = ["postgres"] }
|
||||||
|
tower-http = { version = "0.5", features = ["trace"] }
|
||||||
|
bytes = "1.7"
|
||||||
|
once_cell = "1.19"
|
||||||
|
|
||||||
|
tracing-opentelemetry = "0.25.0"
|
||||||
|
opentelemetry = "0.24.0"
|
||||||
|
opentelemetry_sdk = { version = "0.24.1", features = ["rt-tokio"] }
|
||||||
|
opentelemetry-http = { version = "0.13.0", features = ["reqwest"] }
|
||||||
|
opentelemetry-otlp = { version = "0.17.0", features = ["grpc-tonic", "http-json", "tokio"] }
|
||||||
|
opentelemetry-semantic-conventions = "0.16.0"
|
||||||
|
reqwest = "0.12.7"
|
||||||
|
|
||||||
[workspace]
|
[workspace]
|
||||||
members = [".", "migration", "entity"]
|
members = [".", "migration", "entity"]
|
||||||
|
|||||||
+14
-76
@@ -1,78 +1,16 @@
|
|||||||
# syntax=docker/dockerfile:1
|
FROM --platform=$BUILDPLATFORM debian:bullseye-slim AS builder
|
||||||
|
ARG TARGETPLATFORM
|
||||||
# Comments are provided throughout this file to help you get started.
|
ARG BINARY_NAME
|
||||||
# If you need more help, visit the Dockerfile reference guide at
|
|
||||||
# https://docs.docker.com/engine/reference/builder/
|
|
||||||
|
|
||||||
################################################################################
|
|
||||||
# Create a stage for building the application.
|
|
||||||
|
|
||||||
ARG RUST_VERSION=1.76.0
|
|
||||||
ARG APP_NAME=monzo-ingestion
|
|
||||||
FROM rust:${RUST_VERSION}-slim-bullseye AS build
|
|
||||||
ARG APP_NAME
|
|
||||||
WORKDIR /app
|
WORKDIR /app
|
||||||
|
COPY . .
|
||||||
|
RUN case "$TARGETPLATFORM" in \
|
||||||
|
"linux/amd64") BINARY_PATH="target/release/${BINARY_NAME}" ;; \
|
||||||
|
"linux/arm64") BINARY_PATH="target/aarch64-unknown-linux-gnu/release/${BINARY_NAME}" ;; \
|
||||||
|
*) exit 1 ;; \
|
||||||
|
esac && \
|
||||||
|
mv "$BINARY_PATH" /usr/local/bin/${BINARY_NAME}
|
||||||
|
|
||||||
# Build the application.
|
FROM --platform=$TARGETPLATFORM debian:bullseye-slim
|
||||||
# Leverage a cache mount to /usr/local/cargo/registry/
|
ARG BINARY_NAME
|
||||||
# for downloaded dependencies and a cache mount to /app/target/ for
|
COPY --from=builder /usr/local/bin/${BINARY_NAME} /usr/local/bin/
|
||||||
# compiled dependencies which will speed up subsequent builds.
|
CMD ["${BINARY_NAME}"]
|
||||||
# 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"]
|
|
||||||
|
|||||||
@@ -11,11 +11,13 @@ pub struct Model {
|
|||||||
pub transaction_type: String,
|
pub transaction_type: String,
|
||||||
pub total_amount: Decimal,
|
pub total_amount: Decimal,
|
||||||
pub timestamp: DateTime,
|
pub timestamp: DateTime,
|
||||||
pub title: String,
|
pub title: Option<String>,
|
||||||
pub emoji: Option<String>,
|
pub emoji: Option<String>,
|
||||||
pub notes: Option<String>,
|
pub notes: Option<String>,
|
||||||
pub receipt: Option<String>,
|
pub receipt: Option<String>,
|
||||||
pub description: Option<String>,
|
pub description: Option<String>,
|
||||||
|
#[sea_orm(unique)]
|
||||||
|
pub identity_hash: Option<i64>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
|
||||||
|
|||||||
@@ -0,0 +1,3 @@
|
|||||||
|
Transaction ID,Date,Time,Type,Name,Emoji,Category,Amount,Currency,Local amount,Local currency,Notes and #tags,Address,Receipt,Description,Category split
|
||||||
|
tx_0000AhmcbFKVlkAFbeZuAk,11/05/2024,13:41:47,Pot transfer,Savings Pot,,Savings,-1.4,GBP,-1.4,GBP,,,,Round up,
|
||||||
|
tx_0000AhmgdZ0FzkVZou77Ev,11/05/2024,14:27:01,Card payment,D,💊,C,-100.42,GBP,-100.42,GBP,,B,,A,"Groceries:15.00,Personal care:10.00"
|
||||||
|
@@ -0,0 +1,38 @@
|
|||||||
|
[
|
||||||
|
[
|
||||||
|
"tx_0000AhmcbFKVlkAFbeZuAk",
|
||||||
|
"2024-05-11T00:00:00.000Z",
|
||||||
|
"1899-12-30T13:41:47.000Z",
|
||||||
|
"Pot transfer",
|
||||||
|
"Savings Pot",
|
||||||
|
"",
|
||||||
|
"Savings",
|
||||||
|
-1.4,
|
||||||
|
"GBP",
|
||||||
|
-1.4,
|
||||||
|
"GBP",
|
||||||
|
"",
|
||||||
|
"",
|
||||||
|
"",
|
||||||
|
"Round up",
|
||||||
|
""
|
||||||
|
],
|
||||||
|
[
|
||||||
|
"tx_0000AhmgdZ0FzkVZou77Ev",
|
||||||
|
"2024-05-11T00:00:00.000Z",
|
||||||
|
"1899-12-30T14:27:01.000Z",
|
||||||
|
"Card payment",
|
||||||
|
"D",
|
||||||
|
"💊",
|
||||||
|
"C",
|
||||||
|
-100.42,
|
||||||
|
"GBP",
|
||||||
|
-100.42,
|
||||||
|
"GBP",
|
||||||
|
"",
|
||||||
|
"B",
|
||||||
|
"",
|
||||||
|
"A",
|
||||||
|
"Groceries:15.00,Personal care:10.00"
|
||||||
|
]
|
||||||
|
]
|
||||||
@@ -0,0 +1,12 @@
|
|||||||
|
#!/usr/bin/env just --justfile
|
||||||
|
|
||||||
|
generate-migration:
|
||||||
|
sea-orm-cli migrate generate $1
|
||||||
|
|
||||||
|
migrate:
|
||||||
|
sea-orm-cli migrate up
|
||||||
|
|
||||||
|
generate-entities:
|
||||||
|
sea-orm-cli generate entity --with-serde both \
|
||||||
|
-l -o entity/src \
|
||||||
|
--ignore-tables monzo_ingestion_seaql_migrations
|
||||||
@@ -1,6 +1,8 @@
|
|||||||
pub use sea_orm_migration::prelude::*;
|
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 m20240603_162500_make_title_optional;
|
||||||
|
|
||||||
pub struct Migrator;
|
pub struct Migrator;
|
||||||
|
|
||||||
@@ -11,6 +13,10 @@ impl MigratorTrait for Migrator {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn migrations() -> Vec<Box<dyn MigrationTrait>> {
|
fn migrations() -> Vec<Box<dyn MigrationTrait>> {
|
||||||
vec![Box::new(m20230904_141851_create_monzo_tables::Migration)]
|
vec![
|
||||||
|
Box::new(m20230904_141851_create_monzo_tables::Migration),
|
||||||
|
Box::new(m20240529_195030_add_transaction_identity_hash::Migration),
|
||||||
|
Box::new(m20240603_162500_make_title_optional::Migration),
|
||||||
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,39 @@
|
|||||||
|
use sea_orm_migration::prelude::*;
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
TableAlterStatement::new()
|
||||||
|
.table(Transaction::Table)
|
||||||
|
.add_column(
|
||||||
|
ColumnDef::new(Transaction::IdentityHash)
|
||||||
|
.big_integer()
|
||||||
|
.unique_key(),
|
||||||
|
)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
TableAlterStatement::new()
|
||||||
|
.table(Transaction::Table)
|
||||||
|
.drop_column(Transaction::IdentityHash)
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Transaction {
|
||||||
|
Table,
|
||||||
|
IdentityHash,
|
||||||
|
}
|
||||||
@@ -0,0 +1,61 @@
|
|||||||
|
use sea_orm_migration::prelude::*;
|
||||||
|
|
||||||
|
#[derive(DeriveMigrationName)]
|
||||||
|
pub struct Migration;
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl MigrationTrait for Migration {
|
||||||
|
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
TableAlterStatement::new()
|
||||||
|
.table(Transaction::Table)
|
||||||
|
.modify_column(ColumnDef::new(Transaction::Title).string().null())
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
// Set all empty string titles to null
|
||||||
|
manager
|
||||||
|
.get_connection()
|
||||||
|
.execute_unprepared(
|
||||||
|
r#"
|
||||||
|
update transaction
|
||||||
|
set title = null
|
||||||
|
where title = ''
|
||||||
|
"#,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
|
||||||
|
// Set all null titles to empty string when reverting
|
||||||
|
manager
|
||||||
|
.get_connection()
|
||||||
|
.execute_unprepared(
|
||||||
|
r#"
|
||||||
|
update transaction
|
||||||
|
set title = ''
|
||||||
|
where title is null
|
||||||
|
"#,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
manager
|
||||||
|
.alter_table(
|
||||||
|
TableAlterStatement::new()
|
||||||
|
.table(Transaction::Table)
|
||||||
|
.modify_column(ColumnDef::new(Transaction::Title).string().not_null())
|
||||||
|
.to_owned(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(DeriveIden)]
|
||||||
|
enum Transaction {
|
||||||
|
Table,
|
||||||
|
Title,
|
||||||
|
}
|
||||||
+13
-5
@@ -6,25 +6,33 @@ use tracing::log::error;
|
|||||||
#[derive(thiserror::Error, Debug)]
|
#[derive(thiserror::Error, Debug)]
|
||||||
pub enum AppError {
|
pub enum AppError {
|
||||||
/// SeaORM error, separated for ease of use allowing us to `?` db operations.
|
/// SeaORM error, separated for ease of use allowing us to `?` db operations.
|
||||||
#[error("Internal error")]
|
#[error("Database error: {0}")]
|
||||||
DbError(#[from] DbErr),
|
DbError(#[from] DbErr),
|
||||||
|
|
||||||
#[error("Invalid request {0}")]
|
#[error("Invalid request {0}")]
|
||||||
BadRequest(anyhow::Error),
|
BadRequest(anyhow::Error),
|
||||||
|
|
||||||
/// Catch all for error we don't care to expose publicly.
|
/// Catch all for error we don't care to expose publicly.
|
||||||
#[error("Internal error")]
|
#[error("An error occurred: {0}")]
|
||||||
Anyhow(#[from] anyhow::Error),
|
Anyhow(#[from] anyhow::Error),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl AppError {
|
||||||
|
fn to_response_string(&self) -> String {
|
||||||
|
match self {
|
||||||
|
AppError::BadRequest(e) => e.to_string(),
|
||||||
|
_ => "Internal server error".to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl IntoResponse for AppError {
|
impl IntoResponse for AppError {
|
||||||
fn into_response(self) -> Response {
|
fn into_response(self) -> Response {
|
||||||
error!("Internal server error: {self:?}");
|
|
||||||
|
|
||||||
let status_code = match self {
|
let status_code = match self {
|
||||||
|
AppError::BadRequest(_) => StatusCode::BAD_REQUEST,
|
||||||
_ => StatusCode::INTERNAL_SERVER_ERROR,
|
_ => StatusCode::INTERNAL_SERVER_ERROR,
|
||||||
};
|
};
|
||||||
|
|
||||||
(status_code, self.to_string()).into_response()
|
(status_code, self.to_response_string()).into_response()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+254
-40
@@ -1,56 +1,74 @@
|
|||||||
|
use crate::error::AppError;
|
||||||
|
use crate::ingestion::ingestion_logic::MonzoRow;
|
||||||
use anyhow::anyhow;
|
use anyhow::anyhow;
|
||||||
use entity::{expenditure, transaction};
|
use entity::{expenditure, transaction};
|
||||||
use sea_orm::sea_query::OnConflict;
|
use sea_orm::sea_query::OnConflict;
|
||||||
use sea_orm::{ConnectionTrait, DatabaseBackend, QueryFilter, QueryTrait, Statement};
|
use sea_orm::{
|
||||||
use sea_orm::{ColumnTrait, DatabaseConnection, EntityTrait, Iterable, TransactionTrait};
|
ColumnTrait, DatabaseConnection, DbErr, EntityTrait, Iterable, QuerySelect, TransactionTrait,
|
||||||
use migration::PostgresQueryBuilder;
|
};
|
||||||
use crate::error::AppError;
|
use sea_orm::{ConnectionTrait, DatabaseTransaction, QueryFilter};
|
||||||
|
|
||||||
pub struct Insertion {
|
pub struct Insertion {
|
||||||
pub transaction: transaction::ActiveModel,
|
pub transaction: transaction::ActiveModel,
|
||||||
pub contained_expenditures: Vec<expenditure::ActiveModel>,
|
pub contained_expenditures: Vec<expenditure::ActiveModel>,
|
||||||
|
pub identity_hash: i64,
|
||||||
}
|
}
|
||||||
|
|
||||||
// Note while this is more efficient in db calls, it does bind together the entire group.
|
// Note while this is more efficient in db calls, it does bind together the entire group.
|
||||||
// We employ a batching process for now to try balance speed and failure rate, but it is worth
|
// We employ a batching process for now to try balance speed and failure rate, but it is worth
|
||||||
// trying to move failures earlier and improve reporting.
|
// trying to move failures earlier and improve reporting.
|
||||||
pub async fn insert(db: &DatabaseConnection, insertions: Vec<Insertion>) -> Result<Vec<String>, AppError> {
|
pub async fn insert(
|
||||||
|
db: &DatabaseConnection,
|
||||||
|
monzo_rows: Vec<MonzoRow>,
|
||||||
|
) -> Result<Vec<String>, AppError> {
|
||||||
let mut new_transaction_ids = Vec::new();
|
let mut new_transaction_ids = Vec::new();
|
||||||
|
let insertions = monzo_rows
|
||||||
|
.into_iter()
|
||||||
|
.map(MonzoRow::into_insertion)
|
||||||
|
.collect::<Result<Vec<_>, _>>()?;
|
||||||
|
|
||||||
for insertions in insertions.chunks(400) {
|
for insertions in insertions.chunks(400) {
|
||||||
|
let (new_or_updated_insertions, inserted_transaction_ids) =
|
||||||
|
whittle_insertions(insertions, db).await?;
|
||||||
|
|
||||||
|
if new_or_updated_insertions.is_empty() {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
let tx = db.begin().await?;
|
let tx = db.begin().await?;
|
||||||
|
update_transactions(&tx, &new_or_updated_insertions).await?;
|
||||||
|
update_expenditures(&tx, &new_or_updated_insertions, &inserted_transaction_ids).await?;
|
||||||
|
tx.commit().await?;
|
||||||
|
|
||||||
let insert = transaction::Entity::insert_many(insertions.iter().map(|i| &i.transaction).cloned())
|
// We wait until the transaction is committed before adding the new transaction ids to the
|
||||||
.on_conflict(
|
// list to avoid issues with the transaction being rolled back.
|
||||||
OnConflict::column(transaction::Column::Id)
|
new_transaction_ids.extend(inserted_transaction_ids);
|
||||||
.update_columns(transaction::Column::iter())
|
}
|
||||||
.to_owned(),
|
|
||||||
)
|
|
||||||
.into_query()
|
|
||||||
.returning_col(transaction::Column::Id)
|
|
||||||
.build(PostgresQueryBuilder);
|
|
||||||
|
|
||||||
let inserted_transaction_ids = tx.query_all(Statement::from_sql_and_values(
|
// Notify the new transactions once everything is committed.
|
||||||
DatabaseBackend::Postgres,
|
notify_new_transactions(db, &new_transaction_ids).await?;
|
||||||
insert.0,
|
|
||||||
insert.1,
|
Ok(new_transaction_ids)
|
||||||
)).await?
|
}
|
||||||
.iter()
|
|
||||||
.map(|r| r.try_get_by("id"))
|
async fn update_expenditures(
|
||||||
.collect::<Result<Vec<String>, _>>()?;
|
tx: &DatabaseTransaction,
|
||||||
|
new_or_updated_insertions: &[&Insertion],
|
||||||
|
inserted_transaction_ids: &[String],
|
||||||
|
) -> Result<(), AppError> {
|
||||||
|
if new_or_updated_insertions.is_empty() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
// Expenditures can change as we re-categorise them, so we delete all the old ones and
|
// Expenditures can change as we re-categorise them, so we delete all the old ones and
|
||||||
// insert an entirely new set to ensure we don't end up leaving old ones around.
|
// insert an entirely new set to ensure we don't end up leaving old ones around.
|
||||||
expenditure::Entity::delete_many()
|
expenditure::Entity::delete_many()
|
||||||
.filter(
|
.filter(expenditure::Column::TransactionId.is_in(inserted_transaction_ids))
|
||||||
expenditure::Column::TransactionId
|
.exec(tx)
|
||||||
.is_in(insertions.iter().map(|i| i.transaction.id.as_ref())),
|
|
||||||
)
|
|
||||||
.exec(&tx)
|
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
expenditure::Entity::insert_many(
|
expenditure::Entity::insert_many(
|
||||||
insertions
|
new_or_updated_insertions
|
||||||
.iter()
|
.iter()
|
||||||
.flat_map(|i| &i.contained_expenditures)
|
.flat_map(|i| &i.contained_expenditures)
|
||||||
.cloned(),
|
.cloned(),
|
||||||
@@ -63,23 +81,219 @@ pub async fn insert(db: &DatabaseConnection, insertions: Vec<Insertion>) -> Resu
|
|||||||
.update_columns(expenditure::Column::iter())
|
.update_columns(expenditure::Column::iter())
|
||||||
.to_owned(),
|
.to_owned(),
|
||||||
)
|
)
|
||||||
.exec(&tx)
|
.exec(tx)
|
||||||
.await?;
|
.await?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
tx.commit().await?;
|
async fn update_transactions(
|
||||||
new_transaction_ids.extend(inserted_transaction_ids);
|
tx: &DatabaseTransaction,
|
||||||
|
new_or_updated_insertions: &[&Insertion],
|
||||||
|
) -> Result<(), DbErr> {
|
||||||
|
if new_or_updated_insertions.is_empty() {
|
||||||
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
||||||
let payload = serde_json::to_string(&new_transaction_ids)
|
let transactions = new_or_updated_insertions
|
||||||
.map_err(|e| anyhow!(e))?;
|
.iter()
|
||||||
|
.map(|i| &i.transaction)
|
||||||
|
.cloned();
|
||||||
|
|
||||||
db.execute(
|
transaction::Entity::insert_many(transactions)
|
||||||
Statement::from_sql_and_values(
|
.on_conflict(
|
||||||
DatabaseBackend::Postgres,
|
OnConflict::column(transaction::Column::Id)
|
||||||
"NOTIFY monzo_new_transactions, $1",
|
.update_columns(transaction::Column::iter())
|
||||||
vec![sea_orm::Value::from(payload)],
|
.to_owned(),
|
||||||
)
|
)
|
||||||
).await?;
|
.exec(tx)
|
||||||
|
.await?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
Ok(new_transaction_ids)
|
async fn whittle_insertions<'a>(
|
||||||
|
insertions: &'a [Insertion],
|
||||||
|
tx: &DatabaseConnection,
|
||||||
|
) -> Result<(Vec<&'a Insertion>, Vec<String>), AppError> {
|
||||||
|
let existing_hashes = transaction::Entity::find()
|
||||||
|
.select_only()
|
||||||
|
.columns([transaction::Column::IdentityHash])
|
||||||
|
.filter(transaction::Column::IdentityHash.is_not_null())
|
||||||
|
.into_tuple::<(i64, )>()
|
||||||
|
.all(tx)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
tracing::debug!("Found existing entries: {existing_hashes:?}");
|
||||||
|
|
||||||
|
// We will only update those where the hash is different to avoid unnecessary updates and
|
||||||
|
// notifications.
|
||||||
|
let new_or_updated_insertions = insertions
|
||||||
|
.iter()
|
||||||
|
.filter(|i| {
|
||||||
|
let hash = i.identity_hash;
|
||||||
|
!existing_hashes
|
||||||
|
.iter()
|
||||||
|
.any(|(existing_hash, )| *existing_hash == hash)
|
||||||
|
})
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
let inserted_transaction_ids = new_or_updated_insertions
|
||||||
|
.iter()
|
||||||
|
.map(|i| i.transaction.id.clone().unwrap())
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
Ok((new_or_updated_insertions, inserted_transaction_ids))
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn notify_new_transactions(
|
||||||
|
db: &DatabaseConnection,
|
||||||
|
new_transaction_ids: &[String],
|
||||||
|
) -> Result<(), AppError> {
|
||||||
|
let payload = serde_json::to_string(&new_transaction_ids).map_err(|e| anyhow!(e))?;
|
||||||
|
|
||||||
|
db.execute_unprepared(&format!(r#"NOTIFY monzo_new_transactions, '{payload}'"#))
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
mod tests {
|
||||||
|
use super::{insert, notify_new_transactions, update_expenditures, update_transactions};
|
||||||
|
use anyhow::Error;
|
||||||
|
use tokio::sync::OnceCell;
|
||||||
|
use migration::MigratorTrait;
|
||||||
|
use sea_orm::{DatabaseConnection, TransactionTrait};
|
||||||
|
use serde_json::Value;
|
||||||
|
use sqlx::postgres::PgListener;
|
||||||
|
use sqlx::PgPool;
|
||||||
|
use testcontainers::runners::AsyncRunner;
|
||||||
|
use testcontainers::ContainerAsync;
|
||||||
|
use crate::ingestion::ingestion_logic::from_json_row;
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct DatabaseInstance {
|
||||||
|
container: ContainerAsync<testcontainers_modules::postgres::Postgres>,
|
||||||
|
db: DatabaseConnection,
|
||||||
|
pool: PgPool,
|
||||||
|
}
|
||||||
|
|
||||||
|
static INSTANCE: OnceCell<DatabaseInstance> = OnceCell::const_new();
|
||||||
|
|
||||||
|
async fn initialise_db() -> Result<DatabaseInstance, Error> {
|
||||||
|
let container = testcontainers_modules::postgres::Postgres::default()
|
||||||
|
.start()
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
// prepare connection string
|
||||||
|
let connection_string = &format!(
|
||||||
|
"postgres://postgres:postgres@127.0.0.1:{}/postgres",
|
||||||
|
container.get_host_port_ipv4(5432).await?
|
||||||
|
);
|
||||||
|
|
||||||
|
let db: DatabaseConnection = sea_orm::Database::connect(connection_string).await?;
|
||||||
|
migration::Migrator::up(&db, None).await?;
|
||||||
|
|
||||||
|
let pool = PgPool::connect(connection_string).await?;
|
||||||
|
let instance = DatabaseInstance {
|
||||||
|
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)
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_empty_insertion_list() -> Result<(), Error> {
|
||||||
|
let db = get_or_initialize_db_instance().await?;
|
||||||
|
let insertions = vec![];
|
||||||
|
let tx = db.db.begin().await?;
|
||||||
|
update_transactions(&tx, &insertions).await?;
|
||||||
|
update_expenditures(&tx, &insertions, &vec![]).await?;
|
||||||
|
tx.commit().await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_notify() -> Result<(), Error> {
|
||||||
|
let dbi = get_or_initialize_db_instance().await?;
|
||||||
|
let mut listener = PgListener::connect_with(&dbi.pool).await?;
|
||||||
|
listener.listen("monzo_new_transactions").await?;
|
||||||
|
|
||||||
|
let ids = vec![
|
||||||
|
"test1".to_string(),
|
||||||
|
"test2".to_string(),
|
||||||
|
"test3".to_string(),
|
||||||
|
];
|
||||||
|
|
||||||
|
notify_new_transactions(&dbi.db, &ids).await?;
|
||||||
|
|
||||||
|
let notification = listener.recv().await?;
|
||||||
|
let payload = notification.payload();
|
||||||
|
println!("Payload: {}", payload);
|
||||||
|
assert_eq!(
|
||||||
|
serde_json::from_str::<Vec<String>>(&payload)?,
|
||||||
|
ids,
|
||||||
|
"Payloads do not match"
|
||||||
|
);
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_notify_on_insert() -> Result<(), Error> {
|
||||||
|
let dbi = get_or_initialize_db_instance().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)?;
|
||||||
|
let data = json
|
||||||
|
.iter()
|
||||||
|
.map(|row| from_json_row(row.clone()))
|
||||||
|
.collect::<Result<Vec<_>, anyhow::Error>>()
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
insert(&dbi.db, data.clone()).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()).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()).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(())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,9 @@
|
|||||||
|
#[allow(dead_code)]
|
||||||
|
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;
|
||||||
|
}
|
||||||
@@ -1,16 +1,17 @@
|
|||||||
use crate::ingestion::db::Insertion;
|
use crate::ingestion::db::Insertion;
|
||||||
use anyhow::Context;
|
use anyhow::{anyhow, Context};
|
||||||
use chrono::{DateTime, NaiveDate, NaiveDateTime, NaiveTime};
|
use chrono::{DateTime, NaiveDate, NaiveDateTime, NaiveTime};
|
||||||
use csv::StringRecord;
|
use csv::StringRecord;
|
||||||
use entity::expenditure::ActiveModel;
|
use entity::expenditure::ActiveModel;
|
||||||
use entity::transaction;
|
use entity::transaction;
|
||||||
use num_traits::FromPrimitive;
|
use num_traits::FromPrimitive;
|
||||||
use sea_orm::prelude::Decimal;
|
use sea_orm::prelude::Decimal;
|
||||||
use sea_orm::ActiveValue::*;
|
|
||||||
use sea_orm::IntoActiveModel;
|
use sea_orm::IntoActiveModel;
|
||||||
|
use serde_json::Value;
|
||||||
|
use std::hash::Hash;
|
||||||
|
|
||||||
#[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;
|
||||||
@@ -29,7 +30,23 @@ mod headings {
|
|||||||
pub const CATEGORY_SPLIT: usize = 15;
|
pub const CATEGORY_SPLIT: usize = 15;
|
||||||
}
|
}
|
||||||
|
|
||||||
fn parse_section(monzo_transaction_id: &str, section: &str) -> anyhow::Result<ActiveModel> {
|
#[derive(Debug, Eq, PartialEq, Hash, Clone)]
|
||||||
|
pub struct MonzoRow {
|
||||||
|
pub category_split: Option<String>,
|
||||||
|
pub primary_category: String,
|
||||||
|
pub total_amount: Decimal,
|
||||||
|
pub receipt: Option<String>,
|
||||||
|
pub notes: Option<String>,
|
||||||
|
pub emoji: Option<String>,
|
||||||
|
pub description: Option<String>,
|
||||||
|
pub transaction_type: String,
|
||||||
|
pub title: Option<String>,
|
||||||
|
pub timestamp: NaiveDateTime,
|
||||||
|
pub transaction_id: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl MonzoRow {
|
||||||
|
fn parse_section(monzo_transaction_id: &str, section: &str) -> anyhow::Result<ActiveModel> {
|
||||||
let mut components = section.split(':');
|
let mut components = section.split(':');
|
||||||
let category: String = components
|
let category: String = components
|
||||||
.next()
|
.next()
|
||||||
@@ -47,23 +64,81 @@ fn parse_section(monzo_transaction_id: &str, section: &str) -> anyhow::Result<Ac
|
|||||||
amount,
|
amount,
|
||||||
}
|
}
|
||||||
.into_active_model())
|
.into_active_model())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Compute a hash of this row, returning the number as an i64 to be used as a unique constraint
|
||||||
|
/// in the database.
|
||||||
|
pub fn compute_hash(&self) -> i64 {
|
||||||
|
use std::collections::hash_map::DefaultHasher;
|
||||||
|
use std::hash::Hasher;
|
||||||
|
|
||||||
|
let mut hasher = DefaultHasher::new();
|
||||||
|
self.hash(&mut hasher);
|
||||||
|
hasher.finish() as i64
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn into_insertion(self) -> Result<Insertion, anyhow::Error> {
|
||||||
|
let identity_hash = self.compute_hash();
|
||||||
|
|
||||||
|
let expenditures: Vec<_> = match &self.category_split {
|
||||||
|
Some(split) if !split.is_empty() => split
|
||||||
|
.split(',')
|
||||||
|
.map(|section| Self::parse_section(&self.transaction_id, section))
|
||||||
|
.collect::<Result<Vec<_>, anyhow::Error>>()?,
|
||||||
|
|
||||||
|
_ => vec![entity::expenditure::Model {
|
||||||
|
category: self.primary_category.clone(),
|
||||||
|
amount: self.total_amount,
|
||||||
|
transaction_id: self.transaction_id.clone(),
|
||||||
|
}
|
||||||
|
.into_active_model()],
|
||||||
|
};
|
||||||
|
|
||||||
|
Ok(Insertion {
|
||||||
|
transaction: transaction::Model {
|
||||||
|
id: self.transaction_id,
|
||||||
|
transaction_type: self.transaction_type,
|
||||||
|
timestamp: self.timestamp,
|
||||||
|
title: self.title,
|
||||||
|
emoji: self.emoji,
|
||||||
|
notes: self.notes,
|
||||||
|
receipt: self.receipt,
|
||||||
|
total_amount: self.total_amount,
|
||||||
|
description: self.description,
|
||||||
|
identity_hash: Some(identity_hash),
|
||||||
|
}
|
||||||
|
.into_active_model(),
|
||||||
|
|
||||||
|
contained_expenditures: expenditures,
|
||||||
|
identity_hash,
|
||||||
|
})
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn json_opt(value: &serde_json::Value) -> Option<String> {
|
fn json_opt(value: &Value) -> Option<String> {
|
||||||
match value {
|
match value {
|
||||||
serde_json::Value::String(string) if string.is_empty() => None,
|
Value::String(string) if string.is_empty() => None,
|
||||||
serde_json::Value::String(string) => Some(string.to_string()),
|
Value::String(string) => Some(string.to_string()),
|
||||||
_ => None,
|
_ => None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn from_json_row(row: Vec<serde_json::Value>) -> anyhow::Result<Insertion> {
|
fn json_required_str(value: &Value, label: &str) -> anyhow::Result<String> {
|
||||||
use serde_json::Value;
|
match value {
|
||||||
let monzo_transaction_id = row[headings::TRANSACTION_ID]
|
Value::String(string) if string.is_empty() => Err(anyhow!("{} is empty", label)),
|
||||||
.as_str()
|
Value::String(string) => Ok(string.to_string()),
|
||||||
.context("No transaction id")?
|
_ => Err(anyhow!("{} is not a string", label)),
|
||||||
.to_string();
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn parse_timestamp(date: &str, time: &str) -> anyhow::Result<NaiveDateTime> {
|
||||||
|
let date = NaiveDate::parse_from_str(date, "%Y-%m-%d")?;
|
||||||
|
let time = NaiveTime::parse_from_str(time, "%H:%M:%S")?;
|
||||||
|
|
||||||
|
Ok(date.and_time(time))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn from_json_row(row: Vec<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")?;
|
||||||
|
|
||||||
@@ -73,57 +148,57 @@ pub fn from_json_row(row: Vec<serde_json::Value>) -> anyhow::Result<Insertion> {
|
|||||||
|
|
||||||
let timestamp = date.date_naive().and_time(time);
|
let timestamp = date.date_naive().and_time(time);
|
||||||
|
|
||||||
let title = row[headings::NAME]
|
|
||||||
.as_str()
|
|
||||||
.context("No title")?
|
|
||||||
.to_string();
|
|
||||||
|
|
||||||
let monzo_transaction_type = row[headings::TYPE]
|
|
||||||
.as_str()
|
|
||||||
.context("No transaction type")?
|
|
||||||
.to_string();
|
|
||||||
|
|
||||||
let description = json_opt(&row[headings::DESCRIPTION]);
|
|
||||||
let emoji = json_opt(&row[headings::EMOJI]);
|
|
||||||
let notes = json_opt(&row[headings::NOTES_AND_TAGS]);
|
|
||||||
let receipt = json_opt(&row[headings::RECEIPT]);
|
|
||||||
let total_amount = Decimal::from_f64(row[headings::AMOUNT].as_f64().context("No amount")?)
|
let total_amount = Decimal::from_f64(row[headings::AMOUNT].as_f64().context("No amount")?)
|
||||||
.context("Failed to parse date")?;
|
.context("Failed to parse date")?;
|
||||||
|
|
||||||
let expenditures: Vec<_> = match row.get(headings::CATEGORY_SPLIT) {
|
Ok(MonzoRow {
|
||||||
Some(Value::String(split)) if !split.is_empty() => split
|
transaction_id: json_required_str(&row[headings::TRANSACTION_ID], "Transaction ID")?,
|
||||||
.split(',')
|
title: json_opt(&row[headings::NAME]),
|
||||||
.map(|section| parse_section(&monzo_transaction_id, section))
|
transaction_type: json_required_str(&row[headings::TYPE], "Transaction type")?,
|
||||||
.collect::<Result<Vec<_>, anyhow::Error>>()?,
|
description: json_opt(&row[headings::DESCRIPTION]),
|
||||||
|
emoji: json_opt(&row[headings::EMOJI]),
|
||||||
_ => vec![entity::expenditure::Model {
|
notes: json_opt(&row[headings::NOTES_AND_TAGS]),
|
||||||
category: row[headings::CATEGORY]
|
receipt: json_opt(&row[headings::RECEIPT]),
|
||||||
.as_str()
|
primary_category: json_required_str(&row[headings::CATEGORY], "Primary Category")?,
|
||||||
.context("No context")?
|
category_split: json_opt(&row[headings::CATEGORY_SPLIT]),
|
||||||
.to_string(),
|
total_amount,
|
||||||
amount: total_amount,
|
timestamp,
|
||||||
transaction_id: monzo_transaction_id.clone(),
|
|
||||||
}
|
|
||||||
.into_active_model()],
|
|
||||||
};
|
|
||||||
|
|
||||||
Ok(Insertion {
|
|
||||||
transaction: transaction::ActiveModel {
|
|
||||||
id: Set(monzo_transaction_id),
|
|
||||||
transaction_type: Set(monzo_transaction_type),
|
|
||||||
timestamp: Set(timestamp),
|
|
||||||
title: Set(title),
|
|
||||||
emoji: Set(emoji),
|
|
||||||
notes: Set(notes),
|
|
||||||
receipt: Set(receipt),
|
|
||||||
total_amount: Set(total_amount),
|
|
||||||
description: Set(description),
|
|
||||||
},
|
|
||||||
|
|
||||||
contained_expenditures: expenditures,
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_json() {
|
||||||
|
let json = include_str!("../../fixtures/transactions.json");
|
||||||
|
let csv = include_str!("../../fixtures/transactions.csv");
|
||||||
|
|
||||||
|
let json: Vec<Vec<Value>> = serde_json::from_str(json).unwrap();
|
||||||
|
let mut csv_reader = csv::Reader::from_reader(csv.as_bytes());
|
||||||
|
|
||||||
|
let json_rows = json
|
||||||
|
.iter()
|
||||||
|
.map(|row| from_json_row(row.clone()))
|
||||||
|
.collect::<Result<Vec<_>, anyhow::Error>>()
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let csv_rows = csv_reader
|
||||||
|
.records()
|
||||||
|
.map(|record| from_csv_row(record.unwrap()))
|
||||||
|
.collect::<Result<Vec<_>, anyhow::Error>>()
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert_eq!(csv_rows.len(), json_rows.len(), "Different number of rows");
|
||||||
|
|
||||||
|
for (i, (json_row, csv_row)) in json_rows.iter().zip(csv_rows.iter()).enumerate() {
|
||||||
|
assert_eq!(json_row, csv_row, "Row {} is different", i);
|
||||||
|
assert_eq!(
|
||||||
|
json_row.compute_hash(),
|
||||||
|
csv_row.compute_hash(),
|
||||||
|
"Row {} hash are different",
|
||||||
|
i
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn csv_opt(s: &str) -> Option<String> {
|
fn csv_opt(s: &str) -> Option<String> {
|
||||||
match s {
|
match s {
|
||||||
"" => None,
|
"" => None,
|
||||||
@@ -131,9 +206,7 @@ fn csv_opt(s: &str) -> Option<String> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn from_csv_row(row: StringRecord) -> anyhow::Result<Insertion> {
|
pub fn from_csv_row(row: StringRecord) -> anyhow::Result<MonzoRow> {
|
||||||
let monzo_transaction_id = row[headings::TRANSACTION_ID].to_string();
|
|
||||||
|
|
||||||
let date = NaiveDate::parse_from_str(&row[headings::DATE], "%d/%m/%Y")
|
let date = NaiveDate::parse_from_str(&row[headings::DATE], "%d/%m/%Y")
|
||||||
.context("Failed to parse date from csv")?;
|
.context("Failed to parse date from csv")?;
|
||||||
|
|
||||||
@@ -142,41 +215,17 @@ pub fn from_csv_row(row: StringRecord) -> anyhow::Result<Insertion> {
|
|||||||
|
|
||||||
let timestamp = NaiveDateTime::new(date, time);
|
let timestamp = NaiveDateTime::new(date, time);
|
||||||
|
|
||||||
let title = row[headings::NAME].to_string();
|
Ok(MonzoRow {
|
||||||
let monzo_transaction_type = row[headings::TYPE].to_string();
|
timestamp,
|
||||||
|
transaction_id: row[headings::TRANSACTION_ID].to_string(),
|
||||||
let description = csv_opt(&row[headings::DESCRIPTION]);
|
title: csv_opt(&row[headings::NAME]),
|
||||||
let emoji = csv_opt(&row[headings::EMOJI]);
|
transaction_type: row[headings::TYPE].to_string(),
|
||||||
let notes = csv_opt(&row[headings::NOTES_AND_TAGS]);
|
description: csv_opt(&row[headings::DESCRIPTION]),
|
||||||
let receipt = csv_opt(&row[headings::RECEIPT]);
|
emoji: csv_opt(&row[headings::EMOJI]),
|
||||||
let total_amount = row[headings::AMOUNT].parse::<Decimal>()?;
|
notes: csv_opt(&row[headings::NOTES_AND_TAGS]),
|
||||||
|
receipt: csv_opt(&row[headings::RECEIPT]),
|
||||||
let expenditures: Vec<_> = match row.get(headings::CATEGORY_SPLIT) {
|
total_amount: row[headings::AMOUNT].parse::<Decimal>()?,
|
||||||
Some(split) if !split.is_empty() => split
|
category_split: csv_opt(&row[headings::CATEGORY_SPLIT]),
|
||||||
.split(',')
|
primary_category: row[headings::CATEGORY].to_string(),
|
||||||
.map(|section| parse_section(&monzo_transaction_id, section))
|
|
||||||
.collect::<Result<Vec<_>, anyhow::Error>>()?,
|
|
||||||
|
|
||||||
_ => vec![entity::expenditure::Model {
|
|
||||||
transaction_id: monzo_transaction_id.clone(),
|
|
||||||
category: row[headings::CATEGORY].to_string(),
|
|
||||||
amount: total_amount,
|
|
||||||
}
|
|
||||||
.into_active_model()],
|
|
||||||
};
|
|
||||||
|
|
||||||
Ok(Insertion {
|
|
||||||
transaction: transaction::ActiveModel {
|
|
||||||
id: Set(monzo_transaction_id),
|
|
||||||
transaction_type: Set(monzo_transaction_type),
|
|
||||||
timestamp: Set(timestamp),
|
|
||||||
title: Set(title),
|
|
||||||
emoji: Set(emoji),
|
|
||||||
notes: Set(notes),
|
|
||||||
receipt: Set(receipt),
|
|
||||||
total_amount: Set(total_amount),
|
|
||||||
description: Set(description),
|
|
||||||
},
|
|
||||||
contained_expenditures: expenditures,
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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;
|
||||||
|
|||||||
+27
-30
@@ -2,43 +2,33 @@ use crate::error::AppError;
|
|||||||
use crate::ingestion::db;
|
use crate::ingestion::db;
|
||||||
use crate::ingestion::ingestion_logic::{from_csv_row, from_json_row};
|
use crate::ingestion::ingestion_logic::{from_csv_row, from_json_row};
|
||||||
use anyhow::anyhow;
|
use anyhow::anyhow;
|
||||||
|
use axum::extract::multipart::MultipartError;
|
||||||
use axum::extract::{Extension, Json, Multipart};
|
use axum::extract::{Extension, Json, Multipart};
|
||||||
|
use bytes::Bytes;
|
||||||
use sea_orm::DatabaseConnection;
|
use sea_orm::DatabaseConnection;
|
||||||
use serde_json::Value;
|
use serde_json::Value;
|
||||||
use std::io::Cursor;
|
use std::io::Cursor;
|
||||||
|
|
||||||
pub async fn monzo_updated(
|
|
||||||
Extension(db): Extension<DatabaseConnection>,
|
|
||||||
Json(row): Json<Vec<Value>>,
|
|
||||||
) -> Result<&'static str, AppError> {
|
|
||||||
db::insert(&db, vec![from_json_row(row)?]).await.unwrap();
|
|
||||||
|
|
||||||
Ok("Ok")
|
|
||||||
}
|
|
||||||
|
|
||||||
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<Vec<Vec<Value>>>,
|
||||||
) -> Result<&'static str, AppError> {
|
) -> Result<&'static str, AppError> {
|
||||||
let insertions = data
|
let data = data
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.skip(1)
|
.skip(1) // Skip the header row.
|
||||||
.map(|row| from_json_row(row))
|
.map(from_json_row)
|
||||||
.collect::<Result<_, _>>()?;
|
.collect::<Result<_, _>>()?;
|
||||||
|
|
||||||
db::insert(&db, insertions).await.unwrap();
|
db::insert(&db, data).await?;
|
||||||
|
|
||||||
Ok("Ok")
|
Ok("Ok")
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn monzo_batched_csv(
|
async fn extract_csv(mut multipart: Multipart) -> Result<Option<Bytes>, MultipartError> {
|
||||||
Extension(db): Extension<DatabaseConnection>,
|
|
||||||
mut multipart: Multipart,
|
|
||||||
) -> Result<&'static str, AppError> {
|
|
||||||
let csv = loop {
|
let csv = loop {
|
||||||
match multipart.next_field().await.unwrap() {
|
match multipart.next_field().await? {
|
||||||
Some(field) if field.name() == Some("csv") => {
|
Some(field) if field.name() == Some("csv") => {
|
||||||
break Some(field.bytes().await.unwrap());
|
break Some(field.bytes().await?);
|
||||||
}
|
}
|
||||||
|
|
||||||
Some(_) => {}
|
Some(_) => {}
|
||||||
@@ -46,22 +36,29 @@ pub async fn monzo_batched_csv(
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
let Some(csv) = csv else {
|
Ok(csv)
|
||||||
return Err(AppError::BadRequest(anyhow!("No CSV file provided")));
|
}
|
||||||
};
|
|
||||||
|
pub async fn monzo_batched_csv(
|
||||||
|
Extension(db): Extension<DatabaseConnection>,
|
||||||
|
multipart: Multipart,
|
||||||
|
) -> Result<&'static str, AppError> {
|
||||||
|
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)
|
||||||
|
.await
|
||||||
|
.map_err(|e| AppError::BadRequest(anyhow!(e)))
|
||||||
|
.and_then(|csv| csv.ok_or(AppError::BadRequest(anyhow!(CSV_MISSING_ERR_MSG))))?;
|
||||||
|
|
||||||
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);
|
||||||
let data = csv.records();
|
let data = csv.records();
|
||||||
|
let data = data
|
||||||
db::insert(
|
.filter_map(|f| f.ok())
|
||||||
&db,
|
|
||||||
data.filter_map(|f| f.ok())
|
|
||||||
.map(from_csv_row)
|
.map(from_csv_row)
|
||||||
.collect::<Result<_, _>>()?,
|
.collect::<Result<_, _>>()?;
|
||||||
)
|
|
||||||
.await
|
db::insert(&db, data).await?;
|
||||||
.unwrap();
|
|
||||||
|
|
||||||
Ok("Ok")
|
Ok("Ok")
|
||||||
}
|
}
|
||||||
|
|||||||
+95
-18
@@ -2,27 +2,70 @@ mod error;
|
|||||||
mod ingestion;
|
mod ingestion;
|
||||||
|
|
||||||
use crate::error::AppError;
|
use crate::error::AppError;
|
||||||
use crate::ingestion::routes::{monzo_batched_csv, monzo_batched_json, monzo_updated};
|
use crate::ingestion::db;
|
||||||
|
use crate::ingestion::ingestion_logic::from_csv_row;
|
||||||
|
use crate::ingestion::routes::{monzo_batched_csv, monzo_batched_json};
|
||||||
use axum::routing::{get, post};
|
use axum::routing::{get, post};
|
||||||
use axum::{Extension, Router};
|
use axum::{Extension, Router};
|
||||||
use clap::Parser;
|
use clap::{Parser, Subcommand};
|
||||||
use migration::{Migrator, MigratorTrait};
|
use migration::{Migrator, MigratorTrait};
|
||||||
use sea_orm::{ConnectionTrait, DatabaseConnection};
|
use sea_orm::{ConnectionTrait, DatabaseConnection};
|
||||||
|
use std::fs::File;
|
||||||
use std::net::SocketAddr;
|
use std::net::SocketAddr;
|
||||||
|
use std::path::PathBuf;
|
||||||
|
use tower_http::trace::TraceLayer;
|
||||||
|
use tracing::log::LevelFilter;
|
||||||
|
|
||||||
#[derive(Debug, clap::Parser)]
|
use opentelemetry::trace::TracerProvider as _;
|
||||||
struct Config {
|
use opentelemetry_sdk::trace::TracerProvider;
|
||||||
/// If we should perform migration on startup.
|
use opentelemetry_otlp;
|
||||||
|
use tracing::{error, span};
|
||||||
|
use tracing_subscriber::layer::SubscriberExt;
|
||||||
|
use tracing_subscriber::Registry;
|
||||||
|
use tracin_support::init_tracing_subscriber;
|
||||||
|
|
||||||
|
mod tracin_support;
|
||||||
|
|
||||||
|
#[derive(Debug, Subcommand)]
|
||||||
|
enum Commands {
|
||||||
|
/// Manually run database migrations.
|
||||||
|
Migrate {
|
||||||
|
/// Number of migration steps to perform. If not provided, all migrations will be run.
|
||||||
|
#[arg(long)]
|
||||||
|
steps: Option<u32>,
|
||||||
|
|
||||||
|
/// If we should perform migration down.
|
||||||
|
#[arg(long)]
|
||||||
|
down: bool,
|
||||||
|
},
|
||||||
|
|
||||||
|
/// Start web app for to the google-sheets app script
|
||||||
|
Serve {
|
||||||
|
/// If we should perform migration at startup.
|
||||||
#[clap(short, long, env, default_value_t = true)]
|
#[clap(short, long, env, default_value_t = true)]
|
||||||
migrate: bool,
|
migrate: bool,
|
||||||
|
|
||||||
/// The server address to bind to.
|
/// The server address to bind to.
|
||||||
#[clap(short, long, env, default_value = "0.0.0.0:3000")]
|
#[clap(short, long, env, default_value = "0.0.0.0:3000")]
|
||||||
addr: SocketAddr,
|
addr: SocketAddr,
|
||||||
|
},
|
||||||
|
|
||||||
|
/// Ingest a google-sheets CSV export into the database.
|
||||||
|
Csv {
|
||||||
|
/// The path of the CSV file to ingest.
|
||||||
|
csv_file: PathBuf,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Ingest and manage monzo transactions from the Monzo Plus auto-export feature
|
||||||
|
#[derive(Debug, clap::Parser)]
|
||||||
|
struct Cli {
|
||||||
/// URL to PostgreSQL database.
|
/// URL to PostgreSQL database.
|
||||||
#[clap(short, long = "db", env)]
|
#[clap(short, long = "db", env)]
|
||||||
database_url: String,
|
database_url: String,
|
||||||
|
|
||||||
|
#[command(subcommand)]
|
||||||
|
command: Commands,
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn health_check(
|
async fn health_check(
|
||||||
@@ -35,24 +78,58 @@ async fn health_check(
|
|||||||
|
|
||||||
#[tokio::main]
|
#[tokio::main]
|
||||||
async fn main() -> anyhow::Result<()> {
|
async fn main() -> anyhow::Result<()> {
|
||||||
let config: Config = Config::parse();
|
let _guard = init_tracing_subscriber();
|
||||||
let connection = sea_orm::Database::connect(&config.database_url).await?;
|
|
||||||
|
|
||||||
if config.migrate {
|
let cli: Cli = Cli::parse();
|
||||||
|
let connection = sea_orm::ConnectOptions::new(&cli.database_url)
|
||||||
|
.sqlx_logging_level(LevelFilter::Debug)
|
||||||
|
.to_owned();
|
||||||
|
|
||||||
|
let connection = sea_orm::Database::connect(connection).await?;
|
||||||
|
|
||||||
|
match cli.command {
|
||||||
|
Commands::Migrate { steps, down } => {
|
||||||
|
if down {
|
||||||
|
Migrator::down(&connection, steps).await?;
|
||||||
|
} else {
|
||||||
|
Migrator::up(&connection, steps).await?
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Commands::Serve { migrate, addr } => {
|
||||||
|
if migrate {
|
||||||
Migrator::up(&connection, None).await?;
|
Migrator::up(&connection, None).await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
tracing_subscriber::fmt::init();
|
serve_web(addr, connection).await?;
|
||||||
let app = Router::new()
|
}
|
||||||
.route("/health", get(health_check))
|
|
||||||
.route("/monzo-updated", post(monzo_updated))
|
|
||||||
.route("/monzo-batch-export", post(monzo_batched_json))
|
|
||||||
.route("/monzo-csv-ingestion", post(monzo_batched_csv))
|
|
||||||
.layer(Extension(connection.clone()));
|
|
||||||
|
|
||||||
tracing::debug!("listening on {}", &config.addr);
|
Commands::Csv { csv_file } => {
|
||||||
let listener = tokio::net::TcpListener::bind(&config.addr).await.unwrap();
|
let mut csv = csv::Reader::from_reader(File::open(csv_file)?);
|
||||||
axum::serve(listener, app).await.unwrap();
|
let data = csv.records();
|
||||||
|
let data = data
|
||||||
|
.filter_map(|f| f.ok())
|
||||||
|
.map(from_csv_row)
|
||||||
|
.collect::<Result<_, _>>()?;
|
||||||
|
|
||||||
|
db::insert(&connection, data).await?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn serve_web(address: SocketAddr, connection: DatabaseConnection) -> anyhow::Result<()> {
|
||||||
|
let app = Router::new()
|
||||||
|
.route("/health", get(health_check))
|
||||||
|
.route("/monzo-batch-export", post(monzo_batched_json))
|
||||||
|
.route("/monzo-csv-ingestion", post(monzo_batched_csv))
|
||||||
|
.layer(Extension(connection.clone()))
|
||||||
|
.layer(TraceLayer::new_for_http());
|
||||||
|
|
||||||
|
tracing::info!("listening on {}", &address);
|
||||||
|
let listener = tokio::net::TcpListener::bind(&address).await?;
|
||||||
|
axum::serve(listener, app).await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,108 @@
|
|||||||
|
use opentelemetry::{global, trace::TracerProvider, Key, KeyValue};
|
||||||
|
use opentelemetry_sdk::{
|
||||||
|
metrics::{
|
||||||
|
reader::{DefaultAggregationSelector, DefaultTemporalitySelector},
|
||||||
|
Aggregation, Instrument, MeterProviderBuilder, PeriodicReader, SdkMeterProvider, Stream,
|
||||||
|
},
|
||||||
|
runtime,
|
||||||
|
trace::{BatchConfig, RandomIdGenerator, Sampler, Tracer},
|
||||||
|
Resource,
|
||||||
|
};
|
||||||
|
use opentelemetry_semantic_conventions::{
|
||||||
|
resource::{DEPLOYMENT_ENVIRONMENT, SERVICE_NAME, SERVICE_VERSION},
|
||||||
|
SCHEMA_URL,
|
||||||
|
};
|
||||||
|
use tracing::Level;
|
||||||
|
use tracing_opentelemetry::{MetricsLayer, OpenTelemetryLayer};
|
||||||
|
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
||||||
|
|
||||||
|
// Create a Resource that captures information about the entity for which telemetry is recorded.
|
||||||
|
fn resource() -> Resource {
|
||||||
|
Resource::from_schema_url(
|
||||||
|
[
|
||||||
|
KeyValue::new(SERVICE_NAME, env!("CARGO_PKG_NAME")),
|
||||||
|
KeyValue::new(SERVICE_VERSION, env!("CARGO_PKG_VERSION")),
|
||||||
|
KeyValue::new(DEPLOYMENT_ENVIRONMENT, "develop"),
|
||||||
|
],
|
||||||
|
SCHEMA_URL,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Construct MeterProvider for MetricsLayer
|
||||||
|
fn init_meter_provider() -> SdkMeterProvider {
|
||||||
|
let exporter = opentelemetry_otlp::new_exporter()
|
||||||
|
.http()
|
||||||
|
.with_http_client(reqwest::Client::new())
|
||||||
|
.build_metrics_exporter(
|
||||||
|
Box::new(DefaultAggregationSelector::new()),
|
||||||
|
Box::new(DefaultTemporalitySelector::new()),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let reader = PeriodicReader::builder(exporter, runtime::Tokio)
|
||||||
|
.with_interval(std::time::Duration::from_secs(30))
|
||||||
|
.build();
|
||||||
|
|
||||||
|
let meter_provider = MeterProviderBuilder::default()
|
||||||
|
.with_resource(resource())
|
||||||
|
.with_reader(reader)
|
||||||
|
.build();
|
||||||
|
|
||||||
|
global::set_meter_provider(meter_provider.clone());
|
||||||
|
|
||||||
|
meter_provider
|
||||||
|
}
|
||||||
|
|
||||||
|
// Construct Tracer for OpenTelemetryLayer
|
||||||
|
fn init_tracer() -> Tracer {
|
||||||
|
let provider = opentelemetry_otlp::new_pipeline()
|
||||||
|
.tracing()
|
||||||
|
.with_trace_config(
|
||||||
|
opentelemetry_sdk::trace::Config::default()
|
||||||
|
// Customize sampling strategy
|
||||||
|
.with_sampler(Sampler::ParentBased(Box::new(Sampler::TraceIdRatioBased(
|
||||||
|
1.0,
|
||||||
|
))))
|
||||||
|
// If export trace to AWS X-Ray, you can use XrayIdGenerator
|
||||||
|
.with_id_generator(RandomIdGenerator::default())
|
||||||
|
.with_resource(resource()),
|
||||||
|
)
|
||||||
|
.with_batch_config(BatchConfig::default())
|
||||||
|
.with_exporter(opentelemetry_otlp::new_exporter().tonic())
|
||||||
|
.install_batch(runtime::Tokio)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
global::set_tracer_provider(provider.clone());
|
||||||
|
provider.tracer("tracing-otel-subscriber")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Initialize tracing-subscriber and return OtelGuard for opentelemetry-related termination processing
|
||||||
|
pub(crate) fn init_tracing_subscriber() -> OtelGuard {
|
||||||
|
let meter_provider = init_meter_provider();
|
||||||
|
let tracer = init_tracer();
|
||||||
|
|
||||||
|
tracing_subscriber::registry()
|
||||||
|
.with(tracing_subscriber::filter::LevelFilter::from_level(
|
||||||
|
Level::DEBUG,
|
||||||
|
))
|
||||||
|
.with(tracing_subscriber::fmt::layer())
|
||||||
|
.with(MetricsLayer::new(meter_provider.clone()))
|
||||||
|
.with(OpenTelemetryLayer::new(tracer))
|
||||||
|
.init();
|
||||||
|
|
||||||
|
OtelGuard { meter_provider }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) struct OtelGuard {
|
||||||
|
meter_provider: SdkMeterProvider,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for OtelGuard {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
if let Err(err) = self.meter_provider.shutdown() {
|
||||||
|
eprintln!("{err:?}");
|
||||||
|
}
|
||||||
|
|
||||||
|
global::shutdown_tracer_provider();
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user