Skip to content

Commit 2cdaf59

Browse files
committed
Extract finish lifecycle package
Move finish-aware series coordination behind a dedicated internal package. Export only the supported lifecycle protocol and use phase tokens to make measurement release and collection completion explicit. The lifecycle now owns serialization, preventing callers from reaching packed-state helpers.
1 parent 749aa10 commit 2cdaf59

6 files changed

Lines changed: 505 additions & 325 deletions

File tree

sdk/metric/internal/aggregate/finish_lifecycle.go

Lines changed: 0 additions & 190 deletions
This file was deleted.

sdk/metric/internal/aggregate/finish_lifecycle_test.go

Lines changed: 0 additions & 112 deletions
This file was deleted.

sdk/metric/internal/aggregate/finish_sum.go

Lines changed: 15 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -10,13 +10,12 @@ import (
1010
"time"
1111

1212
"go.opentelemetry.io/otel/attribute"
13+
"go.opentelemetry.io/otel/sdk/metric/internal/finish"
1314
"go.opentelemetry.io/otel/sdk/metric/metricdata"
1415
)
1516

1617
type finishSumValue[N int64 | float64] struct {
17-
mu sync.Mutex
18-
19-
lifecycle finishLifecycle
18+
lifecycle finish.Lifecycle
2019
value atomicCounter[N]
2120
attrs attribute.Set
2221
start time.Time
@@ -45,65 +44,58 @@ func (v *finishSumValue[N]) measure(
4544
value N,
4645
lazy lazyFilteredAttributes,
4746
) bool {
48-
if !v.lifecycle.acquireMeasurement() {
47+
measurement, ok := v.lifecycle.AcquireMeasurement()
48+
if !ok {
4949
return false
5050
}
5151
v.value.add(value)
5252
if v.dropExemplars {
53-
v.lifecycle.releaseMeasurement()
53+
measurement.Release()
5454
return true
5555
}
56-
defer v.lifecycle.releaseMeasurement()
56+
defer measurement.Release()
5757
v.reservoir.Offer(ctx, value, lazy)
5858
return true
5959
}
6060

6161
func (v *finishSumValue[N]) finish(t time.Time) {
62-
v.mu.Lock()
63-
defer v.mu.Unlock()
6462
if !v.overflow {
65-
v.lifecycle.finish(t)
63+
v.lifecycle.Finish(t)
6664
}
6765
}
6866

6967
func (v *finishSumValue[N]) collectCumulative(
7068
t time.Time,
7169
) (metricdata.DataPoint[N], bool, bool) {
72-
v.mu.Lock()
73-
defer v.mu.Unlock()
74-
return v.collect(v.lifecycle.startCumulativeCollection(t))
70+
return v.collect(v.lifecycle.BeginCumulativeCollection(t))
7571
}
7672

7773
func (v *finishSumValue[N]) collectDelta(
7874
t time.Time,
7975
) (metricdata.DataPoint[N], bool, bool) {
80-
v.mu.Lock()
81-
defer v.mu.Unlock()
82-
return v.collect(v.lifecycle.startDeltaCollection(t))
76+
return v.collect(v.lifecycle.BeginDeltaCollection(t))
8377
}
8478

8579
func (v *finishSumValue[N]) collect(
86-
decision collectionDecision,
80+
collection finish.Collection,
8781
) (metricdata.DataPoint[N], bool, bool) {
88-
if !decision.emit {
82+
if !collection.ShouldEmit() {
8983
return metricdata.DataPoint[N]{}, false, false
9084
}
91-
defer v.lifecycle.completeCollection(decision)
85+
defer collection.Complete()
9286

9387
dp := metricdata.DataPoint[N]{
9488
Attributes: v.attrs,
9589
StartTime: v.start,
96-
Time: decision.time,
90+
Time: collection.Time(),
9791
Value: v.value.load(),
9892
}
9993
collectExemplars(&dp.Exemplars, v.reservoir.Collect)
100-
return dp, true, decision.retire
94+
return dp, true, collection.ShouldRetire()
10195
}
10296

10397
func (v *finishSumValue[N]) shutdown() {
104-
v.mu.Lock()
105-
defer v.mu.Unlock()
106-
v.lifecycle.retire()
98+
v.lifecycle.Retire()
10799
}
108100

109101
// FinishSum contains the operations of a finish-aware Sum aggregation.

0 commit comments

Comments
 (0)