107
}
108
109
>
func newTestContext(t *resourcetest.Test, eventsCache events.Cache, config ContextConfigOverrides) *ContextImpl {
context_testutil.go
110
>
hostInfoProvider := t.GetHostInfoProvider()
111
>
lifecycleCtx, lifecycleCancel := context.WithCancel(context.Background())
112
>
if config.ShardInfo.QueueStates == nil {
113
config.ShardInfo.QueueStates = make(map[int32]*persistencespb.QueueState)
114
}
116
>
if registry == nil {
117
>
registry = t.GetNamespaceRegistry()
118
>
}
119
>
clusterMetadata := config.ClusterMetadata
120
>
if clusterMetadata == nil {
121
>
clusterMetadata = t.GetClusterMetadata()
122
>
}
123
>
executionManager := config.ExecutionManager
124
>
if executionManager == nil {
125
>
executionManager = t.ExecutionMgr
126
>
}
127
>
taskCategoryRegistry := tasks.NewDefaultTaskCategoryRegistry()
128
>
taskCategoryRegistry.AddCategory(tasks.CategoryArchival)
129
>
130
>
ctx := &ContextImpl{
131
>
shardID: config.ShardInfo.GetShardId(),
132
>
owner: config.ShardInfo.GetOwner(),
133
>
stringRepr: fmt.Sprintf("Shard(%d)", config.ShardInfo.GetShardId()),
134
>
executionManager: executionManager,
135
>
metricsHandler: t.MetricsHandler,
136
>
eventsCache: eventsCache,
137
>
config: config.Config,
138
>
contextTaggedLogger: t.GetLogger(),
139
>
throttledLogger: t.GetThrottledLogger(),
140
>
lifecycleCtx: lifecycleCtx,
141
>
lifecycleCancel: lifecycleCancel,
142
>
queueMetricEmitter: sync.Once{},
143
>
144
>
state: contextStateAcquired,
145
>
engineFuture: future.NewFuture[historyi.Engine](),
146
>
shardInfo: config.ShardInfo,
147
>
remoteClusterInfos: make(map[string]*remoteClusterInfo),
148
>
149
>
clusterMetadata: clusterMetadata,
150
>
timeSource: t.TimeSource,
151
>
namespaceRegistry: registry,
152
>
stateMachineRegistry: hsm.NewRegistry(),
153
>
chasmRegistry: chasm.NewRegistry(t.GetLogger()),
154
>
businessIDRateLimiters: cache.New(
155
>
config.Config.BusinessIDReuseLimiterCacheSize(),
156
>
&cache.Options{TTL: config.Config.BusinessIDReuseLimiterCacheTTL()},
157
>
),
158
>
persistenceShardManager: t.GetShardManager(),
159
>
clientBean: t.GetClientBean(),
160
>
saProvider: t.GetSearchAttributesProvider(),
161
>
saMapperProvider: t.GetSearchAttributesMapperProvider(),
162
>
historyClient: t.GetHistoryClient(),
163
>
payloadSerializer: t.GetPayloadSerializer(),
164
>
archivalMetadata: t.GetArchivalMetadata(),
165
>
hostInfoProvider: hostInfoProvider,
166
>
taskCategoryRegistry: taskCategoryRegistry,
167
>
ioSemaphore: locks.NewPrioritySemaphore(1),
168
>
}
169
>
ctx.taskKeyManager = newTaskKeyManager(
170
>
ctx.taskCategoryRegistry,
171
>
ctx.timeSource,
172
>
config.Config,
173
>
ctx.GetLogger(),
174
>
func() error {
175
return ctx.renewRangeLocked(false)
176
},
177
)
179
>
ctx.handoverTracker = NewDefaultHandoverTrackerFactory()(HandoverTrackerParams{
180
>
ClusterMetadata: clusterMetadata,
181
>
GetMaxReplicationTaskID: ctx.getMaxReplicationTaskID,
182
>
ErrorByStateFn: ctx.errorByState,
183
>
NotifyReplicationFn: ctx.notifyReplicationQueueProcessor,
184
>
NamespaceRegistry: registry,
185
>
Logger: ctx.contextTaggedLogger,
186
>
})
187
>
return ctx
188
}
189