raw_store_events store started

This commit is contained in:
2025-12-24 00:55:22 -07:00
parent cdbc08dea2
commit 3823fa7331
17 changed files with 363 additions and 73 deletions
@@ -0,0 +1 @@
DROP TABLE raw_store_events;
+10
View File
@@ -0,0 +1,10 @@
CREATE TABLE raw_store_events (
platform TEXT NOT NULL,
store_id TEXT NOT NULL,
event_timestamp TIMESTAMPTZ NOT NULL,
event_id TEXT NOT NULL,
raw_payload JSONB NOT NULL,
parsed BOOLEAN NOT NULL DEFAULT FALSE,
PRIMARY KEY (platform, store_id, event_id, event_timestamp)
);
@@ -0,0 +1,3 @@
DROP TABLE tiktok_store_events;
DROP TABLE etsy_store_events;
DROP TABLE wix_store_events;
@@ -0,0 +1,29 @@
CREATE TABLE tiktok_store_events (
platform TEXT NOT NULL DEFAULT 'tiktok' CHECK (platform = 'tiktok'),
store_id TEXT NOT NULL,
event_timestamp TIMESTAMPTZ NOT NULL,
event_id TEXT NOT NULL,
PRIMARY KEY (store_id, event_id, event_timestamp),
FOREIGN KEY (platform, store_id, event_id, event_timestamp) REFERENCES raw_store_events
);
CREATE TABLE etsy_store_events (
platform TEXT NOT NULL DEFAULT 'etsy' CHECK (platform = 'etsy'),
store_id TEXT NOT NULL,
event_timestamp TIMESTAMPTZ NOT NULL,
event_id TEXT NOT NULL,
PRIMARY KEY (store_id, event_id, event_timestamp),
FOREIGN KEY (platform, store_id, event_id, event_timestamp) REFERENCES raw_store_events
);
CREATE TABLE wix_store_events (
platform TEXT NOT NULL DEFAULT 'wix' CHECK (platform = 'wix'),
store_id TEXT NOT NULL,
event_timestamp TIMESTAMPTZ NOT NULL,
event_id TEXT NOT NULL,
PRIMARY KEY (store_id, event_id, event_timestamp),
FOREIGN KEY (platform, store_id, event_id, event_timestamp) REFERENCES raw_store_events
);
File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 11 KiB

After

Width:  |  Height:  |  Size: 7.1 KiB

+28 -44
View File
@@ -1,56 +1,40 @@
@startuml @startuml
object store_events object raw_store_events
store_events : platform raw_store_events : platform text
store_events : store_id raw_store_events : store_id text
store_events : event_id raw_store_events : event_id text
store_events : event_timestamp raw_store_events : event_timestamp timestamptz
store_events : tiktok_store_id raw_store_events : raw_payload jsonb
store_events : tiktok_event_id raw_store_events : parsed boolean
store_events : <platform B>_store_id
store_events : <platform B>_event_id
store_events : <platform C>_store_id
store_events : <platform C>_event_id
store_events : <platform ..>_store_id
store_events : <platform ..>_event_id
store_events : raw payload?
note right of store_events::"tiktok_event_id" note right of raw_store_events::"parsed"
reference tiktok_store_event_details indicates if the raw_payload has been parsed
end note and stored in one of the referencing tables
note right of store_events::"<platform B>_event_id"
reference <platform B>_store_event_details
end note
note right of store_events::"<platform C>_event_id"
reference <platform C>_store_event_details
end note
note right of store_events::"<platform ..>_event_id"
only one reference can exist at a time
(eg event CANNOT be both from tiktok AND etsy)
end note end note
object tiktok_store_event_details object tiktok_store_events
tiktok_store_event_details : store_id tiktok_store_events : store_id text
tiktok_store_event_details : event_id tiktok_store_events : event_id text
tiktok_store_event_details : event_timestamp tiktok_store_events : event_timestamp timestamptz
tiktok_store_event_details : [attribute: value ...] tiktok_store_events : [attribute: value ...]
object "<platform B>_store_event_details" as platform_b_specific_store_event_details object etsy_store_events
platform_b_specific_store_event_details : store_id etsy_store_events : store_id text
platform_b_specific_store_event_details : event_id etsy_store_events : event_id text
platform_b_specific_store_event_details : event_timestamp etsy_store_events : event_timestamp timestamptz
platform_b_specific_store_event_details : [attribute: value ...] etsy_store_events : [attribute: value ...]
object "<platform C>_store_event_details" as platform_c_specific_store_event_details object etc_store_events
platform_c_specific_store_event_details : store_id etc_store_events : store_id text
platform_c_specific_store_event_details : event_id etc_store_events : event_id text
platform_c_specific_store_event_details : event_timestamp etc_store_events : event_timestamp timestamptz
platform_c_specific_store_event_details : [attribute: value ...] etc_store_events : [attribute: value ...]
store_events <|-- tiktok_store_event_details raw_store_events <|-- tiktok_store_events : (store_id, event_id)
store_events <|-- platform_b_specific_store_event_details raw_store_events <|-- etsy_store_events : (store_id, event_id)
store_events <|-- platform_c_specific_store_event_details raw_store_events <|-- etc_store_events : (store_id, event_id)
@enduml @enduml
File diff suppressed because one or more lines are too long

After

Width:  |  Height:  |  Size: 11 KiB

+56
View File
@@ -0,0 +1,56 @@
@startuml
object store_events
store_events : platform
store_events : store_id
store_events : event_id
store_events : event_timestamp
store_events : tiktok_store_id
store_events : tiktok_event_id
store_events : <platform B>_store_id
store_events : <platform B>_event_id
store_events : <platform C>_store_id
store_events : <platform C>_event_id
store_events : <platform ..>_store_id
store_events : <platform ..>_event_id
store_events : raw payload?
note right of store_events::"tiktok_event_id"
reference tiktok_store_event_details
end note
note right of store_events::"<platform B>_event_id"
reference <platform B>_store_event_details
end note
note right of store_events::"<platform C>_event_id"
reference <platform C>_store_event_details
end note
note right of store_events::"<platform ..>_event_id"
only one reference can exist at a time
(eg event CANNOT be both from tiktok AND etsy)
end note
object tiktok_store_event_details
tiktok_store_event_details : store_id
tiktok_store_event_details : event_id
tiktok_store_event_details : event_timestamp
tiktok_store_event_details : [attribute: value ...]
object "<platform B>_store_event_details" as platform_b_specific_store_event_details
platform_b_specific_store_event_details : store_id
platform_b_specific_store_event_details : event_id
platform_b_specific_store_event_details : event_timestamp
platform_b_specific_store_event_details : [attribute: value ...]
object "<platform C>_store_event_details" as platform_c_specific_store_event_details
platform_c_specific_store_event_details : store_id
platform_c_specific_store_event_details : event_id
platform_c_specific_store_event_details : event_timestamp
platform_c_specific_store_event_details : [attribute: value ...]
store_events <|-- tiktok_store_event_details
store_events <|-- platform_b_specific_store_event_details
store_events <|-- platform_c_specific_store_event_details
@enduml
+11
View File
@@ -1,3 +1,14 @@
module ruben/inventory2 module ruben/inventory2
go 1.24.2 go 1.24.2
require github.com/jackc/pgx/v5 v5.7.6
require (
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect
golang.org/x/crypto v0.41.0 // indirect
golang.org/x/sync v0.16.0 // indirect
golang.org/x/text v0.28.0 // indirect
)
+28
View File
@@ -0,0 +1,28 @@
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.7.6 h1:rWQc5FwZSPX58r1OQmkuaNicxdmExaEz5A2DO2hUuTk=
github.com/jackc/pgx/v5 v5.7.6/go.mod h1:aruU7o91Tc2q2cFp5h4uP3f6ztExVpyVv88Xl/8Vl8M=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk=
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
golang.org/x/crypto v0.41.0 h1:WKYxWedPGCTVVl5+WHSSrOBT0O8lx32+zxmHxijgXp4=
golang.org/x/crypto v0.41.0/go.mod h1:pO5AFd7FA68rFak7rOAGVuygIISepHftHnr8dr6+sUc=
golang.org/x/sync v0.16.0 h1:ycBJEhp9p4vXvUZNszeOq0kGTPghopOL8q0fq3vstxw=
golang.org/x/sync v0.16.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA=
golang.org/x/text v0.28.0 h1:rhazDwis8INMIwQ4tpjLDzUhx6RlXqZNPEM0huQojng=
golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
Binary file not shown.
+89
View File
@@ -0,0 +1,89 @@
package raw_events
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgtype"
"github.com/jackc/pgx/v5/pgxpool"
)
type (
Store struct {
db *pgxpool.Pool
}
Event struct {
Platform string
StoreID string
EventID string
EventTimestamp time.Time
Payload json.RawMessage
}
)
func NewStore(ctx context.Context) (*Store, error) {
pool, err := newPool(ctx)
if err != nil {
return nil, err
}
return &Store{
db: pool,
}, nil
}
func newPool(ctx context.Context) (*pgxpool.Pool, error) {
pool, err := pgxpool.New(ctx, "postgres://app_client:app_password@localhost:5432/inventory_2?sslmode=disable")
if err != nil {
return nil, fmt.Errorf("failed to create database client: %w", err)
}
conn, err := pool.Acquire(ctx)
if err != nil {
return nil, fmt.Errorf("failed to create a database connection: %w", err)
}
conn.Release()
return pool, nil
}
func (db *Store) Save(ctx context.Context, e *Event) error {
_, err := db.db.Exec(
ctx,
`
INSERT INTO raw_store_events (
platform,
store_id,
event_timestamp,
event_id,
raw_payload
)
VALUES (
@platform,
@store_id,
@event_timestamp,
@event_id,
@raw_payload
)
`,
pgx.NamedArgs{
"platform": e.Platform,
"store_id": e.StoreID,
"event_timestamp": pgtype.Timestamptz{
Time: e.EventTimestamp,
Valid: !e.EventTimestamp.IsZero(),
},
"event_id": e.EventID,
"raw_payload": string(e.Payload),
},
)
if err != nil {
return fmt.Errorf("failed to perform query: %w", err)
}
return nil
}
+27 -4
View File
@@ -1,16 +1,39 @@
package etsy package etsy
import ( import (
"encoding/json"
"fmt" "fmt"
"net/http" "net/http"
"time"
"ruben/inventory2/internal/domains/raw_events"
) )
func NewWebhookHandler() http.Handler { func NewWebhookHandler(db *raw_events.Store) http.Handler {
mux := http.NewServeMux() mux := http.NewServeMux()
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { mux.HandleFunc("POST /test", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain") var body json.RawMessage
w.Write([]byte(fmt.Sprintf("received etsy webhook call: %s %s\n", r.Method, r.URL))) if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, "Failed to decode body as json: "+err.Error(), 500)
return
}
ts := time.Now().UTC()
err := db.Save(r.Context(), &raw_events.Event{
Platform: "etsy",
StoreID: "test-store-1",
EventID: fmt.Sprint(ts.Unix()),
EventTimestamp: ts,
Payload: body,
})
if err != nil {
http.Error(w, "Error occurred saving the body as the event payload: "+err.Error(), 500)
return
}
w.WriteHeader(201)
}) })
return mux return mux
+27 -4
View File
@@ -1,16 +1,39 @@
package tiktok package tiktok
import ( import (
"encoding/json"
"fmt" "fmt"
"net/http" "net/http"
"time"
"ruben/inventory2/internal/domains/raw_events"
) )
func NewWebhookHandler() http.Handler { func NewWebhookHandler(db *raw_events.Store) http.Handler {
mux := http.NewServeMux() mux := http.NewServeMux()
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { mux.HandleFunc("POST /test", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain") var body json.RawMessage
w.Write([]byte(fmt.Sprintf("received tiktok webhook call: %s %s\n", r.Method, r.URL))) if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, "Failed to decode body as json: "+err.Error(), 500)
return
}
ts := time.Now().UTC()
err := db.Save(r.Context(), &raw_events.Event{
Platform: "tiktok",
StoreID: "test-store-1",
EventID: fmt.Sprint(ts.Unix()),
EventTimestamp: ts,
Payload: body,
})
if err != nil {
http.Error(w, "Error occurred saving the body as the event payload: "+err.Error(), 500)
return
}
w.WriteHeader(201)
}) })
return mux return mux
+5 -4
View File
@@ -3,17 +3,18 @@ package webhooks
import ( import (
"net/http" "net/http"
"ruben/inventory2/internal/domains/raw_events"
"ruben/inventory2/internal/webhooks/etsy" "ruben/inventory2/internal/webhooks/etsy"
"ruben/inventory2/internal/webhooks/tiktok" "ruben/inventory2/internal/webhooks/tiktok"
"ruben/inventory2/internal/webhooks/wix" "ruben/inventory2/internal/webhooks/wix"
) )
func New() http.Handler { func New(db *raw_events.Store) http.Handler {
wh := http.NewServeMux() wh := http.NewServeMux()
wh.Handle("/etsy/", http.StripPrefix("/etsy", etsy.NewWebhookHandler())) wh.Handle("/etsy/", http.StripPrefix("/etsy", etsy.NewWebhookHandler(db)))
wh.Handle("/tiktok/", http.StripPrefix("/tiktok", tiktok.NewWebhookHandler())) wh.Handle("/tiktok/", http.StripPrefix("/tiktok", tiktok.NewWebhookHandler(db)))
wh.Handle("/wix/", http.StripPrefix("/wix", wix.NewWebhookHandler())) wh.Handle("/wix/", http.StripPrefix("/wix", wix.NewWebhookHandler(db)))
return wh return wh
} }
+27 -4
View File
@@ -1,16 +1,39 @@
package wix package wix
import ( import (
"encoding/json"
"fmt" "fmt"
"net/http" "net/http"
"time"
"ruben/inventory2/internal/domains/raw_events"
) )
func NewWebhookHandler() http.Handler { func NewWebhookHandler(db *raw_events.Store) http.Handler {
mux := http.NewServeMux() mux := http.NewServeMux()
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { mux.HandleFunc("POST /test", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain") var body json.RawMessage
w.Write([]byte(fmt.Sprintf("received wix webhook call: %s %s\n", r.Method, r.URL))) if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, "Failed to decode body as json: "+err.Error(), 500)
return
}
ts := time.Now().UTC()
err := db.Save(r.Context(), &raw_events.Event{
Platform: "wix",
StoreID: "test-store-1",
EventID: fmt.Sprint(ts.Unix()),
EventTimestamp: ts,
Payload: body,
})
if err != nil {
http.Error(w, "Error occurred saving the body as the event payload: "+err.Error(), 500)
return
}
w.WriteHeader(201)
}) })
return mux return mux
+20 -12
View File
@@ -10,6 +10,7 @@ import (
"syscall" "syscall"
"time" "time"
"ruben/inventory2/internal/domains/raw_events"
"ruben/inventory2/internal/site" "ruben/inventory2/internal/site"
"ruben/inventory2/internal/webhooks" "ruben/inventory2/internal/webhooks"
) )
@@ -25,9 +26,16 @@ func runApp(ctx context.Context) error {
ctx, shutdown := context.WithCancel(ctx) ctx, shutdown := context.WithCancel(ctx)
defer shutdown() defer shutdown()
// connect to the database
db, err := raw_events.NewStore(ctx)
if err != nil {
return fmt.Errorf("failed to initialize the raw event store: %w", err)
}
// start http server // start http server
srvErrCh := runServer(ctx) srvErrCh := runServer(ctx, db)
// wait for interrupt signal or unrecoverable failure, then shutdown // wait for interrupt signal or unrecoverable failure, then shutdown
@@ -63,19 +71,10 @@ func runApp(ctx context.Context) error {
return errors.Join(errs...) return errors.Join(errs...)
} }
func buildHTTPHandler() http.Handler { func runServer(ctx context.Context, db *raw_events.Store) <-chan error {
mux := http.NewServeMux()
mux.Handle("/webhooks/", http.StripPrefix("/webhooks", webhooks.New()))
mux.Handle("/site/", http.StripPrefix("/site", site.NewSiteHandler()))
return mux
}
func runServer(ctx context.Context) <-chan error {
srv := &http.Server{ srv := &http.Server{
Addr: ":9000", // local Addr: ":9000", // local
Handler: buildHTTPHandler(), Handler: buildHTTPHandler(db),
} }
ctx, cancel := context.WithCancel(ctx) ctx, cancel := context.WithCancel(ctx)
@@ -125,3 +124,12 @@ func runServer(ctx context.Context) <-chan error {
return errCh return errCh
} }
func buildHTTPHandler(db *raw_events.Store) http.Handler {
mux := http.NewServeMux()
mux.Handle("/webhooks/", http.StripPrefix("/webhooks", webhooks.New(db)))
mux.Handle("/site/", http.StripPrefix("/site", site.NewSiteHandler()))
return mux
}