| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 45bb281 commit 17da946
15 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -37,6 +37,21 @@ jobs: | |||
| 37 | 37 | - name: Wait for services to be healthy | |
| 38 | 38 | working-directory: ./tests | |
| 39 | 39 | run: | | |
| 40 | + echo "Waiting for sqld to be healthy..." | ||
| 41 | + for i in $(seq 1 20); do | ||
| 42 | + if curl -sf http://localhost:8090/health >/dev/null 2>&1; then | ||
| 43 | + echo "sqld is healthy!" | ||
| 44 | + break | ||
| 45 | + fi | ||
| 46 | + if [ $i -eq 20 ]; then | ||
| 47 | + echo "sqld failed to become healthy" | ||
| 48 | + docker compose logs sqld | ||
| 49 | + exit 1 | ||
| 50 | + fi | ||
| 51 | + echo "sqld attempt $i/20 - waiting 3s..." | ||
| 52 | + sleep 3 | ||
| 53 | + done | ||
| 54 | + | ||
| 40 | 55 | echo "Waiting for API to be healthy..." | |
| 41 | 56 | for i in $(seq 1 40); do | |
| 42 | 57 | if docker compose exec api curl -sf http://localhost:8000/health >/dev/null 2>&1; then | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -44,6 +44,7 @@ require ( | |||
| 44 | 44 | github.com/stretchr/testify v1.11.1 | |
| 45 | 45 | github.com/swaggo/swag v1.16.6 | |
| 46 | 46 | github.com/thedevsaddam/govalidator v1.9.10 | |
| 47 | + github.com/tursodatabase/libsql-client-go v0.0.0-20260514053736-a9a8fadfe885 | ||
| 47 | 48 | github.com/uptrace/uptrace-go v1.43.0 | |
| 48 | 49 | github.com/xuri/excelize/v2 v2.10.1 | |
| 49 | 50 | go.opentelemetry.io/otel v1.43.0 | |
@@ -93,11 +94,13 @@ require ( | |||
| 93 | 94 | github.com/PuerkitoBio/goquery v1.12.0 // indirect | |
| 94 | 95 | github.com/andybalholm/brotli v1.2.1 // indirect | |
| 95 | 96 | github.com/andybalholm/cascadia v1.3.3 // indirect | |
| 97 | + github.com/antlr4-go/antlr/v4 v4.13.0 // indirect | ||
| 96 | 98 | github.com/cenkalti/backoff/v5 v5.0.3 // indirect | |
| 97 | 99 | github.com/cespare/xxhash/v2 v2.3.0 // indirect | |
| 98 | 100 | github.com/clipperhouse/displaywidth v0.11.0 // indirect | |
| 99 | 101 | github.com/clipperhouse/uax29/v2 v2.7.0 // indirect | |
| 100 | 102 | github.com/cncf/xds/go v0.0.0-20260202195803-dba9d589def2 // indirect | |
| 103 | + github.com/coder/websocket v1.8.12 // indirect | ||
| 101 | 104 | github.com/envoyproxy/go-control-plane/envoy v1.37.0 // indirect | |
| 102 | 105 | github.com/envoyproxy/protoc-gen-validate v1.3.3 // indirect | |
| 103 | 106 | github.com/fatih/color v1.19.0 // indirect | |
@@ -183,6 +186,7 @@ require ( | |||
| 183 | 186 | go.uber.org/zap v1.28.0 // indirect | |
| 184 | 187 | go.yaml.in/yaml/v3 v3.0.4 // indirect | |
| 185 | 188 | golang.org/x/crypto v0.50.0 // indirect | |
| 189 | + golang.org/x/exp v0.0.0-20240325151524-a685a6edb6d8 // indirect | ||
| 186 | 190 | golang.org/x/mod v0.35.0 // indirect | |
| 187 | 191 | golang.org/x/net v0.53.0 // indirect | |
| 188 | 192 | golang.org/x/oauth2 v0.36.0 // indirect | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -66,6 +66,8 @@ github.com/andybalholm/brotli v1.2.1 h1:R+f5xP285VArJDRgowrfb9DqL18yVK0gKAW/F+eT | |||
| 66 | 66 | github.com/andybalholm/brotli v1.2.1/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY= | |
| 67 | 67 | github.com/andybalholm/cascadia v1.3.3 h1:AG2YHrzJIm4BZ19iwJ/DAua6Btl3IwJX+VI4kktS1LM= | |
| 68 | 68 | github.com/andybalholm/cascadia v1.3.3/go.mod h1:xNd9bqTn98Ln4DwST8/nG+H0yuB8Hmgu1YHNnWw0GeA= | |
| 69 | + github.com/antlr4-go/antlr/v4 v4.13.0 h1:lxCg3LAv+EUK6t1i0y1V6/SLeUi0eKEKdhQAlS8TVTI= | ||
| 70 | + github.com/antlr4-go/antlr/v4 v4.13.0/go.mod h1:pfChB/xh/Unjila75QW7+VU4TSnWnnk9UTnmpPaOR2g= | ||
| 69 | 71 | github.com/avast/retry-go/v5 v5.0.0 h1:kf1Qc2UsTZ4qq8elDymqfbISvkyMuhgRxuJqX2NHP7k= | |
| 70 | 72 | github.com/avast/retry-go/v5 v5.0.0/go.mod h1://d+usmKWio1agtZfS1H/ltTqwtIfBnRq9zEwjc3eH8= | |
| 71 | 73 | github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= | |
@@ -88,6 +90,8 @@ github.com/cncf/xds/go v0.0.0-20260202195803-dba9d589def2 h1:aBangftG7EVZoUb69Os | |||
| 88 | 90 | github.com/cncf/xds/go v0.0.0-20260202195803-dba9d589def2/go.mod h1:qwXFYgsP6T7XnJtbKlf1HP8AjxZZyzxMmc+Lq5GjlU4= | |
| 89 | 91 | github.com/cockroachdb/cockroach-go/v2 v2.4.3 h1:LJO3K3jC5WXvMePRQSJE1NsIGoFGcEx1LW83W6RAlhw= | |
| 90 | 92 | github.com/cockroachdb/cockroach-go/v2 v2.4.3/go.mod h1:9U179XbCx4qFWtNhc7BiWLPfuyMVQ7qdAhfrwLz1vH0= | |
| 93 | + github.com/coder/websocket v1.8.12 h1:5bUXkEPPIbewrnkU8LTCLVaxi4N4J8ahufH2vlo4NAo= | ||
| 94 | + github.com/coder/websocket v1.8.12/go.mod h1:LNVeNrXQZfe5qhS9ALED3uA+l5pPqvwXg3CKoDBB2gs= | ||
| 91 | 95 | github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= | |
| 92 | 96 | github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= | |
| 93 | 97 | github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= | |
@@ -320,6 +324,8 @@ github.com/thedevsaddam/govalidator v1.9.10 h1:m3dLRbSZ5Hts3VUWYe+vxLMG+FdyQuWOj | |||
| 320 | 324 | github.com/thedevsaddam/govalidator v1.9.10/go.mod h1:Ilx8u7cg5g3LXbSS943cx5kczyNuUn7LH/cK5MYuE90= | |
| 321 | 325 | github.com/tiendc/go-deepcopy v1.7.2 h1:Ut2yYR7W9tWjTQitganoIue4UGxZwCcJy3orjrrIj44= | |
| 322 | 326 | github.com/tiendc/go-deepcopy v1.7.2/go.mod h1:4bKjNC2r7boYOkD2IOuZpYjmlDdzjbpTRyCx+goBCJQ= | |
| 327 | + github.com/tursodatabase/libsql-client-go v0.0.0-20260514053736-a9a8fadfe885 h1:YssVXwM/9nUAjGNmUWdgvb05JVcsaBrDn5yr+MaJTn0= | ||
| 328 | + github.com/tursodatabase/libsql-client-go v0.0.0-20260514053736-a9a8fadfe885/go.mod h1:08inkKyguB6CGGssc/JzhmQWwBgFQBgjlYFjxjRh7nU= | ||
| 323 | 329 | github.com/uptrace/uptrace-go v1.43.0 h1:5QuCdyFJdWUEXx6Fr6sYfezdgO6n6lnkOvUTLlyQO7U= | |
| 324 | 330 | github.com/uptrace/uptrace-go v1.43.0/go.mod h1:ehDTIdtBSolg4Z0CCvg1C8yR6VX1YFDqBcg2KmsXWn0= | |
| 325 | 331 | github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= | |
@@ -410,6 +416,8 @@ golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v | |||
| 410 | 416 | golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk= | |
| 411 | 417 | golang.org/x/crypto v0.50.0 h1:zO47/JPrL6vsNkINmLoo/PH1gcxpls50DNogFvB5ZGI= | |
| 412 | 418 | golang.org/x/crypto v0.50.0/go.mod h1:3muZ7vA7PBCE6xgPX7nkzzjiUq87kRItoJQM1Yo8S+Q= | |
| 419 | + golang.org/x/exp v0.0.0-20240325151524-a685a6edb6d8 h1:aAcj0Da7eBAtrTp03QXWvm88pSyOt+UgdZw2BFZ+lEw= | ||
| 420 | + golang.org/x/exp v0.0.0-20240325151524-a685a6edb6d8/go.mod h1:CQ1k9gNrJ50XIzaKCRR2hssIjF07kZFEiieALBM/ARQ= | ||
| 413 | 421 | golang.org/x/image v0.25.0 h1:Y6uW6rH1y5y/LK1J8BPWZtr6yZ7hrsy6hFrXjgsc2fQ= | |
| 414 | 422 | golang.org/x/image v0.25.0/go.mod h1:tCAmOEGthTtkalusGp1g3xa2gke8J6c2N565dTyl9Rs= | |
| 415 | 423 | golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -3,6 +3,7 @@ package di | |||
| 3 | 3 | import ( | |
| 4 | 4 | "context" | |
| 5 | 5 | "crypto/tls" | |
| 6 | + "database/sql" | ||
| 6 | 7 | "fmt" | |
| 7 | 8 | "net/http" | |
| 8 | 9 | "os" | |
@@ -82,6 +83,7 @@ type Container struct { | |||
| 82 | 83 | projectID string | |
| 83 | 84 | db *gorm.DB | |
| 84 | 85 | dedicatedDB *gorm.DB | |
| 86 | + tursoDB *sql.DB | ||
| 85 | 87 | version string | |
| 86 | 88 | app *fiber.App | |
| 87 | 89 | eventDispatcher *services.EventDispatcher | |
@@ -269,19 +271,16 @@ func (container *Container) DedicatedDB() (db *gorm.DB) { | |||
| 269 | 271 | container.logger.Fatal(err) | |
| 270 | 272 | } | |
| 271 | 273 | ||
| 272 | - sqlDB, err := db.DB() | ||
| 273 | - if err != nil { | ||
| 274 | - container.logger.Fatal(stacktrace.Propagate(err, "cannot get sql.DB from GORM")) | ||
| 275 | - } | ||
| 276 | - | ||
| 277 | - sqlDB.SetMaxOpenConns(1) | ||
| 278 | - sqlDB.SetMaxIdleConns(0) | ||
| 279 | - sqlDB.SetConnMaxLifetime(10 * time.Second) | ||
| 280 | - | ||
| 281 | 274 | if err = db.Use(tracing.NewPlugin()); err != nil { | |
| 282 | 275 | container.logger.Fatal(stacktrace.Propagate(err, "cannot use GORM tracing plugin")) | |
| 283 | 276 | } | |
| 284 | 277 | ||
| 278 | + container.dedicatedDB = db | ||
| 279 | + if os.Getenv("DATABASE_MIGRATION_SKIP") != "" { | ||
| 280 | + container.logger.Debug(fmt.Sprintf("skipping migrations for [%T]", db)) | ||
| 281 | + return container.dedicatedDB | ||
| 282 | + } | ||
| 283 | + | ||
| 285 | 284 | container.logger.Debug(fmt.Sprintf("Running migrations for dedicated [%T]", db)) | |
| 286 | 285 | if err = db.AutoMigrate(&entities.Heartbeat{}); err != nil { | |
| 287 | 286 | container.logger.Fatal(stacktrace.Propagate(err, fmt.Sprintf("cannot migrate %T", &entities.Heartbeat{}))) | |
@@ -291,10 +290,43 @@ func (container *Container) DedicatedDB() (db *gorm.DB) { | |||
| 291 | 290 | container.logger.Fatal(stacktrace.Propagate(err, fmt.Sprintf("cannot migrate %T", &entities.HeartbeatMonitor{}))) | |
| 292 | 291 | } | |
| 293 | 292 | ||
| 294 | - container.dedicatedDB = db | ||
| 295 | 293 | return container.dedicatedDB | |
| 296 | 294 | } | |
| 297 | 295 | ||
| 296 | + // TursoDB creates a *sql.DB connection to a Turso/libSQL database | ||
| 297 | + func (container *Container) TursoDB() *sql.DB { | ||
| 298 | + if container.tursoDB != nil { | ||
| 299 | + return container.tursoDB | ||
| 300 | + } | ||
| 301 | + | ||
| 302 | + container.logger.Debug("creating Turso *sql.DB connection") | ||
| 303 | + | ||
| 304 | + db, err := repositories.NewTursoDB(os.Getenv("TURSO_DATABASE_DSN")) | ||
| 305 | + if err != nil { | ||
| 306 | + container.logger.Fatal(err) | ||
| 307 | + } | ||
| 308 | + | ||
| 309 | + container.tursoDB = db | ||
| 310 | + return container.tursoDB | ||
| 311 | + } | ||
| 312 | + | ||
| 313 | + // HedgingFailureCounter creates an OTel counter for hedging secondary write failures | ||
| 314 | + func (container *Container) HedgingFailureCounter() otelMetric.Int64Counter { | ||
| 315 | + meter := otel.GetMeterProvider().Meter( | ||
| 316 | + container.projectID, | ||
| 317 | + otelMetric.WithInstrumentationVersion(otel.Version()), | ||
| 318 | + ) | ||
| 319 | + counter, err := meter.Int64Counter( | ||
| 320 | + "hedging.secondary.write.failures", | ||
| 321 | + otelMetric.WithUnit("1"), | ||
| 322 | + otelMetric.WithDescription("Number of failed secondary writes in hedging repositories"), | ||
| 323 | + ) | ||
| 324 | + if err != nil { | ||
| 325 | + container.logger.Fatal(stacktrace.Propagate(err, "cannot create hedging failure counter")) | ||
| 326 | + } | ||
| 327 | + return counter | ||
| 328 | + } | ||
| 329 | + | ||
| 298 | 330 | // DBWithoutMigration creates an instance of gorm.DB if it has not been created already | |
| 299 | 331 | func (container *Container) DBWithoutMigration() (db *gorm.DB) { | |
| 300 | 332 | if container.db != nil { | |
@@ -889,12 +921,31 @@ func (container *Container) MessageThreadRepository() (repository repositories.M | |||
| 889 | 921 | ||
| 890 | 922 | // HeartbeatMonitorRepository creates a new instance of repositories.HeartbeatMonitorRepository | |
| 891 | 923 | func (container *Container) HeartbeatMonitorRepository() (repository repositories.HeartbeatMonitorRepository) { | |
| 892 | - container.logger.Debug("creating GORM repositories.HeartbeatMonitorRepository") | ||
| 893 | - return repositories.NewGormHeartbeatMonitorRepository( | ||
| 894 | - container.Logger(), | ||
| 895 | - container.Tracer(), | ||
| 896 | - container.DedicatedDB(), | ||
| 897 | - ) | ||
| 924 | + switch os.Getenv("HEARTBEAT_DB_BACKEND") { | ||
| 925 | + case "turso": | ||
| 926 | + container.logger.Debug("creating libSQL repositories.HeartbeatMonitorRepository") | ||
| 927 | + return repositories.NewLibsqlHeartbeatMonitorRepository( | ||
| 928 | + container.Logger(), | ||
| 929 | + container.Tracer(), | ||
| 930 | + container.TursoDB(), | ||
| 931 | + ) | ||
| 932 | + case "hedging": | ||
| 933 | + container.logger.Debug("creating hedging repositories.HeartbeatMonitorRepository") | ||
| 934 | + return repositories.NewHedgingHeartbeatMonitorRepository( | ||
| 935 | + container.Logger(), | ||
| 936 | + container.Tracer(), | ||
| 937 | + repositories.NewGormHeartbeatMonitorRepository(container.Logger(), container.Tracer(), container.DedicatedDB()), | ||
| 938 | + repositories.NewLibsqlHeartbeatMonitorRepository(container.Logger(), container.Tracer(), container.TursoDB()), | ||
| 939 | + container.HedgingFailureCounter(), | ||
| 940 | + ) | ||
| 941 | + default: | ||
| 942 | + container.logger.Debug("creating GORM repositories.HeartbeatMonitorRepository") | ||
| 943 | + return repositories.NewGormHeartbeatMonitorRepository( | ||
| 944 | + container.Logger(), | ||
| 945 | + container.Tracer(), | ||
| 946 | + container.DedicatedDB(), | ||
| 947 | + ) | ||
| 948 | + } | ||
| 898 | 949 | } | |
| 899 | 950 | ||
| 900 | 951 | // HeartbeatService creates a new instance of services.HeartbeatService | |
@@ -1708,12 +1759,31 @@ func (container *Container) RegisterSwaggerRoutes() { | |||
| 1708 | 1759 | ||
| 1709 | 1760 | // HeartbeatRepository registers a new instance of repositories.HeartbeatRepository | |
| 1710 | 1761 | func (container *Container) HeartbeatRepository() repositories.HeartbeatRepository { | |
| 1711 | - container.logger.Debug("creating GORM repositories.HeartbeatRepository") | ||
| 1712 | - return repositories.NewGormHeartbeatRepository( | ||
| 1713 | - container.Logger(), | ||
| 1714 | - container.Tracer(), | ||
| 1715 | - container.DedicatedDB(), | ||
| 1716 | - ) | ||
| 1762 | + switch os.Getenv("HEARTBEAT_DB_BACKEND") { | ||
| 1763 | + case "turso": | ||
| 1764 | + container.logger.Debug("creating libSQL repositories.HeartbeatRepository") | ||
| 1765 | + return repositories.NewLibsqlHeartbeatRepository( | ||
| 1766 | + container.Logger(), | ||
| 1767 | + container.Tracer(), | ||
| 1768 | + container.TursoDB(), | ||
| 1769 | + ) | ||
| 1770 | + case "hedging": | ||
| 1771 | + container.logger.Debug("creating hedging repositories.HeartbeatRepository") | ||
| 1772 | + return repositories.NewHedgingHeartbeatRepository( | ||
| 1773 | + container.Logger(), | ||
| 1774 | + container.Tracer(), | ||
| 1775 | + repositories.NewGormHeartbeatRepository(container.Logger(), container.Tracer(), container.DedicatedDB()), | ||
| 1776 | + repositories.NewLibsqlHeartbeatRepository(container.Logger(), container.Tracer(), container.TursoDB()), | ||
| 1777 | + container.HedgingFailureCounter(), | ||
| 1778 | + ) | ||
| 1779 | + default: | ||
| 1780 | + container.logger.Debug("creating GORM repositories.HeartbeatRepository") | ||
| 1781 | + return repositories.NewGormHeartbeatRepository( | ||
| 1782 | + container.Logger(), | ||
| 1783 | + container.Tracer(), | ||
| 1784 | + container.DedicatedDB(), | ||
| 1785 | + ) | ||
| 1786 | + } | ||
| 1717 | 1787 | } | |
| 1718 | 1788 | ||
| 1719 | 1789 | // UserRepository registers a new instance of repositories.UserRepository | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -0,0 +1,128 @@ | |||
| 1 | + package repositories | ||
| 2 | + | ||
| 3 | + import ( | ||
| 4 | + "context" | ||
| 5 | + "fmt" | ||
| 6 | + | ||
| 7 | + "github.com/google/uuid" | ||
| 8 | + otelMetric "go.opentelemetry.io/otel/metric" | ||
| 9 | + | ||
| 10 | + "github.com/NdoleStudio/httpsms/pkg/entities" | ||
| 11 | + "github.com/NdoleStudio/httpsms/pkg/telemetry" | ||
| 12 | + "github.com/palantir/stacktrace" | ||
| 13 | + ) | ||
| 14 | + | ||
| 15 | + // hedgingHeartbeatMonitorRepository writes to both primary and secondary repositories. | ||
| 16 | + // Reads only hit primary. Secondary writes are fail-open. | ||
| 17 | + type hedgingHeartbeatMonitorRepository struct { | ||
| 18 | + logger telemetry.Logger | ||
| 19 | + tracer telemetry.Tracer | ||
| 20 | + primary HeartbeatMonitorRepository | ||
| 21 | + secondary HeartbeatMonitorRepository | ||
| 22 | + failureCounter otelMetric.Int64Counter | ||
| 23 | + } | ||
| 24 | + | ||
| 25 | + // NewHedgingHeartbeatMonitorRepository creates a hedging HeartbeatMonitorRepository | ||
| 26 | + func NewHedgingHeartbeatMonitorRepository( | ||
| 27 | + logger telemetry.Logger, | ||
| 28 | + tracer telemetry.Tracer, | ||
| 29 | + primary HeartbeatMonitorRepository, | ||
| 30 | + secondary HeartbeatMonitorRepository, | ||
| 31 | + failureCounter otelMetric.Int64Counter, | ||
| 32 | + ) HeartbeatMonitorRepository { | ||
| 33 | + return &hedgingHeartbeatMonitorRepository{ | ||
| 34 | + logger: logger.WithService(fmt.Sprintf("%T", &hedgingHeartbeatMonitorRepository{})), | ||
| 35 | + tracer: tracer, | ||
| 36 | + primary: primary, | ||
| 37 | + secondary: secondary, | ||
| 38 | + failureCounter: failureCounter, | ||
| 39 | + } | ||
| 40 | + } | ||
| 41 | + | ||
| 42 | + func (repository *hedgingHeartbeatMonitorRepository) Store(ctx context.Context, monitor *entities.HeartbeatMonitor) error { | ||
| 43 | + ctx, span := repository.tracer.Start(ctx) | ||
| 44 | + defer span.End() | ||
| 45 | + | ||
| 46 | + if err := repository.primary.Store(ctx, monitor); err != nil { | ||
| 47 | + return err | ||
| 48 | + } | ||
| 49 | + | ||
| 50 | + if err := repository.secondary.Store(ctx, monitor); err != nil { | ||
| 51 | + repository.logger.Error(stacktrace.Propagate(err, fmt.Sprintf("hedging: secondary write failed for monitor [%s]", monitor.ID))) | ||
| 52 | + repository.failureCounter.Add(ctx, 1) | ||
| 53 | + } | ||
| 54 | + | ||
| 55 | + return nil | ||
| 56 | + } | ||
| 57 | + | ||
| 58 | + func (repository *hedgingHeartbeatMonitorRepository) Load(ctx context.Context, userID entities.UserID, phoneNumber string) (*entities.HeartbeatMonitor, error) { | ||
| 59 | + return repository.primary.Load(ctx, userID, phoneNumber) | ||
| 60 | + } | ||
| 61 | + | ||
| 62 | + func (repository *hedgingHeartbeatMonitorRepository) Exists(ctx context.Context, userID entities.UserID, monitorID uuid.UUID) (bool, error) { | ||
| 63 | + return repository.primary.Exists(ctx, userID, monitorID) | ||
| 64 | + } | ||
| 65 | + | ||
| 66 | + func (repository *hedgingHeartbeatMonitorRepository) UpdateQueueID(ctx context.Context, monitorID uuid.UUID, queueID string) error { | ||
| 67 | + ctx, span := repository.tracer.Start(ctx) | ||
| 68 | + defer span.End() | ||
| 69 | + | ||
| 70 | + if err := repository.primary.UpdateQueueID(ctx, monitorID, queueID); err != nil { | ||
| 71 | + return err | ||
| 72 | + } | ||
| 73 | + | ||
| 74 | + if err := repository.secondary.UpdateQueueID(ctx, monitorID, queueID); err != nil { | ||
| 75 | + repository.logger.Error(stacktrace.Propagate(err, fmt.Sprintf("hedging: secondary UpdateQueueID failed for monitor [%s]", monitorID))) | ||
| 76 | + repository.failureCounter.Add(ctx, 1) | ||
| 77 | + } | ||
| 78 | + | ||
| 79 | + return nil | ||
| 80 | + } | ||
| 81 | + | ||
| 82 | + func (repository *hedgingHeartbeatMonitorRepository) Delete(ctx context.Context, userID entities.UserID, phoneNumber string) error { | ||
| 83 | + ctx, span := repository.tracer.Start(ctx) | ||
| 84 | + defer span.End() | ||
| 85 | + | ||
| 86 | + if err := repository.primary.Delete(ctx, userID, phoneNumber); err != nil { | ||
| 87 | + return err | ||
| 88 | + } | ||
| 89 | + | ||
| 90 | + if err := repository.secondary.Delete(ctx, userID, phoneNumber); err != nil { | ||
| 91 | + repository.logger.Error(stacktrace.Propagate(err, fmt.Sprintf("hedging: secondary delete failed for monitor with owner [%s]", phoneNumber))) | ||
| 92 | + repository.failureCounter.Add(ctx, 1) | ||
| 93 | + } | ||
| 94 | + | ||
| 95 | + return nil | ||
| 96 | + } | ||
| 97 | + | ||
| 98 | + func (repository *hedgingHeartbeatMonitorRepository) UpdatePhoneOnline(ctx context.Context, userID entities.UserID, monitorID uuid.UUID, online bool) error { | ||
| 99 | + ctx, span := repository.tracer.Start(ctx) | ||
| 100 | + defer span.End() | ||
| 101 | + | ||
| 102 | + if err := repository.primary.UpdatePhoneOnline(ctx, userID, monitorID, online); err != nil { | ||
| 103 | + return err | ||
| 104 | + } | ||
| 105 | + | ||
| 106 | + if err := repository.secondary.UpdatePhoneOnline(ctx, userID, monitorID, online); err != nil { | ||
| 107 | + repository.logger.Error(stacktrace.Propagate(err, fmt.Sprintf("hedging: secondary UpdatePhoneOnline failed for monitor [%s]", monitorID))) | ||
| 108 | + repository.failureCounter.Add(ctx, 1) | ||
| 109 | + } | ||
| 110 | + | ||
| 111 | + return nil | ||
| 112 | + } | ||
| 113 | + | ||
| 114 | + func (repository *hedgingHeartbeatMonitorRepository) DeleteAllForUser(ctx context.Context, userID entities.UserID) error { | ||
| 115 | + ctx, span := repository.tracer.Start(ctx) | ||
| 116 | + defer span.End() | ||
| 117 | + | ||
| 118 | + if err := repository.primary.DeleteAllForUser(ctx, userID); err != nil { | ||
| 119 | + return err | ||
| 120 | + } | ||
| 121 | + | ||
| 122 | + if err := repository.secondary.DeleteAllForUser(ctx, userID); err != nil { | ||
| 123 | + repository.logger.Error(stacktrace.Propagate(err, fmt.Sprintf("hedging: secondary delete all failed for user [%s]", userID))) | ||
| 124 | + repository.failureCounter.Add(ctx, 1) | ||
| 125 | + } | ||
| 126 | + | ||
| 127 | + return nil | ||
| 128 | + } | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments