Skip to content

Commit 2ff263a

Browse files
authored
Merge branch 'main' into conf-workflows-e2e-add-chain-write
2 parents 9f1cbcb + 2036f9a commit 2ff263a

4 files changed

Lines changed: 61 additions & 12 deletions

File tree

enclave/nitro/host/host.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1056,7 +1056,10 @@ func main() {
10561056
if err := mainServer.Shutdown(shutdownCtx); err != nil {
10571057
lggr.Errorw("main server shutdown error", "error", err)
10581058
}
1059-
if err := telemetry.close(shutdownCtx); err != nil {
1059+
// Give telemetry its own timeout budget
1060+
flushCtx, flushCancel := context.WithTimeout(context.Background(), *shutdownTimeout)
1061+
defer flushCancel()
1062+
if err := telemetry.close(flushCtx); err != nil {
10601063
lggr.Errorw("telemetry shutdown error", "error", err)
10611064
}
10621065

enclave/server/response_emitter.go

Lines changed: 30 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,10 @@ package server
22

33
import (
44
"encoding/json"
5+
"maps"
56
"math"
67
"net/http"
8+
"sync"
79
"time"
810

911
"github.com/smartcontractkit/chainlink-confidential-compute/types"
@@ -16,6 +18,7 @@ const durationBucketSeconds = 0.01
1618
// ResponseEmitter collects metrics to be included in the response payload
1719
// instead of emitting them via OTel.
1820
type ResponseEmitter struct {
21+
mu sync.RWMutex
1922
metrics []types.MetricEvent
2023
// Captured at construction so WriteErrorResponse can report the total request
2124
// duration on error, not just on the success path.
@@ -33,9 +36,13 @@ func NewResponseEmitter() *ResponseEmitter {
3336
}
3437

3538
func (e *ResponseEmitter) Emit(event string, details map[string]any) {
39+
details = roundDurations(maps.Clone(details))
40+
41+
e.mu.Lock()
42+
defer e.mu.Unlock()
3643
e.metrics = append(e.metrics, types.MetricEvent{
3744
Event: event,
38-
Details: roundDurations(details),
45+
Details: details,
3946
})
4047
}
4148

@@ -59,16 +66,29 @@ func roundDurations(details map[string]any) map[string]any {
5966
// Duplicate event names collapse to the last occurrence; use GetMetricEvents
6067
// when per-occurrence counts matter (e.g. repeated capability calls).
6168
func (e *ResponseEmitter) GetMetrics() map[string]any {
62-
result := make(map[string]any)
63-
for _, m := range e.metrics {
64-
result[m.Event] = m.Details
65-
}
66-
return result
69+
metrics, _ := e.Snapshot()
70+
return metrics
6771
}
6872

6973
// GetMetricEvents returns every collected event in order, preserving duplicates.
7074
func (e *ResponseEmitter) GetMetricEvents() []types.MetricEvent {
71-
return e.metrics
75+
_, events := e.Snapshot()
76+
return events
77+
}
78+
79+
// Snapshot returns consistent map and ordered representations of all events.
80+
func (e *ResponseEmitter) Snapshot() (map[string]any, []types.MetricEvent) {
81+
e.mu.RLock()
82+
defer e.mu.RUnlock()
83+
84+
metrics := make(map[string]any)
85+
events := make([]types.MetricEvent, len(e.metrics))
86+
for i, event := range e.metrics {
87+
details := maps.Clone(event.Details)
88+
events[i] = types.MetricEvent{Event: event.Event, Details: details}
89+
metrics[event.Event] = details
90+
}
91+
return metrics, events
7292
}
7393

7494
// WriteErrorResponse writes a JSON error response that includes the accumulated
@@ -87,10 +107,11 @@ func (e *ResponseEmitter) WriteErrorResponse(w http.ResponseWriter, msg string,
87107

88108
w.Header().Set("Content-Type", "application/json")
89109
w.WriteHeader(statusCode)
110+
metrics, metricEvents := e.Snapshot()
90111
resp := types.EnclaveErrorResponse{
91112
Error: msg,
92-
Metrics: e.GetMetrics(),
93-
MetricEvents: e.GetMetricEvents(),
113+
Metrics: metrics,
114+
MetricEvents: metricEvents,
94115
}
95116
_ = json.NewEncoder(w).Encode(resp)
96117
}

enclave/server/response_emitter_test.go

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package server
22

33
import (
4+
"sync"
45
"testing"
56

67
"github.com/stretchr/testify/assert"
@@ -35,3 +36,28 @@ func TestResponseEmitter_ErrorResponseCarriesMetricEvents(t *testing.T) {
3536
// the full ordered list (2 capability_execution + 1 request_completed).
3637
require.Len(t, e.GetMetricEvents(), 2)
3738
}
39+
40+
func TestResponseEmitter_ConcurrentEmit(t *testing.T) {
41+
const (
42+
parallelCalls = 30 // chainlink-common's default per-execution pending-call limit
43+
eventsPerCall = 3 // capability_started, capability_execution, capability_finished
44+
)
45+
46+
e := NewResponseEmitter()
47+
start := make(chan struct{})
48+
var wg sync.WaitGroup
49+
wg.Add(parallelCalls)
50+
for call := range parallelCalls {
51+
go func() {
52+
defer wg.Done()
53+
<-start
54+
e.Emit("capability_started", map[string]any{"step_ref": call})
55+
e.Emit("capability_execution", map[string]any{"step_ref": call})
56+
e.Emit("capability_finished", map[string]any{"step_ref": call})
57+
}()
58+
}
59+
close(start)
60+
wg.Wait()
61+
62+
assert.Len(t, e.GetMetricEvents(), parallelCalls*eventsPerCall)
63+
}

enclave/server/server.go

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -756,8 +756,7 @@ func (s *enclaveServer) handleExecute(w http.ResponseWriter, r *http.Request) {
756756
responseEmitter.WriteErrorResponse(w, fmt.Sprintf("error creating attestation: %v", err), http.StatusInternalServerError)
757757
return
758758
}
759-
resp.Metrics = responseEmitter.GetMetrics()
760-
resp.MetricEvents = responseEmitter.GetMetricEvents()
759+
resp.Metrics, resp.MetricEvents = responseEmitter.Snapshot()
761760
resp.Attestation = att
762761

763762
respBytes, err := json.Marshal(resp)

0 commit comments

Comments
 (0)