-
Notifications
You must be signed in to change notification settings - Fork 7
Expand file tree
/
Copy pathprojection_runner_option.go
More file actions
130 lines (111 loc) · 4.31 KB
/
Copy pathprojection_runner_option.go
File metadata and controls
130 lines (111 loc) · 4.31 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
// MIT License
//
// Copyright (c) 2022-2026 Arsene Tochemey Gandote
//
// Permission is hereby granted, free of charge, to any person obtaining a copy
// of this software and associated documentation files (the "Software"), to deal
// in the Software without restriction, including without limitation the rights
// to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
// copies of the Software, and to permit persons to whom the Software is
// furnished to do so, subject to the following conditions:
//
// The above copyright notice and this permission notice shall be included in all
// copies or substantial portions of the Software.
//
// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
// SOFTWARE.
package ego
import (
"time"
"github.com/tochemey/goakt/v4/log"
"github.com/tochemey/ego/v4/encryption"
"github.com/tochemey/ego/v4/eventadapter"
"github.com/tochemey/ego/v4/eventstream"
"github.com/tochemey/ego/v4/projection"
)
// Option is the interface that applies a configuration option.
type runnerOption interface {
// Apply sets the Option value of a config.
Apply(runner *projectionRunner)
}
var _ runnerOption = runnerOptionFunc(nil)
// OptionFunc implements the Option interface.
type runnerOptionFunc func(*projectionRunner)
// Apply applies the options to Engine
func (f runnerOptionFunc) Apply(runner *projectionRunner) {
f(runner)
}
// WithPullInterval sets the events pull interval
// This defines how often the projection will fetch events
func withPullInterval(interval time.Duration) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.pullInterval = interval
})
}
// withMaxBufferSize sets the max buffer size.
// This defines how many events are fetched on a single run of the projection
func withMaxBufferSize(bufferSize int) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.maxBufferSize = bufferSize
})
}
// withStartOffset sets the starting point where to read the events
func withStartOffset(startOffset time.Time) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.startingOffset = startOffset
})
}
// withResetOffset helps reset the offset to a given timestamp.
func withResetOffset(resetOffset time.Time) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.resetOffsetTo = resetOffset
})
}
// WithLogger sets the actor system custom log
func withLogger(logger log.Logger) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.logger = logger
})
}
// withRecoveryStrategy sets the recovery strategy
func withRecoveryStrategy(strategy *projection.Recovery) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.recovery = strategy
})
}
// withDeadLetterHandler sets the dead letter handler for the projection runner
func withDeadLetterHandler(handler projection.DeadLetterHandler) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.deadLetterHandler = handler
})
}
// withEventAdapters sets the event adapters for the projection runner
func withEventAdapters(adapters []eventadapter.EventAdapter) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.eventAdapters = adapters
})
}
// withMetrics sets the metrics for the projection runner
func withMetrics(m *metrics) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.metrics = m
})
}
// withEncryptor sets the encryptor for the projection runner
func withEncryptor(enc encryption.Encryptor) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.encryptor = enc
})
}
// withEventsStream sets the in-process events stream that triggers an
// immediate pull when events are persisted on the local node.
func withEventsStream(stream eventstream.Stream) runnerOption {
return runnerOptionFunc(func(runner *projectionRunner) {
runner.eventsStream = stream
})
}