diff --git a/.github/workflows/crates-postgres.yaml b/.github/workflows/crates-postgres.yaml index f9ee7c543..3960a608f 100644 --- a/.github/workflows/crates-postgres.yaml +++ b/.github/workflows/crates-postgres.yaml @@ -78,7 +78,9 @@ jobs: run: make -C ./crates/pg test-lint - name: crates/pg unit tests run: make -C ./crates/pg test-unit + - name: start postgres + run: make -C migrations start-nobuild - name: crates/pg db_tests - run: make -C ./crates/pg test-db START_TARGET=start-nobuild + run: make -C ./crates/pg test-db - name: crates/pg clean up run: make -C ./crates/pg stop diff --git a/.github/workflows/prod-services-deploy.yaml b/.github/workflows/prod-services-deploy.yaml index d816155f3..adb7ef97c 100644 --- a/.github/workflows/prod-services-deploy.yaml +++ b/.github/workflows/prod-services-deploy.yaml @@ -27,9 +27,16 @@ jobs: - uses: actions/checkout@v4 - name: mask values run: echo "::add-mask::${{ secrets.AWS_ACCOUNT_ID }}" - - name: deploy services to prod cloud environment + - name: deploy lambda services to prod cloud environment run: | - for app in $(bash scripts/list-deployments.sh | grep -v client | xargs basename -a); do + for app_dir in $(bash scripts/list-deployments.sh | grep -v client); do + app=$(basename $app_dir) + CONF_PATH=$(bash scripts/dir-to-conf-path.sh "$app_dir") + DEPLOY_TARGET=$(yq "$CONF_PATH.deploy_target" project.yaml) + if [[ "$DEPLOY_TARGET" != "lambda" ]]; then + echo "*** skipping $app (deploy_target: $DEPLOY_TARGET)" + continue + fi # tag newly pushed prod images with master merge commit sha bash scripts/tag-merge-commit.sh --app-name $app --env prod --env-id ${{ secrets.PROD_ENV_ID }}; # deploy newly tagged images to prod cloud diff --git a/.github/workflows/services.yaml b/.github/workflows/services.yaml index e123b0f71..cc8a810f8 100644 --- a/.github/workflows/services.yaml +++ b/.github/workflows/services.yaml @@ -139,8 +139,10 @@ jobs: - uses: Swatinem/rust-cache@v2 - name: stop runner postgres run: sudo systemctl stop postgresql + - name: start postgres + run: make -C migrations start-nobuild - name: test database - run: make -C crates/pg test-db START_TARGET=start-nobuild + run: make -C crates/pg test-db compile: name: compile ${{ matrix.service }} diff --git a/Cargo.lock b/Cargo.lock index c9ced2f21..49dfef670 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1498,6 +1498,7 @@ dependencies = [ "tokio-postgres", "tracing", "tracing-subscriber", + "types", ] [[package]] diff --git a/crates/pg/makefile b/crates/pg/makefile index f63dcb4a0..2cbe85370 100644 --- a/crates/pg/makefile +++ b/crates/pg/makefile @@ -3,8 +3,6 @@ RELATIVE_PROJECT_ROOT_PATH=$(shell REL_PATH="."; while [ $$(ls "$$REL_PATH" | gr APP_NAME=$(shell basename $(CURDIR)) PROJECT_CONF_FILE_NAME=project.yaml TEST_NAME?=model::integration_tests:: -MIGRATIONS_DIR=$(RELATIVE_PROJECT_ROOT_PATH)/migrations -START_TARGET?=start test-unit: cargo test @@ -25,7 +23,6 @@ watch-db: test-db: @$(MAKE) get-secrets - @$(MAKE) -C $(MIGRATIONS_DIR) $(START_TARGET) @cd $(RELATIVE_PROJECT_ROOT_PATH); bash scripts/test-reset.sh --set cargo test $(TEST_NAME) --features db_tests -- --test-threads=1 diff --git a/crates/types/src/transaction.rs b/crates/types/src/transaction.rs index d24396f72..51b09824d 100644 --- a/crates/types/src/transaction.rs +++ b/crates/types/src/transaction.rs @@ -84,6 +84,26 @@ impl Transaction { } impl Transaction { + pub fn from_rule_instance( + rule_instance_id: String, + author: Option, + author_role: Option, + ) -> Self { + Transaction { + id: None, + rule_instance_id: Some(rule_instance_id), + author, + author_device_id: None, + author_device_latlng: None, + author_role, + equilibrium_time: None, + event_time: None, + debitor_first: None, + sum_value: "0".to_string(), + transaction_items: TransactionItems(vec![]), + } + } + pub fn new( author: String, equilibrium_time: Option, @@ -459,6 +479,46 @@ pub mod tests { } } + #[test] + fn it_creates_transaction_from_rule_instance_with_author() { + let got = Transaction::from_rule_instance( + "42".to_string(), + Some("StateOfCalifornia".to_string()), + Some(AccountRole::Creditor), + ); + assert_eq!(got.rule_instance_id, Some("42".to_string())); + assert_eq!(got.author, Some("StateOfCalifornia".to_string())); + assert_eq!(got.author_role, Some(AccountRole::Creditor)); + assert_eq!(got.id, None); + assert_eq!(got.author_device_id, None); + assert_eq!(got.author_device_latlng, None); + assert_eq!(got.equilibrium_time, None); + assert_eq!(got.event_time, None); + assert_eq!(got.debitor_first, None); + assert_eq!(got.sum_value, "0"); + assert!(got.transaction_items.0.is_empty()); + } + + #[test] + fn it_creates_transaction_from_rule_instance_without_author() { + let got = Transaction::from_rule_instance("7".to_string(), None, None); + assert_eq!(got.rule_instance_id, Some("7".to_string())); + assert_eq!(got.author, None); + assert_eq!(got.author_role, None); + } + + #[test] + fn it_serializes_transaction_from_rule_instance() { + let transaction = Transaction::from_rule_instance( + "42".to_string(), + Some("StateOfCalifornia".to_string()), + Some(AccountRole::Creditor), + ); + let json = serde_json::to_string(&transaction).unwrap(); + let deserialized: Transaction = serde_json::from_str(&json).unwrap(); + assert_eq!(transaction, deserialized); + } + // cadet todo: test remaining branches of get_author_role #[test] fn it_returns_author_role_found_in_transaction_items() { diff --git a/docker/go-migrate.Dockerfile b/docker/go-migrate.Dockerfile index 27557bd58..29228aa87 100644 --- a/docker/go-migrate.Dockerfile +++ b/docker/go-migrate.Dockerfile @@ -18,6 +18,7 @@ COPY migrations/schema task/migrations/schema COPY migrations/seed task/migrations/seed COPY migrations/testseed task/migrations/testseed COPY migrations/testseedthresh task/migrations/testseedthresh +COPY migrations/testseedcron task/migrations/testseedcron FROM public.ecr.aws/lambda/provided:al2023 diff --git a/migrations/go-migrate/migrate.sh b/migrations/go-migrate/migrate.sh index 80e4c9dc4..44f60e363 100644 --- a/migrations/go-migrate/migrate.sh +++ b/migrations/go-migrate/migrate.sh @@ -33,7 +33,7 @@ if [[ -z "$SUBDIRS" ]]; then exit 1 fi -ALLOWED_SUBDIRS=(schema seed testseed testseedthresh) +ALLOWED_SUBDIRS=(schema seed testseed testseedthresh testseedcron) IFS=',' read -ra SUBDIR_ARGS <<< "$SUBDIRS" for s in "${SUBDIR_ARGS[@]}"; do if [[ ! " ${ALLOWED_SUBDIRS[*]} " =~ " ${s} " ]]; then diff --git a/migrations/makefile b/migrations/makefile index 0bd15b25e..e3240e4e8 100644 --- a/migrations/makefile +++ b/migrations/makefile @@ -116,6 +116,34 @@ thresh-up: thresh-down: @set -a && source $(ENV_FILE) && set +a && SQL_TYPE=postgresql bash $(RELATIVE_PROJECT_ROOT_PATH)/migrations/go-migrate/migrate.sh --dir $(RELATIVE_PROJECT_ROOT_PATH)/migrations --subdirs testseedthresh --cmd down +cron-up: + @set -a && source $(ENV_FILE) && set +a && SQL_TYPE=postgresql bash $(RELATIVE_PROJECT_ROOT_PATH)/migrations/go-migrate/migrate.sh --dir $(RELATIVE_PROJECT_ROOT_PATH)/migrations --subdirs testseedcron --cmd up + @$(MAKE) warm-cache + +cron-down: + @set -a && source $(ENV_FILE) && set +a && SQL_TYPE=postgresql bash $(RELATIVE_PROJECT_ROOT_PATH)/migrations/go-migrate/migrate.sh --dir $(RELATIVE_PROJECT_ROOT_PATH)/migrations --subdirs testseedcron --cmd down + +# poll for cron transactions (waits 15 seconds for pg_cron to fire) +wait-cron: + @for i in $$(seq 1 15); do printf '#'; sleep 1; done; echo + +# poll redis accumulator until threshold fires (value resets) +wait-thresh: + @key=$$(docker exec mxf-redis-1 redis-cli -a test --no-auth-warning KEYS '*accumulator' 2>/dev/null | head -1); \ + while [ -z "$$key" ]; do \ + sleep 2; \ + key=$$(docker exec mxf-redis-1 redis-cli -a test --no-auth-warning KEYS '*accumulator' 2>/dev/null | head -1); \ + done; \ + prev=0; \ + while true; do \ + val=$$(docker exec mxf-redis-1 redis-cli -a test --no-auth-warning GET "$$key" 2>/dev/null | tr -d '"'); \ + if [ -z "$$val" ] || [ "$$val" = "(nil)" ]; then sleep 2; continue; fi; \ + echo "accumulator: $$val"; \ + if [ "$$(echo "$$val < $$prev" | bc -l)" = "1" ]; then echo "threshold fired"; break; fi; \ + prev=$$val; \ + sleep 2; \ + done + ###################### env ###################### env: diff --git a/migrations/schema/000010_notify.down.sql b/migrations/schema/000010_notify.down.sql new file mode 100644 index 000000000..17dbc68b4 --- /dev/null +++ b/migrations/schema/000010_notify.down.sql @@ -0,0 +1 @@ +DROP FUNCTION IF EXISTS notify_event; diff --git a/migrations/schema/000010_notify.up.sql b/migrations/schema/000010_notify.up.sql new file mode 100644 index 000000000..6575fd28e --- /dev/null +++ b/migrations/schema/000010_notify.up.sql @@ -0,0 +1,9 @@ +CREATE OR REPLACE FUNCTION notify_event(event_name TEXT, event_id TEXT) +RETURNS void AS $$ +BEGIN + PERFORM pg_notify('event', json_build_object( + 'event', event_name, + 'id', event_id + )::text); +END; +$$ LANGUAGE plpgsql; diff --git a/migrations/schema/000010_equilibrium.down.sql b/migrations/schema/000011_equilibrium.down.sql similarity index 100% rename from migrations/schema/000010_equilibrium.down.sql rename to migrations/schema/000011_equilibrium.down.sql diff --git a/migrations/schema/000010_equilibrium.up.sql b/migrations/schema/000011_equilibrium.up.sql similarity index 96% rename from migrations/schema/000010_equilibrium.up.sql rename to migrations/schema/000011_equilibrium.up.sql index db3fc64f2..9199ef0db 100644 --- a/migrations/schema/000010_equilibrium.up.sql +++ b/migrations/schema/000011_equilibrium.up.sql @@ -33,7 +33,7 @@ CREATE INDEX redis_name_key_idx ON redis_name(key); CREATE OR REPLACE FUNCTION notify_equilibrium() RETURNS TRIGGER AS $$ BEGIN - PERFORM pg_notify('equilibrium', NEW.id::text); + PERFORM notify_event('equilibrium', NEW.id::text); RETURN NEW; END; $$ LANGUAGE plpgsql; diff --git a/migrations/schema/000012_cron.down.sql b/migrations/schema/000012_cron.down.sql new file mode 100644 index 000000000..63a5cf6ac --- /dev/null +++ b/migrations/schema/000012_cron.down.sql @@ -0,0 +1,2 @@ +DROP FUNCTION IF EXISTS notify_cron; +DROP EXTENSION IF EXISTS pg_cron; \ No newline at end of file diff --git a/migrations/schema/000012_cron.up.sql b/migrations/schema/000012_cron.up.sql new file mode 100644 index 000000000..fac81b1e4 --- /dev/null +++ b/migrations/schema/000012_cron.up.sql @@ -0,0 +1,8 @@ +CREATE EXTENSION IF NOT EXISTS pg_cron; + +CREATE OR REPLACE FUNCTION notify_cron(instance_id INT) +RETURNS void AS $$ +BEGIN + PERFORM notify_event('cron', instance_id::text); +END; +$$ LANGUAGE plpgsql; \ No newline at end of file diff --git a/migrations/testseedcron/000001_interest.down.sql b/migrations/testseedcron/000001_interest.down.sql new file mode 100644 index 000000000..339104b1b --- /dev/null +++ b/migrations/testseedcron/000001_interest.down.sql @@ -0,0 +1,31 @@ +SELECT cron.unschedule('EnergyCoInterest'); + +DELETE FROM approval WHERE transaction_item_id IN ( + SELECT id FROM transaction_item WHERE rule_instance_id IN ( + SELECT id FROM transaction_item_rule_instance WHERE transaction_rule_instance_id = ( + SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoInterest' + ) + ) +); + +DELETE FROM transaction_item WHERE rule_instance_id IN ( + SELECT id FROM transaction_item_rule_instance WHERE transaction_rule_instance_id = ( + SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoInterest' + ) +); + +DELETE FROM transaction WHERE id NOT IN (SELECT DISTINCT transaction_id FROM transaction_item); + +DELETE FROM approval WHERE rule_instance_id IN ( + SELECT id FROM approval_rule_instance + WHERE rule_instance_name IN ('ApproveInterestEnergyCo', 'ApproveInterestJoeCarter') +); + +DELETE FROM approval_rule_instance WHERE rule_instance_name IN ('ApproveInterestEnergyCo', 'ApproveInterestJoeCarter'); + +DELETE FROM transaction_item_rule_instance +WHERE transaction_rule_instance_id = ( + SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoInterest' +); + +DELETE FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoInterest'; diff --git a/migrations/testseedcron/000001_interest.up.sql b/migrations/testseedcron/000001_interest.up.sql new file mode 100644 index 000000000..28c4bfda0 --- /dev/null +++ b/migrations/testseedcron/000001_interest.up.sql @@ -0,0 +1,44 @@ +-- cron debt: EnergyCo pays interest to JoeCarter every 10 seconds +-- +-- same debitor, creditor and amount as cron equity (000002_dividend) +-- modigliani-miller: "debt" and "equity" are just rule instance names +-- when the underlying value transfer is identical + +insert into transaction_rule_instance + (rule_name, rule_instance_name, author, author_role, cron) +values + ('createTransaction', 'EnergyCoInterest', 'EnergyCo', 'debitor', '10 seconds'); + +-- transaction_item_rule_instance: linked line item +insert into transaction_item_rule_instance + (rule_name, rule_instance_name, account_role, account_name, + item_id, price, quantity, variable_values, transaction_rule_instance_id) +values + ('addTransactionItem', 'EnergyCoInterest', 'debitor', 'EnergyCo', + '1,000 x 1% monthly interest', 10.000, 1, '{ "EnergyCo", "JoeCarter" }', + (select id from transaction_rule_instance where rule_instance_name = 'EnergyCoInterest')); + +-- approval_rule_instance: AaronHill approves EnergyCo debits (EnergyCo owner) +insert into approval_rule_instance + (rule_name, rule_instance_name, account_role, account_name, variable_values) +values + ('approveItemBetweenAccounts', 'ApproveInterestEnergyCo', 'debitor', 'AaronHill', + '{ "EnergyCo", "JoeCarter", "1,000 x 1% monthly interest", "debitor", "AaronHill" }'); + +-- approval_rule_instance: JoeCarter approves credits (self-owner) +insert into approval_rule_instance + (rule_name, rule_instance_name, account_role, account_name, variable_values) +values + ('approveItemBetweenAccounts', 'ApproveInterestJoeCarter', 'creditor', 'JoeCarter', + '{ "EnergyCo", "JoeCarter", "1,000 x 1% monthly interest", "creditor", "JoeCarter" }'); + +-- schedule pg_cron job and store jobid +UPDATE transaction_rule_instance +SET cron_job_id = cron.schedule( + 'EnergyCoInterest', + '10 seconds', + $$SELECT notify_cron( + (SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoInterest') + )$$ +) +WHERE rule_instance_name = 'EnergyCoInterest'; diff --git a/migrations/testseedcron/000002_dividend.down.sql b/migrations/testseedcron/000002_dividend.down.sql new file mode 100644 index 000000000..df858f0cd --- /dev/null +++ b/migrations/testseedcron/000002_dividend.down.sql @@ -0,0 +1,31 @@ +SELECT cron.unschedule('EnergyCoDividend'); + +DELETE FROM approval WHERE transaction_item_id IN ( + SELECT id FROM transaction_item WHERE rule_instance_id IN ( + SELECT id FROM transaction_item_rule_instance WHERE transaction_rule_instance_id = ( + SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoDividend' + ) + ) +); + +DELETE FROM transaction_item WHERE rule_instance_id IN ( + SELECT id FROM transaction_item_rule_instance WHERE transaction_rule_instance_id = ( + SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoDividend' + ) +); + +DELETE FROM transaction WHERE id NOT IN (SELECT DISTINCT transaction_id FROM transaction_item); + +DELETE FROM approval WHERE rule_instance_id IN ( + SELECT id FROM approval_rule_instance + WHERE rule_instance_name IN ('ApproveEnergyDividendEnergyCo', 'ApproveEnergyDividendJoeCarter') +); + +DELETE FROM approval_rule_instance WHERE rule_instance_name IN ('ApproveEnergyDividendEnergyCo', 'ApproveEnergyDividendJoeCarter'); + +DELETE FROM transaction_item_rule_instance +WHERE transaction_rule_instance_id = ( + SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoDividend' +); + +DELETE FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoDividend'; diff --git a/migrations/testseedcron/000002_dividend.up.sql b/migrations/testseedcron/000002_dividend.up.sql new file mode 100644 index 000000000..6233688f6 --- /dev/null +++ b/migrations/testseedcron/000002_dividend.up.sql @@ -0,0 +1,44 @@ +-- cron equity: EnergyCo pays dividend to JoeCarter every 10 seconds +-- +-- same debitor, creditor and amount as cron debt (000001_interest) +-- modigliani-miller: "debt" and "equity" are just rule instance names +-- when the underlying value transfer is identical + +insert into transaction_rule_instance + (rule_name, rule_instance_name, author, author_role, cron) +values + ('createTransaction', 'EnergyCoDividend', 'EnergyCo', 'debitor', '10 seconds'); + +-- transaction_item_rule_instance: linked line item +insert into transaction_item_rule_instance + (rule_name, rule_instance_name, account_role, account_name, + item_id, price, quantity, variable_values, transaction_rule_instance_id) +values + ('addTransactionItem', 'EnergyCoDividend', 'debitor', 'EnergyCo', + 'monthly 10.000 dividend', 10.000, 1, '{ "EnergyCo", "JoeCarter" }', + (select id from transaction_rule_instance where rule_instance_name = 'EnergyCoDividend')); + +-- approval_rule_instance: AaronHill approves EnergyCo debits (EnergyCo owner) +insert into approval_rule_instance + (rule_name, rule_instance_name, account_role, account_name, variable_values) +values + ('approveItemBetweenAccounts', 'ApproveEnergyDividendEnergyCo', 'debitor', 'AaronHill', + '{ "EnergyCo", "JoeCarter", "monthly 10.000 dividend", "debitor", "AaronHill" }'); + +-- approval_rule_instance: JoeCarter approves credits (self-owner) +insert into approval_rule_instance + (rule_name, rule_instance_name, account_role, account_name, variable_values) +values + ('approveItemBetweenAccounts', 'ApproveEnergyDividendJoeCarter', 'creditor', 'JoeCarter', + '{ "EnergyCo", "JoeCarter", "monthly 10.000 dividend", "creditor", "JoeCarter" }'); + +-- schedule pg_cron job and store jobid +UPDATE transaction_rule_instance +SET cron_job_id = cron.schedule( + 'EnergyCoDividend', + '10 seconds', + $$SELECT notify_cron( + (SELECT id FROM transaction_rule_instance WHERE rule_instance_name = 'EnergyCoDividend') + )$$ +) +WHERE rule_instance_name = 'EnergyCoDividend'; diff --git a/migrations/testseedthresh/000001_dividend.down.sql b/migrations/testseedthresh/000001_dividend.down.sql index cd148924d..ccfd80ae0 100644 --- a/migrations/testseedthresh/000001_dividend.down.sql +++ b/migrations/testseedthresh/000001_dividend.down.sql @@ -1,6 +1,6 @@ delete from approval where transaction_item_id in (select ti.id from transaction_item ti where ti.rule_instance_id in (select id from transaction_item_rule_instance where rule_instance_name = 'GroceryCoProfitDividend')); delete from transaction_item where rule_instance_id in (select id from transaction_item_rule_instance where rule_instance_name = 'GroceryCoProfitDividend'); delete from transaction where id not in (select distinct transaction_id from transaction_item); -delete from approval_rule_instance where rule_instance_name in ('ApproveDividendGroceryCo', 'ApproveDividendTimeCo'); +delete from approval_rule_instance where rule_instance_name in ('ApproveGroceryDividendGroceryCo', 'ApproveGroceryDividendJoeCarter'); delete from transaction_item_rule_instance where rule_instance_name = 'GroceryCoProfitDividend'; delete from transaction_rule_instance where rule_instance_name = 'GroceryCoProfitDividend'; diff --git a/migrations/testseedthresh/000001_dividend.up.sql b/migrations/testseedthresh/000001_dividend.up.sql index 895dea2aa..be28f833d 100644 --- a/migrations/testseedthresh/000001_dividend.up.sql +++ b/migrations/testseedthresh/000001_dividend.up.sql @@ -17,19 +17,19 @@ values insert into approval_rule_instance (rule_name, rule_instance_name, account_role, account_name, variable_values) values - ('approveItemBetweenAccounts', 'ApproveDividendGroceryCo', 'debitor', 'IgorPetrov', + ('approveItemBetweenAccounts', 'ApproveGroceryDividendGroceryCo', 'debitor', 'IgorPetrov', '{ "GroceryCo", "JoeCarter", "1% dividend", "debitor", "IgorPetrov" }'); -- approval_rule_instance: MiriamLevy approves GroceryCo → JoeCarter debits insert into approval_rule_instance (rule_name, rule_instance_name, account_role, account_name, variable_values) values - ('approveItemBetweenAccounts', 'ApproveDividendGroceryCo', 'debitor', 'MiriamLevy', + ('approveItemBetweenAccounts', 'ApproveGroceryDividendGroceryCo', 'debitor', 'MiriamLevy', '{ "GroceryCo", "JoeCarter", "1% dividend", "debitor", "MiriamLevy" }'); -- approval_rule_instance: JoeCarter approves GroceryCo → JoeCarter credits insert into approval_rule_instance (rule_name, rule_instance_name, account_role, account_name, variable_values) values - ('approveItemBetweenAccounts', 'ApproveDividendJoeCarter', 'creditor', 'JoeCarter', + ('approveItemBetweenAccounts', 'ApproveGroceryDividendJoeCarter', 'creditor', 'JoeCarter', '{ "GroceryCo", "JoeCarter", "1% dividend", "creditor", "JoeCarter" }'); diff --git a/scripts/print-env-id.sh b/scripts/print-env-id.sh index 6431fd50c..e2b849e8e 100644 --- a/scripts/print-env-id.sh +++ b/scripts/print-env-id.sh @@ -11,7 +11,7 @@ ENV_FILE=$ENV_FILE_NAME # assumes project root if [[ $ENV == 'prod' ]]; then # use configured prod env id ENV_ID=$(yq '.env_var.set.PROD_ENV_ID.default' $PROJECT_CONF) -elif [[ -z $ENV_ID ]]; then # get ENV_ID from .env file if not set in env +elif [[ -z ${ENV_ID:-} ]]; then # get ENV_ID from .env file if not set in env if [[ -f $ENV_FILE ]] && grep -q "ENV_ID=" $ENV_FILE; then ENV_ID=$(grep "ENV_ID=" $ENV_FILE | cut -d'=' -f2) else diff --git a/scripts/test-reset.sh b/scripts/test-reset.sh index 68ac85d86..5528d8786 100644 --- a/scripts/test-reset.sh +++ b/scripts/test-reset.sh @@ -84,6 +84,7 @@ if [[ "$SET_MODE" == true ]]; then # query initial max IDs and initial balance from postgres IFS='|' read -r APPROVAL_MAX TRANSACTION_ITEM_MAX TRANSACTION_MAX \ ACCOUNT_PROFILE_MAX APPROVAL_RULE_INSTANCE_MAX ACCOUNT_OWNER_MAX \ + TRANSACTION_RULE_INSTANCE_MAX TRANSACTION_ITEM_RULE_INSTANCE_MAX \ INITIAL_BALANCE <<< "$(psql -d "$DBCONN" -t -A <<'SQL' SELECT (SELECT COALESCE(MAX(id), 0) FROM approval), @@ -92,6 +93,8 @@ SELECT (SELECT COALESCE(MAX(id), 0) FROM account_profile), (SELECT COALESCE(MAX(id), 0) FROM approval_rule_instance), (SELECT COALESCE(MAX(id), 0) FROM account_owner), + (SELECT COALESCE(MAX(id), 0) FROM transaction_rule_instance), + (SELECT COALESCE(MAX(id), 0) FROM transaction_item_rule_instance), (SELECT current_balance FROM account_balance LIMIT 1) SQL )" @@ -105,6 +108,8 @@ SQL "account_profile|$ACCOUNT_PROFILE_MAX" \ "approval_rule_instance|$APPROVAL_RULE_INSTANCE_MAX" \ "account_owner|$ACCOUNT_OWNER_MAX" \ + "transaction_rule_instance|$TRANSACTION_RULE_INSTANCE_MAX" \ + "transaction_item_rule_instance|$TRANSACTION_ITEM_RULE_INSTANCE_MAX" \ "initial_balance|$INITIAL_BALANCE"; do KEY="${pair%%|*}" VAL="${pair##*|}" @@ -121,6 +126,8 @@ SQL test_fixture:account_profile "$ACCOUNT_PROFILE_MAX" \ test_fixture:approval_rule_instance "$APPROVAL_RULE_INSTANCE_MAX" \ test_fixture:account_owner "$ACCOUNT_OWNER_MAX" \ + test_fixture:transaction_rule_instance "$TRANSACTION_RULE_INSTANCE_MAX" \ + test_fixture:transaction_item_rule_instance "$TRANSACTION_ITEM_RULE_INSTANCE_MAX" \ test_fixture:initial_balance "$INITIAL_BALANCE" \ > /dev/null fi @@ -137,6 +144,8 @@ if [[ -n "${ENV:-}" ]]; then ACCOUNT_PROFILE_MAX=$(aws dynamodb get-item --table-name "$DDB_TABLE" --key '{"pk":{"S":"test_fixture:account_profile"},"sk":{"S":"_"}}' --query 'Item.val.N' --region "$REGION" --output text) APPROVAL_RULE_INSTANCE_MAX=$(aws dynamodb get-item --table-name "$DDB_TABLE" --key '{"pk":{"S":"test_fixture:approval_rule_instance"},"sk":{"S":"_"}}' --query 'Item.val.N' --region "$REGION" --output text) ACCOUNT_OWNER_MAX=$(aws dynamodb get-item --table-name "$DDB_TABLE" --key '{"pk":{"S":"test_fixture:account_owner"},"sk":{"S":"_"}}' --query 'Item.val.N' --region "$REGION" --output text) + TRANSACTION_RULE_INSTANCE_MAX=$(aws dynamodb get-item --table-name "$DDB_TABLE" --key '{"pk":{"S":"test_fixture:transaction_rule_instance"},"sk":{"S":"_"}}' --query 'Item.val.N' --region "$REGION" --output text) + TRANSACTION_ITEM_RULE_INSTANCE_MAX=$(aws dynamodb get-item --table-name "$DDB_TABLE" --key '{"pk":{"S":"test_fixture:transaction_item_rule_instance"},"sk":{"S":"_"}}' --query 'Item.val.N' --region "$REGION" --output text) INITIAL_BALANCE=$(aws dynamodb get-item --table-name "$DDB_TABLE" --key '{"pk":{"S":"test_fixture:initial_balance"},"sk":{"S":"_"}}' --query 'Item.val.N' --region "$REGION" --output text) else APPROVAL_MAX=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning GET test_fixture:approval 2>/dev/null) @@ -145,12 +154,15 @@ else ACCOUNT_PROFILE_MAX=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning GET test_fixture:account_profile 2>/dev/null) APPROVAL_RULE_INSTANCE_MAX=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning GET test_fixture:approval_rule_instance 2>/dev/null) ACCOUNT_OWNER_MAX=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning GET test_fixture:account_owner 2>/dev/null) + TRANSACTION_RULE_INSTANCE_MAX=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning GET test_fixture:transaction_rule_instance 2>/dev/null) + TRANSACTION_ITEM_RULE_INSTANCE_MAX=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning GET test_fixture:transaction_item_rule_instance 2>/dev/null) INITIAL_BALANCE=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning GET test_fixture:initial_balance 2>/dev/null) fi # validate initial values exist for VAL in "$APPROVAL_MAX" "$TRANSACTION_ITEM_MAX" "$TRANSACTION_MAX" \ "$ACCOUNT_PROFILE_MAX" "$APPROVAL_RULE_INSTANCE_MAX" "$ACCOUNT_OWNER_MAX" \ + "$TRANSACTION_RULE_INSTANCE_MAX" "$TRANSACTION_ITEM_RULE_INSTANCE_MAX" \ "$INITIAL_BALANCE"; do if [[ -z "$VAL" || "$VAL" == "(nil)" || "$VAL" == "None" ]]; then echo "*** error: fixture not set. run: bash scripts/test-reset.sh --set" @@ -169,6 +181,8 @@ UPDATE account_balance SET current_balance = $INITIAL_BALANCE, current_transacti SET session_replication_role = 'origin'; DELETE FROM account_profile WHERE id > $ACCOUNT_PROFILE_MAX; DELETE FROM approval_rule_instance WHERE id > $APPROVAL_RULE_INSTANCE_MAX; +DELETE FROM transaction_item_rule_instance WHERE id > $TRANSACTION_ITEM_RULE_INSTANCE_MAX; +DELETE FROM transaction_rule_instance WHERE id > $TRANSACTION_RULE_INSTANCE_MAX; DELETE FROM account_owner WHERE owner_account = 'test_account'; DELETE FROM account WHERE name = 'test_account'; @@ -180,6 +194,8 @@ SELECT setval('transaction_id_seq', $TRANSACTION_MAX); SELECT setval('transaction_item_id_seq', $TRANSACTION_ITEM_MAX); SELECT setval('approval_id_seq', $APPROVAL_MAX); SELECT setval('approval_rule_instance_id_seq', $APPROVAL_RULE_INSTANCE_MAX); +SELECT setval('transaction_rule_instance_id_seq', GREATEST($TRANSACTION_RULE_INSTANCE_MAX, 1), $TRANSACTION_RULE_INSTANCE_MAX > 0); +SELECT setval('transaction_item_rule_instance_id_seq', GREATEST($TRANSACTION_ITEM_RULE_INSTANCE_MAX, 1), $TRANSACTION_ITEM_RULE_INSTANCE_MAX > 0); SELECT setval('account_profile_id_seq', $ACCOUNT_PROFILE_MAX); SELECT setval('account_owner_id_seq', $ACCOUNT_OWNER_MAX); SQL @@ -189,8 +205,8 @@ if [[ -n "${ENV:-}" ]]; then # scan for gdp and accumulator keys then batch delete ITEMS=$(aws dynamodb scan \ --table-name "$DDB_TABLE" \ - --filter-expression "contains(pk, :gdp) OR begins_with(pk, :acc)" \ - --expression-attribute-values '{":gdp":{"S":":gdp:"},":acc":{"S":"transaction_rule_instance:"}}' \ + --filter-expression "contains(pk, :gdp) OR begins_with(pk, :acc) OR begins_with(pk, :appr)" \ + --expression-attribute-values '{":gdp":{"S":":gdp:"},":acc":{"S":"transaction_rule_instance:"},":appr":{"S":"rules:approval:"}}' \ --projection-expression "pk, sk" \ --region "$REGION" \ --output json) @@ -203,6 +219,10 @@ if [[ -n "${ENV:-}" ]]; then --region "$REGION" > /dev/null fi else + approval_rule_keys=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning KEYS 'rules:approval:*' 2>/dev/null | tr '\n' ' ') + if [[ -n "$approval_rule_keys" ]]; then + docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning DEL $approval_rule_keys > /dev/null + fi gdp_keys=$(docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning KEYS '*:gdp:*' 2>/dev/null | tr '\n' ' ') if [[ -n "$gdp_keys" ]]; then docker exec mxf-redis-1 redis-cli -a "$REDIS_PASSWORD" --no-auth-warning DEL $gdp_keys > /dev/null diff --git a/services/event/Cargo.toml b/services/event/Cargo.toml index 1c5fdbfc4..21d02a08b 100644 --- a/services/event/Cargo.toml +++ b/services/event/Cargo.toml @@ -17,5 +17,6 @@ cache = { path = "../../crates/cache" } pubsub = { path = "../../crates/pubsub" } queue = { path = "../../crates/queue" } pg = { path = "../../crates/pg" } +types = { path = "../../crates/types" } shutdown = { path = "../../crates/shutdown" } chrono = "0.4" diff --git a/services/event/src/events/cron.rs b/services/event/src/events/cron.rs new file mode 100644 index 000000000..565b45b30 --- /dev/null +++ b/services/event/src/events/cron.rs @@ -0,0 +1,39 @@ +use pg::postgres::ConnectionPool; +use std::sync::Arc; +use types::account_role::AccountRole; +use types::transaction::Transaction; + +const RULE_INSTANCE_QUERY: &str = r#" +SELECT author, author_role FROM transaction_rule_instance WHERE id = $1 +"#; + +pub async fn handle_cron( + pool: &ConnectionPool, + queue: &Arc, + rule_instance_id: &str, +) -> Result<(), Box> { + let conn = pool.get_conn().await; + let ri_id: i32 = rule_instance_id.parse()?; + + let rows = conn.0.query(RULE_INSTANCE_QUERY, &[&ri_id]).await?; + let row = rows + .first() + .ok_or_else(|| format!("no rule instance found for id {}", rule_instance_id))?; + + let author: Option = row.get(0); + let author_role: Option = row.get(1); + + let transaction = Transaction::from_rule_instance( + rule_instance_id.to_string(), + author, + author_role.map(AccountRole::from), + ); + + queue.send(&serde_json::to_string(&transaction)?).await?; + tracing::info!( + "queued cron transaction for rule_instance_id {}", + rule_instance_id + ); + + Ok(()) +} diff --git a/services/event/src/events/mod.rs b/services/event/src/events/mod.rs index 1a4320196..d3bc4b8a8 100644 --- a/services/event/src/events/mod.rs +++ b/services/event/src/events/mod.rs @@ -1,5 +1,7 @@ +mod cron; mod gdp; mod threshold_profit; +pub use cron::handle_cron; pub use gdp::handle_gdp; pub use threshold_profit::handle_threshold_profit; diff --git a/services/event/src/events/threshold_profit.rs b/services/event/src/events/threshold_profit.rs index ff41ea7f3..53d4a8098 100644 --- a/services/event/src/events/threshold_profit.rs +++ b/services/event/src/events/threshold_profit.rs @@ -3,6 +3,8 @@ use pg::postgres::ConnectionPool; use rust_decimal::Decimal; use std::collections::HashSet; use std::sync::Arc; +use types::account_role::AccountRole; +use types::transaction::Transaction; // query accounts involved in a transaction (creditors and debitors) const ACCOUNTS_QUERY: &str = r#" @@ -86,15 +88,13 @@ pub async fn handle_threshold_profit( let author: Option = tri.get(1); let author_role: Option = tri.get(2); - let payload = serde_json::json!({ - "rule_instance_id": rule_instance_id, - "author": author, - "author_role": author_role, - "sum_value": "0", - "transaction_items": [], - }); + let transaction = Transaction::from_rule_instance( + rule_instance_id.clone(), + author, + author_role.map(AccountRole::from), + ); - queue.send(&payload.to_string()).await?; + queue.send(&serde_json::to_string(&transaction)?).await?; } } } diff --git a/services/event/src/main.rs b/services/event/src/main.rs index e5567080f..6d8bb472e 100644 --- a/services/event/src/main.rs +++ b/services/event/src/main.rs @@ -1,11 +1,25 @@ use futures::channel::mpsc; use futures::{stream, FutureExt, StreamExt, TryStreamExt}; use pg::postgres::DB; +use serde::Deserialize; use shutdown::shutdown_signal; use std::{env, sync::Arc}; use tokio_postgres::{AsyncMessage, NoTls}; mod events; +#[derive(Deserialize)] +#[serde(rename_all = "lowercase")] +enum EventType { + Equilibrium, + Cron, +} + +#[derive(Deserialize)] +struct EventPayload { + event: EventType, + id: String, +} + struct AppState { pool: pg::postgres::ConnectionPool, cache: Arc, @@ -99,13 +113,13 @@ async fn main() { let connection = stream.forward(tx).map(|r| r.unwrap()); let handler = tokio::spawn(connection); - if let Err(e) = client.batch_execute("LISTEN equilibrium;").await { + if let Err(e) = client.batch_execute("LISTEN event;").await { tracing::info!("failed to execute LISTEN command: {}", e); handler.abort(); continue; } - tracing::info!("listening on equilibrium channel"); + tracing::info!("listening on event channel"); // drain rows accumulated during reconnection process_pending(&state).await; @@ -122,8 +136,32 @@ async fn main() { match message { Some(message) => match message { - AsyncMessage::Notification(_) => { - process_pending(&state).await; + AsyncMessage::Notification(n) => { + match serde_json::from_str::(n.payload()) { + Ok(payload) => match payload.event { + EventType::Equilibrium => process_pending(&state).await, + EventType::Cron => { + if let Some(ref q) = state.queue { + if let Err(e) = + events::handle_cron(&state.pool, q, &payload.id).await + { + tracing::error!( + "cron handler error for {}: {}", + payload.id, + e + ); + } + } else { + tracing::error!( + "cron event received but no queue configured" + ); + } + } + }, + Err(e) => { + tracing::error!("failed to parse event payload: {}", e); + } + } } _ => { tracing::info!("unhandled message: {:?}", message); diff --git a/services/rule/README.md b/services/rule/README.md index 13277882e..dfe4c9893 100644 --- a/services/rule/README.md +++ b/services/rule/README.md @@ -2,10 +2,34 @@ systemaccounting

-1. invoked by `graphql` after request sent by client, or by `request-create` service to test if all current and expected items are included in client request -1. queries for rules applicable to items and accounts -1. applies transaction_item and approval rules -1. returns current and expected items +axum service that applies transaction item and approval rules to a transaction + +### how it works + +1. receives a `transaction` from `graphql` (client request) or `request-create` (testing current vs expected items) +1. when `rule_instance_id` is set and `transaction_items` is empty, builds items from `transaction_item_rule_instance` templates in the database (used by auto-transact services like cron and threshold) +1. queries for `transaction_item_rule_instance` rules matching the state and account of each debitor and creditor +1. applies transaction item rules per user-defined role sequence (`debitor_first` or creditor first) +1. queries for `approval_rule_instance` rules matching each account owner (approver) +1. applies approval rules, automating `approval_time` when a rule matches +1. labels each `transaction_item` with `debitor_approval_time` or `creditor_approval_time` when all approvals for that role have timestamps +1. returns the transaction with rule-added items, approvals, and updated `sum_value` + +### transaction item rules + +defined in `src/rules/transaction_item.rs`: + +- **multiplyItemValue**: computes `price * factor` and returns only the computed item (e.g. a dividend). `variable_values = [DEBITOR, CREDITOR, ITEM_NAME, FACTOR]` +- **appendMultipliedItemValue**: computes `price * factor` and returns both the original item (with `rule_exec_id` added) and the computed item (e.g. a sales tax). same variable_values + +the `ANY` token in debitor or creditor position matches the corresponding account from the original transaction item + +### approval rules + +defined in `src/rules/approval.rs`: + +- **approveAnyCreditItem**: automates approval for a creditor approver on any credit item. `variable_values = [CREDITOR, APPROVER_ROLE, APPROVER_NAME]` +- **approveItemBetweenAccounts**: automates approval for a specific debitor/creditor/item_id match. `variable_values = [DEBITOR, CREDITOR, ITEM_ID, APPROVER_ROLE, APPROVER_NAME]` ### request @@ -192,38 +216,16 @@ the rule service returns a `transaction` object with `transaction_items` listing ### dev 1. install [rust](https://doc.rust-lang.org/book/ch01-01-installation.html#installing-rustup-on-linux-or-macos), [cargo-watch](https://crates.io/crates/cargo-watch) and [docker](https://docs.docker.com/get-docker/) 1. start docker -1. `make dev` to start: - 1. the rule service - 1. postgres in docker -1. cntrl + c && `make -C migrations clean` OR `make stop-dev` in a separate shell to stop dev process - -### deploy to lambda -1. `make compile` to build for lambda -1. `make zip-only` to zip lambda binary -1. `make put-object ENV=dev` to put zip in s3 -1. `make update-function ENV=dev` to deploy zip to lambda from s3 - -### deploy to lambda FAST -1. `make deploy ENV=dev` - -### invoke lambda -1. set desired values in `TEST_EVENT` makefile variable -1. `make invoke-function ENV=dev` +1. `make dev` to start the rule service and postgres in docker +1. cntrl + c to stop, then `make -C migrations clean` to clean up ### testing 1. `make test-lint` for [clippy](https://github.com/rust-lang/rust-clippy) 1. `make test-unit` for unit tests +1. `make water` to curl a bottled water request to a running local rule service 1. run integration tests from project root: - * docker: - 1. `make compose-up` - 1. `make test-docker` - - cloud: - 1. `make test-cloud` - -### run unit & integration tests FAST -`make test ENV=dev` - -### prepare for terraform -`make initial-deploy ENV=dev` to zip and put source in s3 only + * docker: `make compose-up && make test-docker` + * cloud: `make -C tests test-cloud ENV=dev` -terraform: https://github.com/systemaccounting/mxfactorial/blob/develop/infra/terraform/aws/modules/environment/v001/lambda-services.tf#L28-L51 +### deploy +`make deploy ENV=dev` to build, push to ecr, and deploy lambda