From cb605211eaa9d89819e0e848af71ec2bfb52359c Mon Sep 17 00:00:00 2001 From: Matt Van Horn Date: Sat, 27 Jun 2026 04:34:30 -0700 Subject: [PATCH 1/3] feat(collector): log extension startup completion and duration Capture a start timestamp in main.go before the launch log and thread it through lifecycle.NewManager. After the collector starts successfully, emit a single info-level "OpenTelemetry Lambda extension startup complete" log carrying a startup_duration field so operators can see when the extension is ready and measure its cold-start contribution. Signed-off-by: Matt Van Horn --- collector/internal/lifecycle/manager.go | 7 +- collector/internal/lifecycle/manager_test.go | 86 ++++++++++++++++++++ collector/main.go | 4 +- 3 files changed, 95 insertions(+), 2 deletions(-) diff --git a/collector/internal/lifecycle/manager.go b/collector/internal/lifecycle/manager.go index 02203e4941..5e99b7ff82 100644 --- a/collector/internal/lifecycle/manager.go +++ b/collector/internal/lifecycle/manager.go @@ -22,6 +22,7 @@ import ( "path/filepath" "sync" "syscall" + "time" "github.com/open-telemetry/opentelemetry-lambda/collector/lambdalifecycle" @@ -53,9 +54,10 @@ type manager struct { wg sync.WaitGroup lifecycleListeners []lambdalifecycle.Listener initType lambdalifecycle.InitType + startTime time.Time } -func NewManager(ctx context.Context, logger *zap.Logger, version string) (context.Context, *manager) { +func NewManager(ctx context.Context, logger *zap.Logger, version string, startTime time.Time) (context.Context, *manager) { ctx, cancel := context.WithCancel(ctx) sigs := make(chan os.Signal, 1) @@ -102,6 +104,7 @@ func NewManager(ctx context.Context, logger *zap.Logger, version string) (contex extensionClient: extensionClient, listener: listener, initType: initType, + startTime: startTime, } factories, _ := lambdacomponents.Components(res.ExtensionID) @@ -119,6 +122,8 @@ func (lm *manager) Run(ctx context.Context) error { return err } + lm.logger.Info("OpenTelemetry Lambda extension startup complete", zap.Duration("startup_duration", time.Since(lm.startTime))) + lm.wg.Add(1) go func() { if err := lm.processEvents(ctx); err != nil { diff --git a/collector/internal/lifecycle/manager_test.go b/collector/internal/lifecycle/manager_test.go index b52effe8dc..fc407f187f 100644 --- a/collector/internal/lifecycle/manager_test.go +++ b/collector/internal/lifecycle/manager_test.go @@ -24,15 +24,101 @@ import ( "os" "path/filepath" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.uber.org/zap" "go.uber.org/zap/zaptest" + "go.uber.org/zap/zaptest/observer" "github.com/open-telemetry/opentelemetry-lambda/collector/internal/extensionapi" "github.com/open-telemetry/opentelemetry-lambda/collector/internal/telemetryapi" ) +const startupCompleteMsg = "OpenTelemetry Lambda extension startup complete" + +func TestRunLogsStartupDuration(t *testing.T) { + shutdownServer := func(t *testing.T) *httptest.Server { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(200) + _, err := w.Write([]byte(`{"time":"2006-01-02T15:04:05.000Z", "eventType":"SHUTDOWN", "record":{}}`)) + require.NoError(t, err) + _, err = io.ReadAll(r.Body) + require.NoError(t, err, "failed to read request body: %v", err) + })) + t.Cleanup(server.Close) + return server + } + extensionEventTypes := []extensionapi.EventType{extensionapi.Invoke, extensionapi.Shutdown} + + t.Run("emits a single startup-complete log with a startup_duration field", func(t *testing.T) { + core, logs := observer.New(zap.InfoLevel) + logger := zap.New(core) + + server := shutdownServer(t) + u, err := url.Parse(server.URL) + require.NoError(t, err) + + lm := manager{ + collector: &MockCollector{}, + logger: logger, + listener: telemetryapi.NewListener(logger), + extensionClient: extensionapi.NewClient(logger, u.Host, extensionEventTypes), + startTime: time.Now(), + } + require.NoError(t, lm.Run(context.Background())) + + entries := logs.FilterMessage(startupCompleteMsg).All() + require.Len(t, entries, 1, "expected exactly one startup-complete log") + field, ok := entries[0].ContextMap()["startup_duration"] + require.True(t, ok, "startup-complete log must carry a startup_duration field") + duration, ok := field.(time.Duration) + require.True(t, ok, "startup_duration must be a duration field") + assert.GreaterOrEqual(t, duration, time.Duration(0), "startup_duration must be non-negative") + }) + + t.Run("zero-value start time produces a valid duration without panicking", func(t *testing.T) { + core, logs := observer.New(zap.InfoLevel) + logger := zap.New(core) + + server := shutdownServer(t) + u, err := url.Parse(server.URL) + require.NoError(t, err) + + lm := manager{ + collector: &MockCollector{}, + logger: logger, + listener: telemetryapi.NewListener(logger), + extensionClient: extensionapi.NewClient(logger, u.Host, extensionEventTypes), + // startTime intentionally left as the zero value. + } + require.NotPanics(t, func() { + require.NoError(t, lm.Run(context.Background())) + }) + + entries := logs.FilterMessage(startupCompleteMsg).All() + require.Len(t, entries, 1) + duration, ok := entries[0].ContextMap()["startup_duration"].(time.Duration) + require.True(t, ok) + assert.GreaterOrEqual(t, duration, time.Duration(0)) + }) + + t.Run("does not emit startup-complete log when collector start fails", func(t *testing.T) { + core, logs := observer.New(zap.InfoLevel) + logger := zap.New(core) + + lm := manager{ + collector: &MockCollector{err: fmt.Errorf("test start error")}, + logger: logger, + extensionClient: extensionapi.NewClient(logger, "", extensionEventTypes), + startTime: time.Now(), + } + require.Error(t, lm.Run(context.Background())) + assert.Equal(t, 0, logs.FilterMessage(startupCompleteMsg).Len(), "no startup-complete log on failure") + }) +} + type MockCollector struct { err error } diff --git a/collector/main.go b/collector/main.go index 4304997e88..c4430f45f2 100644 --- a/collector/main.go +++ b/collector/main.go @@ -18,6 +18,7 @@ import ( "context" "flag" "fmt" + "time" "github.com/open-telemetry/opentelemetry-lambda/collector/lambdalifecycle" @@ -44,9 +45,10 @@ func main() { } logger := logging.NewLogger() + startTime := time.Now() logger.Info("Launching OpenTelemetry Lambda extension", zap.String("version", Version)) - ctx, lm := lifecycle.NewManager(context.Background(), logger, Version) + ctx, lm := lifecycle.NewManager(context.Background(), logger, Version, startTime) // Set the new lifecycle manager as the lifecycle notifier for all other components. lambdalifecycle.SetNotifier(lm) From be5b17d1a5d1ef3362d01dff1c455bdcdd281af8 Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Wed, 22 Jul 2026 16:08:40 -0700 Subject: [PATCH 2/3] test(lifecycle): make startup-log test deterministic The existing lifecycle test became flaky after the extension startup log line landed; synchronize the assertion so it no longer races the goroutine emitting the log. --- collector/internal/lifecycle/manager_test.go | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/collector/internal/lifecycle/manager_test.go b/collector/internal/lifecycle/manager_test.go index fc407f187f..344aa1c121 100644 --- a/collector/internal/lifecycle/manager_test.go +++ b/collector/internal/lifecycle/manager_test.go @@ -169,11 +169,12 @@ func TestRun(t *testing.T) { extensionClient: extensionapi.NewClient(logger, u.Host, extensionEventTypes), } lm.wg.Add(1) + runErr := make(chan error, 1) go func() { - require.NoError(t, lm.Run(ctx)) + runErr <- lm.Run(ctx) }() lm.wg.Done() - + assert.NoError(t, <-runErr) } func TestProcessEvents(t *testing.T) { From b061e910cb6eb1fea13f1861e032288d9cb93dad Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Wed, 29 Jul 2026 06:29:08 -0700 Subject: [PATCH 3/3] test: synchronize the waitgroup case and drop the tautological assertions @wpessers: the extension startup log made the waitgroup case in TestRun flaky, since wg.Done() raced lm.Run's own use of it. The test server now blocks until the request is in flight, so Done() lands at a deterministic point. Also drops the zero-value startTime subtest, which measured a clamped time.Since(time.Time{}), and the non-negative duration assertion that could not fail. --- collector/internal/lifecycle/manager_test.go | 53 ++++++++------------ 1 file changed, 21 insertions(+), 32 deletions(-) diff --git a/collector/internal/lifecycle/manager_test.go b/collector/internal/lifecycle/manager_test.go index 344aa1c121..4a1eef7174 100644 --- a/collector/internal/lifecycle/manager_test.go +++ b/collector/internal/lifecycle/manager_test.go @@ -73,35 +73,8 @@ func TestRunLogsStartupDuration(t *testing.T) { require.Len(t, entries, 1, "expected exactly one startup-complete log") field, ok := entries[0].ContextMap()["startup_duration"] require.True(t, ok, "startup-complete log must carry a startup_duration field") - duration, ok := field.(time.Duration) + _, ok = field.(time.Duration) require.True(t, ok, "startup_duration must be a duration field") - assert.GreaterOrEqual(t, duration, time.Duration(0), "startup_duration must be non-negative") - }) - - t.Run("zero-value start time produces a valid duration without panicking", func(t *testing.T) { - core, logs := observer.New(zap.InfoLevel) - logger := zap.New(core) - - server := shutdownServer(t) - u, err := url.Parse(server.URL) - require.NoError(t, err) - - lm := manager{ - collector: &MockCollector{}, - logger: logger, - listener: telemetryapi.NewListener(logger), - extensionClient: extensionapi.NewClient(logger, u.Host, extensionEventTypes), - // startTime intentionally left as the zero value. - } - require.NotPanics(t, func() { - require.NoError(t, lm.Run(context.Background())) - }) - - entries := logs.FilterMessage(startupCompleteMsg).All() - require.Len(t, entries, 1) - duration, ok := entries[0].ContextMap()["startup_duration"].(time.Duration) - require.True(t, ok) - assert.GreaterOrEqual(t, duration, time.Duration(0)) }) t.Run("does not emit startup-complete log when collector start fails", func(t *testing.T) { @@ -163,17 +136,33 @@ func TestRun(t *testing.T) { require.NoError(t, lm.Run(ctx)) // test with waitgroup counter incremented lm = manager{ - collector: &MockCollector{}, - logger: logger, - listener: telemetryapi.NewListener(logger), - extensionClient: extensionapi.NewClient(logger, u.Host, extensionEventTypes), + collector: &MockCollector{}, + logger: logger, + listener: telemetryapi.NewListener(logger), } + requestStarted := make(chan struct{}) + releaseResponse := make(chan struct{}) + synchronizedServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + close(requestStarted) + <-releaseResponse + w.WriteHeader(200) + _, err := w.Write([]byte(`{"time":"2006-01-02T15:04:05.000Z", "eventType":"SHUTDOWN", "record":{}}`)) + require.NoError(t, err) + _, err = io.ReadAll(r.Body) + require.NoError(t, err, "failed to read request body: %v", err) + })) + defer synchronizedServer.Close() + synchronizedURL, err := url.Parse(synchronizedServer.URL) + require.NoError(t, err) + lm.extensionClient = extensionapi.NewClient(logger, synchronizedURL.Host, extensionEventTypes) lm.wg.Add(1) runErr := make(chan error, 1) go func() { runErr <- lm.Run(ctx) }() + <-requestStarted lm.wg.Done() + close(releaseResponse) assert.NoError(t, <-runErr) }