From 0a03b2e62ed8d24e2eee64e59ac470dda34d0449 Mon Sep 17 00:00:00 2001 From: Angel Beltran Date: Tue, 3 Mar 2026 20:51:00 -0700 Subject: [PATCH] mock amazon event processing stubbed --- ...8_mock_event_processing_trackings.down.sql | 58 +++ ...028_mock_event_processing_trackings.up.sql | 345 ++++++++++++++++++ domains/amazon/mock.go | 183 ++++++++++ main.go | 39 +- 4 files changed, 622 insertions(+), 3 deletions(-) create mode 100644 database_migrations/000028_mock_event_processing_trackings.down.sql create mode 100644 database_migrations/000028_mock_event_processing_trackings.up.sql create mode 100644 domains/amazon/mock.go diff --git a/database_migrations/000028_mock_event_processing_trackings.down.sql b/database_migrations/000028_mock_event_processing_trackings.down.sql new file mode 100644 index 0000000..da0b4bc --- /dev/null +++ b/database_migrations/000028_mock_event_processing_trackings.down.sql @@ -0,0 +1,58 @@ +BEGIN; + + +ALTER TABLE mock.shop_amazon_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_amazon_events_by_timestamp; +ALTER TABLE mock.shop_big_cartel_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_big_cartel_events_by_timestamp; +ALTER TABLE mock.shop_ebay_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_ebay_events_by_timestamp; +ALTER TABLE mock.shop_ecwid_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_ecwid_events_by_timestamp; +ALTER TABLE mock.shop_etsy_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_etsy_events_by_timestamp; +ALTER TABLE mock.shop_shopify_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_shopify_events_by_timestamp; +ALTER TABLE mock.shop_square_online_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_square_online_events_by_timestamp; +ALTER TABLE mock.shop_squarespace_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_squarespace_events_by_timestamp; +ALTER TABLE mock.shop_tiktok_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_tiktok_events_by_timestamp; +ALTER TABLE mock.shop_walmart_marketplace_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_walmart_marketplace_events_by_timestamp; +ALTER TABLE mock.shop_wix_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_wix_events_by_timestamp; +ALTER TABLE mock.shop_woo_commerce_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_woo_commerce_events_by_timestamp; +ALTER TABLE mock.shop_zoho_events + DROP COLUMN processed, + DROP COLUMN processed_successfully; +DROP INDEX mock.shop_zoho_events_by_timestamp; + + +COMMIT; diff --git a/database_migrations/000028_mock_event_processing_trackings.up.sql b/database_migrations/000028_mock_event_processing_trackings.up.sql new file mode 100644 index 0000000..786889b --- /dev/null +++ b/database_migrations/000028_mock_event_processing_trackings.up.sql @@ -0,0 +1,345 @@ +BEGIN; + + +ALTER TABLE mock.shop_amazon_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_amazon_events_by_timestamp ON mock.shop_amazon_events (shop_id, event_timestamp); +CREATE INDEX shop_amazon_events_by_timestamp_unprocessed ON mock.shop_amazon_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_big_cartel_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_big_cartel_events_by_timestamp ON mock.shop_big_cartel_events (shop_id, event_timestamp); +CREATE INDEX shop_big_cartel_events_by_timestamp_unprocessed ON mock.shop_big_cartel_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_ebay_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_ebay_events_by_timestamp ON mock.shop_ebay_events (shop_id, event_timestamp); +CREATE INDEX shop_ebay_events_by_timestamp_unprocessed ON mock.shop_ebay_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_ecwid_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_ecwid_events_by_timestamp ON mock.shop_ecwid_events (shop_id, event_timestamp); +CREATE INDEX shop_ecwid_events_by_timestamp_unprocessed ON mock.shop_ecwid_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_etsy_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_etsy_events_by_timestamp ON mock.shop_etsy_events (shop_id, event_timestamp); +CREATE INDEX shop_etsy_events_by_timestamp_unprocessed ON mock.shop_etsy_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_shopify_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_shopify_events_by_timestamp ON mock.shop_shopify_events (shop_id, event_timestamp); +CREATE INDEX shop_shopify_events_by_timestamp_unprocessed ON mock.shop_shopify_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_square_online_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_square_online_events_by_timestamp ON mock.shop_square_online_events (shop_id, event_timestamp); +CREATE INDEX shop_square_online_events_by_timestamp_unprocessed ON mock.shop_square_online_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_squarespace_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_squarespace_events_by_timestamp ON mock.shop_squarespace_events (shop_id, event_timestamp); +CREATE INDEX shop_squarespace_events_by_timestamp_unprocessed ON mock.shop_squarespace_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_tiktok_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_tiktok_events_by_timestamp ON mock.shop_tiktok_events (shop_id, event_timestamp); +CREATE INDEX shop_tiktok_events_by_timestamp_unprocessed ON mock.shop_tiktok_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_walmart_marketplace_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_walmart_marketplace_events_by_timestamp ON mock.shop_walmart_marketplace_events (shop_id, event_timestamp); +CREATE INDEX shop_walmart_marketplace_events_by_timestamp_unprocessed ON mock.shop_walmart_marketplace_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_wix_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_wix_events_by_timestamp ON mock.shop_wix_events (shop_id, event_timestamp); +CREATE INDEX shop_wix_events_by_timestamp_unprocessed ON mock.shop_wix_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_woo_commerce_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_woo_commerce_events_by_timestamp ON mock.shop_woo_commerce_events (shop_id, event_timestamp); +CREATE INDEX shop_woo_commerce_events_by_timestamp_unprocessed ON mock.shop_woo_commerce_events (shop_id, event_timestamp) WHERE NOT processed; + +ALTER TABLE mock.shop_zoho_events + ADD COLUMN processed BOOLEAN DEFAULT FALSE, + ADD COLUMN processed_successfully BOOLEAN DEFAULT FALSE; +CREATE INDEX shop_zoho_events_by_timestamp ON mock.shop_zoho_events (shop_id, event_timestamp); +CREATE INDEX shop_zoho_events_by_timestamp_unprocessed ON mock.shop_zoho_events (shop_id, event_timestamp) WHERE NOT processed; + + +--- event processing + +CREATE OR REPLACE FUNCTION mock.process_raw_amazon_event() RETURNS TRIGGER AS $process_raw_amazon_event$ + BEGIN + INSERT INTO mock.shop_amazon_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_amazon_event_inserted', null); + + RETURN NULL; + END; +$process_raw_amazon_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_big_cartel_event() RETURNS TRIGGER AS $process_raw_big_cartel_event$ + BEGIN + INSERT INTO mock.shop_big_cartel_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_big_cartel_event_inserted', null); + + RETURN NULL; + END; +$process_raw_big_cartel_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_ebay_event() RETURNS TRIGGER AS $process_raw_ebay_event$ + BEGIN + INSERT INTO mock.shop_ebay_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_ebay_event_inserted', null); + + RETURN NULL; + END; +$process_raw_ebay_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_ecwid_event() RETURNS TRIGGER AS $process_raw_ecwid_event$ + BEGIN + INSERT INTO mock.shop_ecwid_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_ecwid_event_inserted', null); + + RETURN NULL; + END; +$process_raw_ecwid_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_etsy_event() RETURNS TRIGGER AS $process_raw_etsy_event$ + BEGIN + INSERT INTO mock.shop_etsy_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_etsy_event_inserted', null); + + RETURN NULL; + END; +$process_raw_etsy_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_shopify_event() RETURNS TRIGGER AS $process_raw_shopify_event$ + BEGIN + INSERT INTO mock.shop_shopify_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_shopify_event_inserted', null); + + RETURN NULL; + END; +$process_raw_shopify_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_square_online_event() RETURNS TRIGGER AS $process_raw_square_online_event$ + BEGIN + INSERT INTO mock.shop_square_online_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_square_online_event_inserted', null); + + RETURN NULL; + END; +$process_raw_square_online_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_squarespace_event() RETURNS TRIGGER AS $process_raw_squarespace_event$ + BEGIN + INSERT INTO mock.shop_squarespace_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_squarespace_event_inserted', null); + + RETURN NULL; + END; +$process_raw_squarespace_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_tiktok_event() RETURNS TRIGGER AS $process_raw_tiktok_event$ + BEGIN + INSERT INTO mock.shop_tiktok_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_tiktok_event_inserted', null); + + RETURN NULL; + END; +$process_raw_tiktok_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_walmart_marketplace_event() RETURNS TRIGGER AS $process_raw_walmart_marketplace_event$ + BEGIN + INSERT INTO mock.shop_walmart_marketplace_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_walmart_marketplace_event_inserted', null); + + RETURN NULL; + END; +$process_raw_walmart_marketplace_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_wix_event() RETURNS TRIGGER AS $process_raw_wix_event$ + BEGIN + INSERT INTO mock.shop_wix_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_wix_event_inserted', null); + + RETURN NULL; + END; +$process_raw_wix_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_woo_commerce_event() RETURNS TRIGGER AS $process_raw_woo_commerce_event$ + BEGIN + INSERT INTO mock.shop_woo_commerce_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_woo_commerce_event_inserted', null); + + RETURN NULL; + END; +$process_raw_woo_commerce_event$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION mock.process_raw_zoho_event() RETURNS TRIGGER AS $process_raw_zoho_event$ + BEGIN + INSERT INTO mock.shop_zoho_events ( + platform, + shop_id, + event_timestamp, + event_id + ) + SELECT + NEW.platform, + NEW.shop_id, + NEW.event_timestamp, + NEW.event_id; + + PERFORM pg_notify('mock_shop_zoho_event_inserted', null); + + RETURN NULL; + END; +$process_raw_zoho_event$ LANGUAGE plpgsql; + +COMMIT; diff --git a/domains/amazon/mock.go b/domains/amazon/mock.go new file mode 100644 index 0000000..5653ac8 --- /dev/null +++ b/domains/amazon/mock.go @@ -0,0 +1,183 @@ +package amazon + +import ( + "context" + "errors" + "fmt" + "ruben/inventory2/logging" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +type ( + Mocks struct { + log *logging.Logger + db *pgxpool.Pool + } + + event struct { + ShopID string + EventTimestamp time.Time + EventID string + } +) + +const ( + eventChannelName = "mock_shop_amazon_event_inserted" +) + +func NewMocks(log *logging.Logger, db *pgxpool.Pool) *Mocks { + return &Mocks{ + log: log, + db: db, + } +} + +func (m *Mocks) ProcessEvents(ctx context.Context) error { + notifCh, errCh := m.listenForNotifications(ctx) + + for { + m.log.Debug("processing events") + if err := m.processUnprocessedEvents(ctx); err != nil { + if errors.Is(err, context.Canceled) { + return nil + } + return fmt.Errorf("error occurred while processing events: %w", err) + } + + select { + case <-ctx.Done(): + return nil + case err := <-errCh: + return err + case _, ok := <-notifCh: + if !ok { + return <-errCh + } + case <-time.After(time.Minute): + } + } +} + +func (m *Mocks) listenForNotifications(ctx context.Context) (<-chan struct{}, <-chan error) { + errCh := make(chan error, 1) + + pc, err := m.db.Acquire(ctx) + if err != nil { + errCh <- fmt.Errorf("failed to acquire a connection: %w", err) + return nil, errCh + } + + conn := pc.Conn() + + if _, err := conn.Exec(ctx, fmt.Sprintf("LISTEN %s", eventChannelName)); err != nil { + errCh <- fmt.Errorf("failed to start listening for notifications: %w", err) + return nil, errCh + } + + ch := make(chan struct{}) + + go func() (err error) { + defer func() { + pc.Release() + if err != nil { + errCh <- err + } + close(ch) + close(errCh) + }() + + for { + if _, err := conn.WaitForNotification(ctx); err != nil { + if errors.Is(err, context.Canceled) { + return nil + } + return fmt.Errorf("error occurred while waiting for the next notification: %w", err) + } + + ch <- struct{}{} + } + }() + + return ch, errCh +} + +func (m *Mocks) processUnprocessedEvents(ctx context.Context) error { + for done := false; !done; { + err := pgx.BeginFunc(ctx, m.db, func(tx pgx.Tx) error { + rows, err := m.db.Query( + ctx, + ` + SELECT + shop_id AS shopid, + event_id AS eventid, + event_timestamp AS eventtimestamp + FROM + mock.shop_amazon_events + WHERE + NOT processed + ORDER BY + shop_id, event_timestamp + LIMIT + 100 + `, + ) + if err != nil { + return fmt.Errorf("failed to to perform query: %w", err) + } + + evts, err := pgx.CollectRows(rows, pgx.RowToStructByNameLax[event]) + if err != nil { + return fmt.Errorf("failed to scan rows: %w", err) + } + + for _, e := range evts { + if err := m.processEvent(ctx, tx, e); err != nil { + return fmt.Errorf("failed to process event: %w", err) + } + } + + if done = len(evts) == 0; !done { + m.log.Infof("processed %d events", len(evts)) + } + + return nil + }) + if err != nil { + return err + } + } + + return nil +} + +func (m *Mocks) processEvent(ctx context.Context, tx pgx.Tx, e event) error { + m.log.Debugf("processing (mock) event: %#v", e) + + _, err := tx.Exec( + ctx, + ` + UPDATE + mock.shop_amazon_events + SET + processed = true, + processed_successfully = true + WHERE + shop_id = @shop_id + AND event_id = @event_id + AND event_timestamp = @event_timestamp + `, + pgx.NamedArgs{ + "shop_id": e.ShopID, + "event_id": e.EventID, + "event_timestamp": e.EventTimestamp, + }, + ) + if err != nil { + return fmt.Errorf("failed to perform query: %w", err) + } + + return nil +} diff --git a/main.go b/main.go index c24a8c1..b9f2edd 100644 --- a/main.go +++ b/main.go @@ -22,6 +22,7 @@ import ( "github.com/lmittmann/tint" "ruben/inventory2/domains/accounts" + "ruben/inventory2/domains/amazon" "ruben/inventory2/domains/authentication" etsy_platform "ruben/inventory2/domains/platforms/etsy" "ruben/inventory2/domains/raw_events" @@ -82,6 +83,11 @@ func runApp(ctx context.Context, logger *logging.Logger) error { authErrCh := runAuthProcesses(ctx, auth) + eventErrCh := runEventProcessing( + ctx, + amazon.NewMocks(logger.WithGroup("amazon"), connPool), + ) + // start http server srvErrCh := runServer(ctx, logger, connPool, auth) @@ -93,8 +99,9 @@ func runApp(ctx context.Context, logger *logging.Logger) error { var ( alreadyShutdown struct { - server bool - authProcesses bool + server bool + authProcesses bool + eventProcessing bool } ) select { @@ -107,6 +114,12 @@ func runApp(ctx context.Context, logger *logging.Logger) error { if err != nil { logger.Error("auth processes encountered error", "error", err) } + case err := <-eventErrCh: + alreadyShutdown.eventProcessing = true + logger.Error("event processing shutdown unexpectedly") + if err != nil { + logger.Error("event processing encountered error", "error", err) + } case err := <-srvErrCh: alreadyShutdown.server = true logger.Error("server shutdown unexpectedly") @@ -128,6 +141,13 @@ func runApp(ctx context.Context, logger *logging.Logger) error { logger.Info("auth processes shut down") } + if !alreadyShutdown.eventProcessing { + if err := <-authErrCh; err != nil { + errs = append(errs, fmt.Errorf("event processing experienced an error: %w", err)) + } + logger.Info("event processing shut down") + } + if !alreadyShutdown.server { if err := <-srvErrCh; err != nil { errs = append(errs, fmt.Errorf("server experienced an error: %w", err)) @@ -151,8 +171,21 @@ func runAuthProcesses(ctx context.Context, auth *authentication.Authenticator) < return errCh } +func runEventProcessing(ctx context.Context, amz *amazon.Mocks) <-chan error { + errCh := make(chan error, 1) + go func() { + defer close(errCh) + + if err := amz.ProcessEvents(ctx); err != nil { + errCh <- err + } + }() + + return errCh +} + func runServer(ctx context.Context, logger *logging.Logger, connPool *pgxpool.Pool, auth *authentication.Authenticator) <-chan error { - accts := accounts.NewStore(logger, connPool) + accts := accounts.NewStore(logger.WithGroup("accounts"), connPool) r := server.NewRouter( logger.WithGroup("server"),