@@ -2,47 +2,18 @@ package postgres
22
33import (
44 "context"
5- "errors"
65 "fmt"
7- "time"
86
97 "github.com/awslabs/aurora-dsql-connectors/go/pgx/dsql"
10- "github.com/cenkalti/backoff/v4"
11- "github.com/jackc/pgx/v5/pgconn"
8+ "github.com/awslabs/aurora-dsql-connectors/go/pgx/occretry"
129 "github.com/jackc/pgx/v5/pgxpool"
1310
1411 "github.com/openfga/openfga/pkg/storage/sqlcommon"
1512)
1613
17- // isOCCError checks if the error is a DSQL optimistic concurrency control conflict.
18- // DSQL returns OC000 for mutation conflicts and OC001 for schema conflicts.
19- func isOCCError (err error ) bool {
20- if err == nil {
21- return false
22- }
23- var pgErr * pgconn.PgError
24- if errors .As (err , & pgErr ) {
25- return pgErr .Code == "OC000" || pgErr .Code == "OC001" || pgErr .Code == "40001"
26- }
27- return false
28- }
29-
3014// withOCCRetry executes fn with automatic retry on DSQL OCC errors.
3115func withOCCRetry (ctx context.Context , fn func () error ) error {
32- policy := backoff .NewExponentialBackOff ()
33- policy .InitialInterval = 10 * time .Millisecond
34- policy .MaxElapsedTime = 5 * time .Second
35-
36- return backoff .Retry (func () error {
37- err := fn ()
38- if err == nil {
39- return nil
40- }
41- if isOCCError (err ) {
42- return err
43- }
44- return backoff .Permanent (err )
45- }, backoff .WithContext (policy , ctx ))
16+ return occretry .Retry (ctx , occretry .DefaultConfig (), fn )
4617}
4718
4819// initDSQLDB initializes a new Aurora DSQL database connection.
@@ -73,10 +44,5 @@ func initDSQLDB(uri string, cfg *sqlcommon.Config) (*pgxpool.Pool, error) {
7344 poolCfg .MaxConnIdleTime = cfg .ConnMaxIdleTime
7445 }
7546
76- pool , err := dsql .NewPool (context .Background (), dsqlCfg , poolCfg )
77- if err != nil {
78- return nil , fmt .Errorf ("create DSQL pool: %w" , err )
79- }
80-
81- return pool .Pool , nil
47+ return dsql .NewPool (context .Background (), dsqlCfg , poolCfg )
8248}
0 commit comments