persistence_rate_limited_clients.go ×11

Frontier kind: Code frontier

unlabeled · c_3dfa8b183cff

48 tests · 2404 LOC · 109 files · introduces 0 tests · 87 LOC · 2 files

Introduces — evidence that enters the hierarchy at this concept

Code
18 ranges87 lines · 2 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
379 ranges2404 lines · 109 files · Browse complete extent
All tests (intent)
48 testsBrowse complete intent

Neighbourhood graph

The orange circle is the focus. Violet and green circles are every ancestor and descendant, broader and narrower, at any distance; blue squares and pink diamonds are the introduced files and exact introduced tests of every visible concept, not only the focus's. Arrows point from broader to narrower concepts and bridge only concepts omitted from this view. Undirected links show source or test introduction. Concept and file size follows LOC; exact test nodes use test-count units.

Introduced files, introduced tests, and structurally relevant concept specialization

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the native relationship evidence on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the native relationship evidence on this page while the interactive map is unavailable.

Native relationship evidence

Every exact file and test below is linked only from the concept that introduces it.

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

No tests are introduced at this concept. Its intent tests are introduced by other concepts.

Introduced code

Every collected source range enters the hierarchy at exactly one concept.

2 files ranked by introduced lines: 87 introduced LOC across 18 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/persistence/persistence_rate_limited_clients.go 73 introduced LOC · 11 ranges

Open complete file

112 shardRateLimiter quotas.RequestRateLimiter,
113 logger log.Logger,
114 > ) ShardManager { persistence_rate_limited_clients.go
115 > return &shardRateLimitedPersistenceClient{
116 > persistence: persistence,
117 > systemRateLimiter: rateLimiter,
118 > namespaceRateLimiter: namespaceRateLimiter,
119 > shardRateLimiter: shardRateLimiter,
120 > logger: logger,
121 > }
122 > }
123
124 // NewExecutionPersistenceRateLimitedClient creates a client to manage executions
129 shardRateLimiter quotas.RequestRateLimiter,
130 logger log.Logger,
131 > ) ExecutionManager { persistence_rate_limited_clients.go
132 > return &executionRateLimitedPersistenceClient{
133 > persistence: persistence,
134 > systemRateLimiter: systemRateLimiter,
135 > namespaceRateLimiter: namespaceRateLimiter,
136 > shardRateLimiter: shardRateLimiter,
137 > logger: logger,
138 > }
139 > }
140
141 // NewTaskPersistenceRateLimitedClient creates a client to manage tasks
163 shardRateLimiter quotas.RequestRateLimiter,
164 logger log.Logger,
165 > ) MetadataManager { persistence_rate_limited_clients.go
166 > return &metadataRateLimitedPersistenceClient{
167 > persistence: persistence,
168 > systemRateLimiter: systemRateLimiter,
169 > namespaceRateLimiter: namespaceRateLimiter,
170 > shardRateLimiter: shardRateLimiter,
171 > logger: logger,
172 > }
173 > }
174
175 // NewClusterMetadataPersistenceRateLimitedClient creates a ClusterMetadataManager client to manage cluster metadata
180 shardRateLimiter quotas.RequestRateLimiter,
181 logger log.Logger,
182 > ) ClusterMetadataManager { persistence_rate_limited_clients.go
183 > return &clusterMetadataRateLimitedPersistenceClient{
184 > persistence: persistence,
185 > systemRateLimiter: systemRateLimiter,
186 > namespaceRateLimiter: namespaceRateLimiter,
187 > shardRateLimiter: shardRateLimiter,
188 > logger: logger,
189 > }
190 > }
191
192 // NewQueuePersistenceRateLimitedClient creates a client to manage queue
197 shardRateLimiter quotas.RequestRateLimiter,
198 logger log.Logger,
200 > return &queueRateLimitedPersistenceClient{
201 > persistence: persistence,
202 > systemRateLimiter: systemRateLimiter,
203 > namespaceRateLimiter: namespaceRateLimiter,
204 > shardRateLimiter: shardRateLimiter,
205 > logger: logger,
206 > }
207 > }
208
209 // NewNexusEndpointPersistenceRateLimitedClient creates a NexusEndpointManager to manage nexus endpoints
214 shardRateLimiter quotas.RequestRateLimiter,
215 logger log.Logger,
216 > ) NexusEndpointManager { persistence_rate_limited_clients.go
217 > return &nexusEndpointRateLimitedPersistenceClient{
218 > persistence: persistence,
219 > systemRateLimiter: systemRateLimiter,
220 > namespaceRateLimiter: namespaceRateLimiter,
221 > shardRateLimiter: shardRateLimiter,
222 > logger: logger,
223 > }
224 > }
225
226 func (p *shardRateLimitedPersistenceClient) GetName() string {
1005 ctx context.Context,
1006 blob *commonpb.DataBlob,
1008 > return p.persistence.Init(ctx, blob)
1009 > }
1010
1011 func (c *clusterMetadataRateLimitedPersistenceClient) Close() {
1151 namespaceRateLimiter quotas.RequestRateLimiter,
1152 shardRateLimiter quotas.RequestRateLimiter,
1154 > callerInfo := headers.GetCallerInfo(ctx)
1155 > // namespace-level rate limits has to be applied before system-level rate limits.
1156 > now := time.Now().UTC()
1157 > quotaRequest := quotas.NewRequest(
1158 > api,
1159 > RateLimitDefaultToken,
1160 > callerInfo.CallerName,
1161 > callerInfo.CallerType,
1162 > shardID,
1163 > callerInfo.CallOrigin,
1164 > )
1165 > if ok := shardRateLimiter.Allow(now, quotaRequest); !ok {
1166 return ErrPersistenceNamespaceShardLimitExceeded
1167 }
1168 > if ok := namespaceRateLimiter.Allow(now, quotaRequest); !ok { persistence_rate_limited_clients.go
1169 return ErrPersistenceNamespaceLimitExceeded
1170 }
1171 > if ok := systemRateLimiter.Allow(now, quotaRequest); !ok { persistence_rate_limited_clients.go
1172 return ErrPersistenceSystemLimitExceeded
1173 }
1175 }
1176
go.temporal.io/server/common/persistence/client/factory.go 14 introduced LOC · 7 ranges

Open complete file

149 result := persistence.NewShardManager(shardStore, f.serializer)
150 if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
151 > result = persistence.NewShardPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger) factory.go
152 > }
153 if f.metricsHandler != nil && f.healthSignals != nil {
154 result = persistence.NewShardPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
167 result := persistence.NewMetadataManagerImpl(store, f.serializer, f.logger, f.clusterName)
168 if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
169 > result = persistence.NewMetadataPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger) factory.go
170 > }
171 if f.metricsHandler != nil && f.healthSignals != nil {
172 result = persistence.NewMetadataPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
185 result := persistence.NewClusterMetadataManagerImpl(store, f.serializer, f.clusterName, f.logger)
186 if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
187 > result = persistence.NewClusterMetadataPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger) factory.go
188 > }
189 if f.metricsHandler != nil && f.healthSignals != nil {
190 result = persistence.NewClusterMetadataPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
210 )
211 if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
212 > result = persistence.NewExecutionPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger) factory.go
213 > }
214 if f.metricsHandler != nil && f.healthSignals != nil {
215 result = persistence.NewExecutionPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
226
227 if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
228 > result = persistence.NewQueuePersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger) factory.go
229 > }
230 if f.metricsHandler != nil && f.healthSignals != nil {
231 result = persistence.NewQueuePersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
251 result := persistence.NewNexusEndpointManager(store, f.serializer, f.logger)
252 if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
253 > result = persistence.NewNexusEndpointPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger) factory.go
254 > }
255 if f.metricsHandler != nil && f.healthSignals != nil {
256 result = persistence.NewNexusEndpointPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
296 }
297 if f.healthSignals == nil {
298 > f.healthSignals = persistence.NoopHealthSignalAggregator factory.go
299 > }
300 f.healthSignals.Start()
301 }