-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmanager.go
More file actions
187 lines (155 loc) · 4.76 KB
/
Copy pathmanager.go
File metadata and controls
187 lines (155 loc) · 4.76 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
187
package manager
import (
"context"
"encoding/json"
"fmt"
"strings"
// Packages
otel "github.com/mutablelogic/go-client/pkg/otel"
pg "github.com/mutablelogic/go-pg"
schema "github.com/mutablelogic/go-pg/pgqueue/schema"
types "github.com/mutablelogic/go-server/pkg/types"
attribute "go.opentelemetry.io/otel/attribute"
)
///////////////////////////////////////////////////////////////////////////////
// TYPES
type Manager struct {
opt
pg.PoolConn
queues *exec
tickers *exec
}
///////////////////////////////////////////////////////////////////////////////
// LIFECYCLE
func New(ctx context.Context, pool pg.PoolConn, opts ...Opt) (*Manager, error) {
self := new(Manager)
// Check arguments
if pool == nil {
return nil, fmt.Errorf("pool is required")
}
// Set default values
if err := self.defaults(); err != nil {
return nil, err
}
// Apply options
if err := self.apply(opts...); err != nil {
return nil, err
}
// Create execution objects after options are applied so they inherit the configured tracer.
self.queues = NewExec(self.tracer)
self.tickers = NewExec(self.tracer)
// Parse and register named queries so bind.Query(...) can resolve them.
queries, err := pg.NewQueries(strings.NewReader(schema.Queries))
if err != nil {
return nil, fmt.Errorf("parse queries.sql: %w", err)
} else {
pool = pool.WithQueries(queries).With(
"schema", self.schema,
"channel", schema.DefaultNotifyChannel,
).(pg.PoolConn)
}
// Create objects in the database schema. This is not done in a transaction
bootstrapCtx, endBootstrapSpan := otel.StartSpan(self.tracer, ctx, "bootstrap",
attribute.String("schema", self.schema),
)
if err := bootstrap(bootstrapCtx, pool, self.schema); err != nil {
endBootstrapSpan(err)
return nil, err
} else {
self.PoolConn = pool
}
// Ensure at least one task partition exists before any task inserts happen.
if _, err := self.CreateNextPartition(bootstrapCtx); err != nil {
endBootstrapSpan(err)
return nil, err
}
// Register a maintenance ticker
if _, err := self.RegisterTicker(bootstrapCtx, schema.DefaultMaintenanceTickerName, schema.TickerMeta{
Interval: types.Ptr(self.maintenancePeriod),
}, self.maintenance); err != nil {
endBootstrapSpan(err)
return nil, err
}
// Register a cleanup ticker
if _, err := self.RegisterTicker(bootstrapCtx, schema.DefaultCleanupTickerName, schema.TickerMeta{
Interval: types.Ptr(schema.DefaultCleanupPeriod),
}, self.cleanup); err != nil {
endBootstrapSpan(err)
return nil, err
}
// Register the metrics
if err := self.registerMetrics(); err != nil {
endBootstrapSpan(err)
return nil, err
}
// Return success
endBootstrapSpan(nil)
return self, nil
}
///////////////////////////////////////////////////////////////////////////////
// PRIVATE METHODS
func bootstrap(ctx context.Context, conn pg.Conn, schemaName string) error {
// Get all objects
objects, err := pg.NewQueries(strings.NewReader(schema.Objects))
if err != nil {
return fmt.Errorf("parse objects.sql: %w", err)
}
// Create the schema
if err := pg.SchemaCreate(ctx, conn, schemaName); err != nil {
return fmt.Errorf("create schema %q: %w", schemaName, err)
}
// Create all objects - not in a transaction
for _, key := range objects.Keys() {
if err := conn.Exec(ctx, objects.Query(key)); err != nil {
return fmt.Errorf("create object %q: %w", key, err)
}
}
// Return success
return nil
}
func (manager *Manager) maintenance(ctx context.Context, _ json.RawMessage) (any, error) {
messages := []string{}
// Create next partition if needed
created, err := manager.CreateNextPartition(ctx)
if err != nil {
return nil, err
} else if created != "" {
messages = append(messages, fmt.Sprintf("created partition %q", created))
}
// Drop old drained partitions
dropped, err := manager.DropDrainedPartition(ctx)
if err != nil {
return nil, err
} else if dropped != "" {
messages = append(messages, fmt.Sprintf("dropped partition %q", dropped))
}
// Return success
return messages, nil
}
func (manager *Manager) cleanup(ctx context.Context, _ json.RawMessage) (any, error) {
messages := make([]string, 0)
pageSize := uint64(schema.QueueListLimit)
offset := uint64(0)
for {
pageLimit := pageSize
request := schema.QueueListRequest{OffsetLimit: pg.OffsetLimit{Offset: offset, Limit: &pageLimit}}
queues, err := manager.ListQueues(ctx, request)
if err != nil {
return nil, err
}
for _, queue := range queues.Body {
removed, err := manager.CleanQueue(ctx, queue.Queue)
if err != nil {
return nil, err
}
if len(removed) > 0 {
messages = append(messages, fmt.Sprintf("cleaned queue %q removed %d tasks", queue.Queue, len(removed)))
}
}
if len(queues.Body) < int(pageSize) {
break
}
offset += uint64(len(queues.Body))
}
return messages, nil
}