-
Notifications
You must be signed in to change notification settings - Fork 22
/
serial_events.go
90 lines (77 loc) · 2.16 KB
/
serial_events.go
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
// Copyright 2022 The Cockroach Authors.
//
// Use of this software is governed by the Business Source License
// included in the file licenses/BSL.txt.
//
// As of the Change Date specified in that file, in accordance with
// the Business Source License, use of this software will be governed
// by the Apache License, Version 2.0, included in the file
// licenses/APL.txt.
package logical
import (
"context"
"github.com/cockroachdb/cdc-sink/internal/types"
"github.com/cockroachdb/cdc-sink/internal/util/ident"
"github.com/cockroachdb/cdc-sink/internal/util/stamp"
"github.com/jackc/pgx/v4"
"github.com/jackc/pgx/v4/pgxpool"
"github.com/pkg/errors"
)
// serialEvents is a transaction-preserving implementation of Events.
type serialEvents struct {
State
appliers types.Appliers
pool *pgxpool.Pool
stamp stamp.Stamp // the latest value passed to OnCommit.
tx pgx.Tx // db transaction created by OnCommit.
}
var _ Events = (*serialEvents)(nil)
// OnBegin implements Events.
func (e *serialEvents) OnBegin(ctx context.Context, point stamp.Stamp) error {
var err error
if e.tx != nil {
return errors.Errorf("OnBegin already called at %s", e.stamp)
}
e.stamp = point
e.tx, err = e.pool.Begin(ctx)
return errors.WithStack(err)
}
// OnCommit implements Events.
func (e *serialEvents) OnCommit(ctx context.Context) error {
if e.tx == nil {
return errors.New("OnCommit called without matching OnBegin")
}
err := e.tx.Commit(ctx)
e.tx = nil
if err != nil {
return errors.WithStack(err)
}
e.setConsistentPoint(e.stamp)
return nil
}
// OnData implements Events.
func (e *serialEvents) OnData(
ctx context.Context, target ident.Table, muts []types.Mutation,
) error {
app, err := e.appliers.Get(ctx, target)
if err != nil {
return err
}
return app.Apply(ctx, e.tx, muts)
}
// OnRollback implements Events.
func (e *serialEvents) OnRollback(_ context.Context, msg Message) error {
if !IsRollback(msg) {
return errors.New("the rollback message must be passed to OnRollback")
}
e.stop()
return nil
}
// reset implements Events.
func (e *serialEvents) stop() {
if e.tx != nil {
_ = e.tx.Rollback(context.Background())
}
e.stamp = nil
e.tx = nil
}