-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpostgresdb.go
More file actions
186 lines (170 loc) · 5.21 KB
/
Copy pathpostgresdb.go
File metadata and controls
186 lines (170 loc) · 5.21 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
package postgresdb
import (
"context"
"errors"
"fmt"
"strings"
"sync"
"time"
"github.com/devctllabs/go-libs/txmanager"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// Config defines one PostgreSQL database instance.
type Config struct {
Writer EndpointConfig
Reader *EndpointConfig
Telemetry Telemetry
}
// EndpointConfig configures one PostgreSQL pool.
type EndpointConfig struct {
DSN string
Pool PoolConfig
}
// PoolConfig overrides a stable subset of pgx pool settings when fields are positive.
type PoolConfig struct {
MaxConnections int32
MinIdleConnections int32
MaxConnectionLifetime time.Duration
MaxConnectionLifetimeJitter time.Duration
MaxConnectionIdleTime time.Duration
HealthCheckPeriod time.Duration
}
// DB owns PostgreSQL reader and writer resources.
type DB struct {
reader *Endpoint
writer *Endpoint
managers txmanager.Managers
readerPool *pgxpool.Pool
writerPool *pgxpool.Pool
closeOnce sync.Once
}
// Open constructs a PostgreSQL database instance without requiring network availability.
func Open(ctx context.Context, cfg Config) (*DB, error) {
if ctx == nil {
return nil, errors.New("postgresdb: context must not be nil")
}
telemetry := newTelemetryConfig(cfg.Telemetry)
writerPool, err := openPool(ctx, cfg.Writer, telemetry, "writer")
if err != nil {
return nil, err
}
readerPool := writerPool
if cfg.Reader != nil {
readerPool, err = openPool(ctx, *cfg.Reader, telemetry, "reader")
if err != nil {
writerPool.Close()
return nil, err
}
}
return buildDB(readerPool, writerPool)
}
func openPool(ctx context.Context, endpoint EndpointConfig, telemetry telemetryConfig, role string) (*pgxpool.Pool, error) {
if strings.TrimSpace(endpoint.DSN) == "" {
return nil, fmt.Errorf("postgresdb: %s DSN must not be blank", role)
}
config, err := parsePoolConfig(endpoint, telemetry, role)
if err != nil {
return nil, fmt.Errorf("postgresdb: %s config: %w", role, err)
}
pool, err := pgxpool.NewWithConfig(ctx, config)
if err != nil {
return nil, fmt.Errorf("postgresdb: open %s pool: %w", role, err)
}
if err := telemetry.recordPoolStats(pool, role); err != nil {
pool.Close()
return nil, fmt.Errorf("postgresdb: instrument %s pool: %w", role, err)
}
return pool, nil
}
func buildDB(readerPool, writerPool *pgxpool.Pool) (*DB, error) {
backend := &transactionBackend{reader: readerPool, writer: writerPool}
coordinator, err := txmanager.NewCoordinator[pgx.Tx](backend)
if err != nil {
closePools(readerPool, writerPool)
return nil, err
}
reader := &Endpoint{pool: readerPool, coordinator: coordinator, manager: coordinator.Reader()}
writer := &Endpoint{pool: writerPool, coordinator: coordinator, manager: coordinator.Writer()}
managers, err := txmanager.NewManagers(reader, writer)
if err != nil {
closePools(readerPool, writerPool)
return nil, err
}
return &DB{
reader: reader,
writer: writer,
managers: managers,
readerPool: readerPool,
writerPool: writerPool,
}, nil
}
// Reader returns the reader endpoint.
func (db *DB) Reader() *Endpoint {
return db.reader
}
// Writer returns the writer endpoint.
func (db *DB) Writer() *Endpoint {
return db.writer
}
// TxManagers returns reader and writer transaction managers backed by this DB.
func (db *DB) TxManagers() txmanager.Managers {
return db.managers
}
// Close releases all distinct pools once.
func (db *DB) Close() error {
db.closeOnce.Do(func() {
if db.readerPool != db.writerPool {
db.readerPool.Close()
}
db.writerPool.Close()
})
return nil
}
func closePools(reader *pgxpool.Pool, writer *pgxpool.Pool) {
if reader != writer {
reader.Close()
}
writer.Close()
}
func parsePoolConfig(endpoint EndpointConfig, telemetry telemetryConfig, role string) (*pgxpool.Config, error) {
if err := validatePoolConfig(endpoint.Pool); err != nil {
return nil, err
}
config, err := pgxpool.ParseConfig(endpoint.DSN)
if err != nil {
return nil, fmt.Errorf("parse DSN: %w", err)
}
config.ConnConfig.DefaultQueryExecMode = pgx.QueryExecModeExec
config.ConnConfig.Tracer = telemetry.tracer(role)
if endpoint.Pool.MaxConnections > 0 {
config.MaxConns = endpoint.Pool.MaxConnections
}
if endpoint.Pool.MinIdleConnections > 0 {
config.MinIdleConns = endpoint.Pool.MinIdleConnections
}
if endpoint.Pool.MaxConnectionLifetime > 0 {
config.MaxConnLifetime = endpoint.Pool.MaxConnectionLifetime
}
if endpoint.Pool.MaxConnectionLifetimeJitter > 0 {
config.MaxConnLifetimeJitter = endpoint.Pool.MaxConnectionLifetimeJitter
}
if endpoint.Pool.MaxConnectionIdleTime > 0 {
config.MaxConnIdleTime = endpoint.Pool.MaxConnectionIdleTime
}
if endpoint.Pool.HealthCheckPeriod > 0 {
config.HealthCheckPeriod = endpoint.Pool.HealthCheckPeriod
}
if config.MinIdleConns > config.MaxConns {
return nil, errors.New("minimum idle connections exceeds maximum connections")
}
return config, nil
}
func validatePoolConfig(config PoolConfig) error {
if config.MaxConnections < 0 || config.MinIdleConnections < 0 ||
config.MaxConnectionLifetime < 0 || config.MaxConnectionLifetimeJitter < 0 ||
config.MaxConnectionIdleTime < 0 || config.HealthCheckPeriod < 0 {
return errors.New("pool values must not be negative")
}
return nil
}