112
components []workercommon.PerNSWorkerComponent,
113
taskQueueName string,
115
>
return &PerNamespaceWorkerManager{
116
>
logger: log.With(logger, tag.ComponentPerNSWorkerManager),
117
>
sdkClientFactory: sdkClientFactory,
118
>
namespaceRegistry: namespaceRegistry,
119
>
hostName: hostName,
120
>
taskQueueName: taskQueueName,
121
>
config: config,
122
>
components: components,
123
>
initialRetry: 1 * time.Second,
124
>
thisClusterName: clusterMetadata.GetCurrentClusterName(),
125
>
startLimiter: quotas.NewDefaultOutgoingRateLimiter(quotas.RateFn(config.PerNamespaceWorkerStartRate)),
126
>
membershipChangedCh: make(chan *membership.ChangedEvent),
127
>
workers: make(map[namespace.ID]*perNamespaceWorker),
128
>
}
129
>
}
130
132
>
return atomic.LoadInt32(&wm.status) == common.DaemonStatusStarted
133
>
}
134
135
func (wm *PerNamespaceWorkerManager) Start(
136
self membership.HostInfo,
137
serviceResolver membership.ServiceResolver,
139
>
if !atomic.CompareAndSwapInt32(
140
>
&wm.status,
141
>
common.DaemonStatusInitialized,
142
>
common.DaemonStatusStarted,
143
>
) {
144
return
145
}
146
148
>
wm.serviceResolver = serviceResolver
149
>
150
>
wm.logger.Info("", tag.LifeCycleStarting)
151
>
152
>
// this will call namespaceCallback with current namespaces
153
>
wm.namespaceRegistry.RegisterStateChangeCallback(wm, wm.namespaceCallback)
154
>
155
>
err := wm.serviceResolver.AddListener(fmt.Sprintf("%p", wm), wm.membershipChangedCh)
156
>
if err != nil {
157
wm.logger.Fatal("Unable to register membership listener", tag.Error(err))
158
}
160
>
wm.backgroundLoops.Go(wm.periodicRefreshLoop)
161
>
162
>
wm.logger.Info("", tag.LifeCycleStarted)
163
}
164
166
>
if !atomic.CompareAndSwapInt32(
167
>
&wm.status,
168
>
common.DaemonStatusStarted,
169
>
common.DaemonStatusStopped,
170
>
) {
171
return
172
}
173
175
>
176
>
wm.namespaceRegistry.UnregisterStateChangeCallback(wm)
177
>
err := wm.serviceResolver.RemoveListener(fmt.Sprintf("%p", wm))
178
>
if err != nil {
179
wm.logger.Error("Unable to unregister membership listener", tag.Error(err))
180
}
182
>
wm.backgroundLoops.Wait()
183
>
184
>
wm.lock.Lock()
185
>
workers := expmaps.Values(wm.workers)
186
>
maps.DeleteFunc(wm.workers, func(_ namespace.ID, _ *perNamespaceWorker) bool { return true })
187
>
wm.lock.Unlock()
188
>
189
>
for _, worker := range workers {
191
>
worker.cancel()
192
>
}
193
195
}
196
197
>
func (wm *PerNamespaceWorkerManager) namespaceCallback(ns *namespace.Namespace, nsDeleted bool) {
pernamespaceworker.go
198
>
go wm.getWorkerByNamespace(ns).update(ns, nsDeleted, nil, nil)
199
>
}
200
202
>
wm.lock.Lock()
203
>
defer wm.lock.Unlock()
204
>
for _, worker := range wm.workers {
205
>
go worker.update(nil, false, nil, nil)
206
>
}
207
}
208
209
>
func (wm *PerNamespaceWorkerManager) membershipChangedListener(ctx context.Context) error {
pernamespaceworker.go
210
>
for {
211
>
select {
213
>
return nil
215
>
wm.refreshAll()
216
}
217
}
218
}
219
220
>
func (wm *PerNamespaceWorkerManager) periodicRefreshLoop(ctx context.Context) error {
pernamespaceworker.go
221
>
ticker := time.NewTicker(refreshInterval)
222
>
defer ticker.Stop()
223
>
224
>
for {
225
>
select {
227
>
return nil
228
case <-ticker.C:
229
wm.refreshAll()