Atlas › Test
TestThrottleLogger
Exact test identity: go.temporal.io/server/common/log/TestThrottleLogger
- Package
go.temporal.io/server/common/log
- Suite / test hierarchy
TestThrottleLogger
- Test
TestThrottleLogger
- Introduced at
- TestThrottleLogger Frontier kind: Test frontier
- Covered ranges
- 80
- Covered lines
- 422
- Covered files
- 20
Covered source
Expand a file to inspect source; the > gutter marks covered lines.
go.temporal.io/server/common/log/zap_logger.go 64 covered LOC · 17 ranges
Open complete file
82
83
// NewZapLogger returns a new zap based logger from zap.Logger
84
>
func NewZapLogger(zl *zap.Logger) *zapLogger {
zap_logger.go
85
>
return &zapLogger{
86
>
zl: zl,
87
>
skip: skipForZapLogger,
88
>
baseZl: zl,
89
>
}
90
>
}
91
92
// BuildZapLogger builds and returns a new zap.Logger for this logging configuration
95
}
96
98
>
_, path, line, ok := runtime.Caller(skip)
99
>
if !ok {
100
return ""
101
}
103
}
104
105
>
func (l *zapLogger) buildFieldsWithCallAt(tags []tag.Tag) []zap.Field {
zap_logger.go
106
>
fields := make([]zap.Field, len(tags)+1)
107
>
l.fillFields(tags, fields)
108
>
fields[len(fields)-1] = zap.String(tag.LoggingCallAtKey, caller(l.skip))
109
>
return fields
110
>
}
111
112
// fillFields fill fields parameter with fields read from tags. Optimized for performance.
113
>
func (l *zapLogger) fillFields(tags []tag.Tag, fields []zap.Field) {
zap_logger.go
114
>
for i, t := range tags {
116
>
fields[i] = zt.Field()
117
>
} else {
118
fields[i] = zap.Any(t.Key(), t.Value())
119
}
121
}
122
124
>
if msg == "" {
125
return defaultMsgForEmpty
126
}
128
}
129
136
}
137
138
>
func (l *zapLogger) Info(msg string, tags ...tag.Tag) {
zap_logger.go
139
>
if l.zl.Core().Enabled(zap.InfoLevel) {
140
>
msg = setDefaultMsg(msg)
141
>
fields := l.buildFieldsWithCallAt(tags)
142
>
l.zl.Info(msg, fields...)
143
>
}
144
}
145
191
//
192
// by deduping "foo" against any existing "foo" tags *only in the former*
193
>
func (l *zapLogger) With(tags ...tag.Tag) Logger {
zap_logger.go
194
>
cloneTags := mergeTags(l.tags, tags)
195
>
if l.baseZl == nil {
196
l.baseZl = l.zl
197
}
199
}
200
201
>
func (l *zapLogger) cloneWithTags(tags []tag.Tag) Logger {
zap_logger.go
202
>
fields := make([]zap.Field, len(tags))
203
>
l.fillFields(tags, fields)
204
>
zl := l.baseZl.With(fields...)
205
>
return &zapLogger{
206
>
zl: zl,
207
>
skip: l.skip,
208
>
baseZl: l.baseZl,
209
>
tags: tags,
210
>
}
211
>
}
212
213
>
func (l *zapLogger) Skip(extraSkip int) Logger {
zap_logger.go
214
>
return &zapLogger{
215
>
zl: l.zl,
216
>
skip: l.skip + extraSkip,
217
>
baseZl: l.baseZl,
218
>
}
219
>
}
220
221
>
func mergeTags(oldTags, newTags []tag.Tag) (outTags []tag.Tag) {
zap_logger.go
222
>
// Even if oldTags empty, we don't just return newTags because we need to de-dupe it.
223
>
outTags = slices.Clone(oldTags)
224
>
for _, t := range newTags {
225
>
if i := slices.IndexFunc(outTags, func(ti tag.Tag) bool {
227
>
}); i >= 0 {
228
outTags[i] = t
230
>
outTags = append(outTags, t)
231
>
}
232
}
234
}
235
go.temporal.io/server/common/log/tag/tags.go 33 covered LOC · 10 ranges
Open complete file
35
36
// Error returns tag for Error
37
>
func Error(err error) ZapTag {
tags.go
38
>
return ZapTag{
39
>
// NOTE: zap already chosen "error" as key
40
>
field: zap.Error(err),
41
>
}
42
>
}
43
44
// ServiceErrorType returns tag for ServiceErrorType
70
71
// WorkflowAction returns tag for WorkflowAction
72
>
func workflowAction(action string) ZapTag {
tags.go
73
>
return NewStringTag("wf-action", action)
74
>
}
75
76
// WorkflowListFilterType returns tag for WorkflowListFilterType
77
>
func workflowListFilterType(listFilterType string) ZapTag {
tags.go
78
>
return NewStringTag("wf-list-filter-type", listFilterType)
79
>
}
80
81
// general
376
377
// Component returns tag for Component
378
>
func component(component string) ZapTag {
tags.go
379
>
return NewStringTag("component", component)
380
>
}
381
382
// Lifecycle returns tag for Lifecycle
383
>
func lifecycle(lifecycle string) ZapTag {
tags.go
384
>
return NewStringTag("lifecycle", lifecycle)
385
>
}
386
387
// StoreOperation returns tag for StoreOperation
388
>
func storeOperation(storeOperation string) ZapTag {
tags.go
389
>
return NewStringTag("store-operation", storeOperation)
390
>
}
391
392
// OperationResult returns tag for OperationResult
393
>
func operationResult(operationResult string) ZapTag {
tags.go
394
>
return NewStringTag("operation-result", operationResult)
395
>
}
396
397
// ErrorType returns tag for ErrorType
401
402
// errorType returns tag for ErrorType given a string
403
>
func errorType(errorType string) ZapTag {
tags.go
404
>
return NewStringTag("error-type", errorType)
405
>
}
406
407
// Shardupdate returns tag for Shardupdate
408
>
func shardupdate(shardupdate string) ZapTag {
tags.go
409
>
return NewStringTag("shard-update", shardupdate)
410
>
}
411
412
// scope returns a tag for scope
413
// Pre-defined scope tags are in values.go.
414
>
func scope(scope string) ZapTag {
tags.go
415
>
return NewStringTag("scope", scope)
416
>
}
417
418
// general
go.temporal.io/server/common/quotas/rate_burst.go 26 covered LOC · 10 ranges
Open complete file
85
rateFn RateFn,
86
burstFn BurstFn,
88
>
return &RateBurstImpl{
89
>
rateFn: rateFn,
90
>
burstFn: burstFn,
91
>
}
92
>
}
93
94
func NewDefaultIncomingRateBurst(
102
func NewDefaultOutgoingRateBurst(
103
rateFn RateFn,
105
>
return NewDefaultRateBurst(rateFn, func() float64 {
106
>
return defaultOutgoingRateBurstRatio
107
>
})
108
}
109
111
rateFn RateFn,
112
rateToBurstRatio BurstRatioFn,
114
>
burstFn := func() int {
116
>
if rate < 0 {
117
rate = 0
118
}
119
121
>
if ratio < 0 {
122
ratio = 0
123
}
125
>
if burst == 0 && rate > 0 && ratio > 0 {
126
burst = 1
127
}
129
}
131
}
132
134
>
return d.rateFn()
135
>
}
136
138
>
return d.burstFn()
139
>
}
140
141
func NewMutableRateBurst(
go.temporal.io/server/common/log/throttle_logger.go 25 covered LOC · 6 ranges
Open complete file
21
//
22
// Fatal/Panic/DPanic logs are always emitted without any throttling
23
>
func NewThrottledLogger(logger Logger, rps quotas.RateFn) *throttledLogger {
throttle_logger.go
24
>
if sl, ok := logger.(SkipLogger); ok {
26
>
}
27
29
>
tl := &throttledLogger{
30
>
limiter: limiter,
31
>
logger: logger,
32
>
}
33
>
return tl
34
}
35
40
}
41
43
>
tl.rateLimit(func() {
44
>
tl.logger.Info(msg, tags...)
45
>
})
46
}
47
81
82
// Return a logger with the specified key-value pairs set, to be included in a subsequent normal logging call
84
>
result := &throttledLogger{
85
>
limiter: tl.limiter,
86
>
logger: With(tl.logger, tags...),
87
>
}
88
>
return result
89
>
}
90
92
>
if tl.limiter.Allow() {
93
>
f()
94
>
}
95
}
96
go.temporal.io/server/common/quotas/dynamic_rate_limiter_impl.go 23 covered LOC · 5 ranges
Open complete file
29
rateBurstFn RateBurst,
30
refreshInterval time.Duration,
32
>
rateLimiter := &DynamicRateLimiterImpl{
33
>
rateBurstFn: rateBurstFn,
34
>
refreshInterval: refreshInterval,
35
>
36
>
refreshTimer: time.NewTimer(refreshInterval),
37
>
rateLimiter: NewRateLimiter(rateBurstFn.Rate(), rateBurstFn.Burst()),
38
>
}
39
>
return rateLimiter
40
>
}
41
42
// NewDefaultIncomingRateLimiter returns a default rate limiter
57
func NewDefaultOutgoingRateLimiter(
58
rateFn RateFn,
60
>
return NewDynamicRateLimiter(
61
>
NewDefaultOutgoingRateBurst(rateFn),
62
>
defaultRefreshInterval,
63
>
)
64
>
}
65
66
// NewDefaultRateLimiter returns a default rate limiter with a dynamic burst ratio
78
// Allow immediately returns with true or false indicating if a rate limit
79
// token is available or not
81
>
d.maybeRefresh()
82
>
return d.rateLimiter.Allow()
83
>
}
84
85
// AllowN immediately returns with true or false indicating if n rate limit
130
}
131
133
>
select {
134
case <-d.refreshTimer.C:
135
d.refreshTimer.Reset(d.refreshInterval)
136
d.Refresh()
137
139
// noop
140
}
go.temporal.io/server/api/enums/v1/cluster.pb.go 20 covered LOC · 2 ranges
Open complete file
209
}
210
211
>
func init() { file_temporal_server_api_enums_v1_cluster_proto_init() }
cluster.pb.go
212
>
func file_temporal_server_api_enums_v1_cluster_proto_init() {
213
>
if File_temporal_server_api_enums_v1_cluster_proto != nil {
214
return
215
}
217
>
out := protoimpl.TypeBuilder{
218
>
File: protoimpl.DescBuilder{
219
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
220
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_cluster_proto_rawDesc), len(file_temporal_server_api_enums_v1_cluster_proto_rawDesc)),
221
>
NumEnums: 2,
222
>
NumMessages: 0,
223
>
NumExtensions: 0,
224
>
NumServices: 0,
225
>
},
226
>
GoTypes: file_temporal_server_api_enums_v1_cluster_proto_goTypes,
227
>
DependencyIndexes: file_temporal_server_api_enums_v1_cluster_proto_depIdxs,
228
>
EnumInfos: file_temporal_server_api_enums_v1_cluster_proto_enumTypes,
229
>
}.Build()
230
>
File_temporal_server_api_enums_v1_cluster_proto = out.File
231
>
file_temporal_server_api_enums_v1_cluster_proto_goTypes = nil
232
>
file_temporal_server_api_enums_v1_cluster_proto_depIdxs = nil
233
}
go.temporal.io/server/api/enums/v1/common.pb.go 20 covered LOC · 2 ranges
Open complete file
264
}
265
266
>
func init() { file_temporal_server_api_enums_v1_common_proto_init() }
common.pb.go
267
>
func file_temporal_server_api_enums_v1_common_proto_init() {
268
>
if File_temporal_server_api_enums_v1_common_proto != nil {
269
return
270
}
272
>
out := protoimpl.TypeBuilder{
273
>
File: protoimpl.DescBuilder{
274
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
275
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_common_proto_rawDesc), len(file_temporal_server_api_enums_v1_common_proto_rawDesc)),
276
>
NumEnums: 3,
277
>
NumMessages: 0,
278
>
NumExtensions: 0,
279
>
NumServices: 0,
280
>
},
281
>
GoTypes: file_temporal_server_api_enums_v1_common_proto_goTypes,
282
>
DependencyIndexes: file_temporal_server_api_enums_v1_common_proto_depIdxs,
283
>
EnumInfos: file_temporal_server_api_enums_v1_common_proto_enumTypes,
284
>
}.Build()
285
>
File_temporal_server_api_enums_v1_common_proto = out.File
286
>
file_temporal_server_api_enums_v1_common_proto_goTypes = nil
287
>
file_temporal_server_api_enums_v1_common_proto_depIdxs = nil
288
}
go.temporal.io/server/api/enums/v1/dlq.pb.go 20 covered LOC · 2 ranges
Open complete file
187
}
188
189
>
func init() { file_temporal_server_api_enums_v1_dlq_proto_init() }
dlq.pb.go
190
>
func file_temporal_server_api_enums_v1_dlq_proto_init() {
191
>
if File_temporal_server_api_enums_v1_dlq_proto != nil {
192
return
193
}
195
>
out := protoimpl.TypeBuilder{
196
>
File: protoimpl.DescBuilder{
197
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
198
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_dlq_proto_rawDesc), len(file_temporal_server_api_enums_v1_dlq_proto_rawDesc)),
199
>
NumEnums: 2,
200
>
NumMessages: 0,
201
>
NumExtensions: 0,
202
>
NumServices: 0,
203
>
},
204
>
GoTypes: file_temporal_server_api_enums_v1_dlq_proto_goTypes,
205
>
DependencyIndexes: file_temporal_server_api_enums_v1_dlq_proto_depIdxs,
206
>
EnumInfos: file_temporal_server_api_enums_v1_dlq_proto_enumTypes,
207
>
}.Build()
208
>
File_temporal_server_api_enums_v1_dlq_proto = out.File
209
>
file_temporal_server_api_enums_v1_dlq_proto_goTypes = nil
210
>
file_temporal_server_api_enums_v1_dlq_proto_depIdxs = nil
211
}
go.temporal.io/server/api/enums/v1/fairness_state.pb.go 20 covered LOC · 2 ranges
Open complete file
123
}
124
125
>
func init() { file_temporal_server_api_enums_v1_fairness_state_proto_init() }
fairness_state.pb.go
126
>
func file_temporal_server_api_enums_v1_fairness_state_proto_init() {
127
>
if File_temporal_server_api_enums_v1_fairness_state_proto != nil {
128
return
129
}
131
>
out := protoimpl.TypeBuilder{
132
>
File: protoimpl.DescBuilder{
133
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
134
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_fairness_state_proto_rawDesc), len(file_temporal_server_api_enums_v1_fairness_state_proto_rawDesc)),
135
>
NumEnums: 1,
136
>
NumMessages: 0,
137
>
NumExtensions: 0,
138
>
NumServices: 0,
139
>
},
140
>
GoTypes: file_temporal_server_api_enums_v1_fairness_state_proto_goTypes,
141
>
DependencyIndexes: file_temporal_server_api_enums_v1_fairness_state_proto_depIdxs,
142
>
EnumInfos: file_temporal_server_api_enums_v1_fairness_state_proto_enumTypes,
143
>
}.Build()
144
>
File_temporal_server_api_enums_v1_fairness_state_proto = out.File
145
>
file_temporal_server_api_enums_v1_fairness_state_proto_goTypes = nil
146
>
file_temporal_server_api_enums_v1_fairness_state_proto_depIdxs = nil
147
}
go.temporal.io/server/api/enums/v1/nexus.pb.go 20 covered LOC · 2 ranges
Open complete file
158
}
159
160
>
func init() { file_temporal_server_api_enums_v1_nexus_proto_init() }
nexus.pb.go
161
>
func file_temporal_server_api_enums_v1_nexus_proto_init() {
162
>
if File_temporal_server_api_enums_v1_nexus_proto != nil {
163
return
164
}
166
>
out := protoimpl.TypeBuilder{
167
>
File: protoimpl.DescBuilder{
168
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
169
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_nexus_proto_rawDesc), len(file_temporal_server_api_enums_v1_nexus_proto_rawDesc)),
170
>
NumEnums: 1,
171
>
NumMessages: 0,
172
>
NumExtensions: 0,
173
>
NumServices: 0,
174
>
},
175
>
GoTypes: file_temporal_server_api_enums_v1_nexus_proto_goTypes,
176
>
DependencyIndexes: file_temporal_server_api_enums_v1_nexus_proto_depIdxs,
177
>
EnumInfos: file_temporal_server_api_enums_v1_nexus_proto_enumTypes,
178
>
}.Build()
179
>
File_temporal_server_api_enums_v1_nexus_proto = out.File
180
>
file_temporal_server_api_enums_v1_nexus_proto_goTypes = nil
181
>
file_temporal_server_api_enums_v1_nexus_proto_depIdxs = nil
182
}
go.temporal.io/server/api/enums/v1/predicate.pb.go 20 covered LOC · 2 ranges
Open complete file
170
}
171
172
>
func init() { file_temporal_server_api_enums_v1_predicate_proto_init() }
predicate.pb.go
173
>
func file_temporal_server_api_enums_v1_predicate_proto_init() {
174
>
if File_temporal_server_api_enums_v1_predicate_proto != nil {
175
return
176
}
178
>
out := protoimpl.TypeBuilder{
179
>
File: protoimpl.DescBuilder{
180
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
181
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_predicate_proto_rawDesc), len(file_temporal_server_api_enums_v1_predicate_proto_rawDesc)),
182
>
NumEnums: 1,
183
>
NumMessages: 0,
184
>
NumExtensions: 0,
185
>
NumServices: 0,
186
>
},
187
>
GoTypes: file_temporal_server_api_enums_v1_predicate_proto_goTypes,
188
>
DependencyIndexes: file_temporal_server_api_enums_v1_predicate_proto_depIdxs,
189
>
EnumInfos: file_temporal_server_api_enums_v1_predicate_proto_enumTypes,
190
>
}.Build()
191
>
File_temporal_server_api_enums_v1_predicate_proto = out.File
192
>
file_temporal_server_api_enums_v1_predicate_proto_goTypes = nil
193
>
file_temporal_server_api_enums_v1_predicate_proto_depIdxs = nil
194
}
go.temporal.io/server/api/enums/v1/replication.pb.go 20 covered LOC · 2 ranges
Open complete file
314
}
315
316
>
func init() { file_temporal_server_api_enums_v1_replication_proto_init() }
replication.pb.go
317
>
func file_temporal_server_api_enums_v1_replication_proto_init() {
318
>
if File_temporal_server_api_enums_v1_replication_proto != nil {
319
return
320
}
322
>
out := protoimpl.TypeBuilder{
323
>
File: protoimpl.DescBuilder{
324
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
325
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_replication_proto_rawDesc), len(file_temporal_server_api_enums_v1_replication_proto_rawDesc)),
326
>
NumEnums: 3,
327
>
NumMessages: 0,
328
>
NumExtensions: 0,
329
>
NumServices: 0,
330
>
},
331
>
GoTypes: file_temporal_server_api_enums_v1_replication_proto_goTypes,
332
>
DependencyIndexes: file_temporal_server_api_enums_v1_replication_proto_depIdxs,
333
>
EnumInfos: file_temporal_server_api_enums_v1_replication_proto_enumTypes,
334
>
}.Build()
335
>
File_temporal_server_api_enums_v1_replication_proto = out.File
336
>
file_temporal_server_api_enums_v1_replication_proto_goTypes = nil
337
>
file_temporal_server_api_enums_v1_replication_proto_depIdxs = nil
338
}
go.temporal.io/server/api/enums/v1/task.pb.go 20 covered LOC · 2 ranges
Open complete file
456
}
457
458
>
func init() { file_temporal_server_api_enums_v1_task_proto_init() }
task.pb.go
459
>
func file_temporal_server_api_enums_v1_task_proto_init() {
460
>
if File_temporal_server_api_enums_v1_task_proto != nil {
461
return
462
}
464
>
out := protoimpl.TypeBuilder{
465
>
File: protoimpl.DescBuilder{
466
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
467
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_task_proto_rawDesc), len(file_temporal_server_api_enums_v1_task_proto_rawDesc)),
468
>
NumEnums: 3,
469
>
NumMessages: 0,
470
>
NumExtensions: 0,
471
>
NumServices: 0,
472
>
},
473
>
GoTypes: file_temporal_server_api_enums_v1_task_proto_goTypes,
474
>
DependencyIndexes: file_temporal_server_api_enums_v1_task_proto_depIdxs,
475
>
EnumInfos: file_temporal_server_api_enums_v1_task_proto_enumTypes,
476
>
}.Build()
477
>
File_temporal_server_api_enums_v1_task_proto = out.File
478
>
file_temporal_server_api_enums_v1_task_proto_goTypes = nil
479
>
file_temporal_server_api_enums_v1_task_proto_depIdxs = nil
480
}
go.temporal.io/server/api/enums/v1/workflow.pb.go 20 covered LOC · 2 ranges
Open complete file
275
}
276
277
>
func init() { file_temporal_server_api_enums_v1_workflow_proto_init() }
workflow.pb.go
278
>
func file_temporal_server_api_enums_v1_workflow_proto_init() {
279
>
if File_temporal_server_api_enums_v1_workflow_proto != nil {
280
return
281
}
283
>
out := protoimpl.TypeBuilder{
284
>
File: protoimpl.DescBuilder{
285
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
286
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_workflow_proto_rawDesc), len(file_temporal_server_api_enums_v1_workflow_proto_rawDesc)),
287
>
NumEnums: 3,
288
>
NumMessages: 0,
289
>
NumExtensions: 0,
290
>
NumServices: 0,
291
>
},
292
>
GoTypes: file_temporal_server_api_enums_v1_workflow_proto_goTypes,
293
>
DependencyIndexes: file_temporal_server_api_enums_v1_workflow_proto_depIdxs,
294
>
EnumInfos: file_temporal_server_api_enums_v1_workflow_proto_enumTypes,
295
>
}.Build()
296
>
File_temporal_server_api_enums_v1_workflow_proto = out.File
297
>
file_temporal_server_api_enums_v1_workflow_proto_goTypes = nil
298
>
file_temporal_server_api_enums_v1_workflow_proto_depIdxs = nil
299
}
go.temporal.io/server/api/enums/v1/workflow_task_type.pb.go 20 covered LOC · 2 ranges
Open complete file
124
}
125
127
>
func file_temporal_server_api_enums_v1_workflow_task_type_proto_init() {
128
>
if File_temporal_server_api_enums_v1_workflow_task_type_proto != nil {
129
return
130
}
132
>
out := protoimpl.TypeBuilder{
133
>
File: protoimpl.DescBuilder{
134
>
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
135
>
RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_enums_v1_workflow_task_type_proto_rawDesc), len(file_temporal_server_api_enums_v1_workflow_task_type_proto_rawDesc)),
136
>
NumEnums: 1,
137
>
NumMessages: 0,
138
>
NumExtensions: 0,
139
>
NumServices: 0,
140
>
},
141
>
GoTypes: file_temporal_server_api_enums_v1_workflow_task_type_proto_goTypes,
142
>
DependencyIndexes: file_temporal_server_api_enums_v1_workflow_task_type_proto_depIdxs,
143
>
EnumInfos: file_temporal_server_api_enums_v1_workflow_task_type_proto_enumTypes,
144
>
}.Build()
145
>
File_temporal_server_api_enums_v1_workflow_task_type_proto = out.File
146
>
file_temporal_server_api_enums_v1_workflow_task_type_proto_goTypes = nil
147
>
file_temporal_server_api_enums_v1_workflow_task_type_proto_depIdxs = nil
148
}
go.temporal.io/server/common/log/tag/zap_tag.go 16 covered LOC · 4 ranges
Open complete file
26
}
27
28
>
func (t ZapTag) Field() zap.Field {
zap_tag.go
29
>
return t.field
30
>
}
31
33
>
return t.field.Key
34
>
}
35
36
func (t ZapTag) Value() any {
44
}
45
46
>
func NewStringTag(key string, value string) ZapTag {
zap_tag.go
47
>
return ZapTag{
48
>
field: zap.String(key, value),
49
>
}
50
>
}
51
52
func NewStringsTag(key string, value []string) ZapTag {
118
}
119
120
>
func NewBoolTag(key string, value bool) ZapTag {
zap_tag.go
121
>
return ZapTag{
122
>
field: zap.Bool(key, value),
123
>
}
124
>
}
125
126
func NewErrorTag(key string, value error) ZapTag {
go.temporal.io/server/common/quotas/clocked_rate_limiter.go 13 covered LOC · 3 ranges
Open complete file
25
)
26
27
>
func NewClockedRateLimiter(rateLimiter *rate.Limiter, timeSource clock.TimeSource) ClockedRateLimiter {
clocked_rate_limiter.go
28
>
return ClockedRateLimiter{
29
>
rateLimiter: rateLimiter,
30
>
timeSource: timeSource,
31
>
recycleCh: make(chan struct{}),
32
>
}
33
>
}
34
36
>
return l.AllowN(l.timeSource.Now(), 1)
37
>
}
38
40
>
return l.rateLimiter.AllowN(now, token)
41
>
}
42
43
// ClockedReservation wraps a rate.Reservation with a clockwork.Clock. It is used to ensure that the reservation
go.temporal.io/server/common/quotas/rate_limiter_impl.go 12 covered LOC · 1 range
Open complete file
24
// NewRateLimiter returns a new rate limiter that can handle dynamic
25
// configuration updates
27
>
limiter := rate.NewLimiter(rate.Limit(newRPS), newBurst)
28
>
ts := clock.NewRealTimeSource()
29
>
rl := &RateLimiterImpl{
30
>
rps: newRPS,
31
>
burst: newBurst,
32
>
timeSource: ts,
33
>
ClockedRateLimiter: NewClockedRateLimiter(limiter, ts),
34
>
}
35
>
36
>
return rl
37
>
}
38
39
// SetRPS sets the rate of the rate limiter
go.temporal.io/server/common/clock/time_source.go 6 covered LOC · 2 ranges
Open complete file
31
32
// NewRealTimeSource returns a timeSource that uses the real wall timeSource time.
34
>
return RealTimeSource{}
35
>
}
36
37
// Now returns the current time, with the location set to UTC.
39
>
return time.Now().UTC()
40
>
}
41
42
// Since returns the time elapsed since t
go.temporal.io/server/common/log/with_logger.go 4 covered LOC · 2 ranges
Open complete file
14
// With returns Logger instance that prepend every log entry with tags. If logger implements
15
// WithLogger it is used, otherwise every log call will be intercepted.
16
>
func With(logger Logger, tags ...tag.Tag) Logger {
with_logger.go
17
>
if wl, ok := logger.(WithLogger); ok {
19
>
}
20
return newWithLogger(logger, tags...)
21
}