fix(aggregator): honor forceFlushAll on MicroVM's on-demand Flush#54118
Conversation
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 86a9eb5b3a
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
There was a problem hiding this comment.
Pull request overview
This PR ensures AWS Lambda MicroVM lifecycle-triggered metric flushes (/suspend, /terminate) can include samples still in the current (not-yet-closed) DogStatsD time bucket, preventing “last-moment” telemetry (including the terminate metric itself) from being dropped with no chance to retry after the VM is gone.
Changes:
- Extend
aggregator.Demultiplexer.ForceFlushToSerializerwith aforceFlushAllflag and thread it throughAgentDemultiplexer. - Add a
forceFlushAllOnFlushoption topkg/serverless/metrics.ServerlessMetricAgentand plumb it fromserverless-initvia a newCloudService.ShouldForceFlushAllOnForceFlushToSerializer()hook (true only for MicroVM). - Update and add unit tests to cover open-bucket flush behavior and adjust existing call sites for the new API.
Reviewed changes
Copilot reviewed 20 out of 20 changed files in this pull request and generated 1 comment.
Show a summary per file
| File | Description |
|---|---|
| pkg/serverless/metrics/metric.go | Add forceFlushAllOnFlush to control whether Flush() includes the current open bucket. |
| pkg/serverless/metrics/metric_test.go | Add tests for open-bucket flushing behavior (with/without forceFlushAll). |
| pkg/clusteragent/admission/validate/kubernetesadmissionevents/kubernetesadmissionevents_test.go | Update mock demux flush call for the new ForceFlushToSerializer signature. |
| pkg/aggregator/demultiplexer.go | Update Demultiplexer interface to accept forceFlushAll. |
| pkg/aggregator/demultiplexer_agent.go | Thread forceFlushAll into the trigger path for manual flushes. |
| pkg/aggregator/demultiplexer_agent_test.go | Update tests for new ForceFlushToSerializer signature. |
| pkg/aggregator/aggregator_test.go | Update tests for new ForceFlushToSerializer signature. |
| cmd/serverless-init/main.go | Pass per-cloud-service forceFlushAll intent into the metric agent. |
| cmd/serverless-init/main_test.go | Update metric agent construction for new signature. |
| cmd/serverless-init/lifecycle/server.go | Comment update to reflect lambda_microvm_id tag naming. |
| cmd/serverless-init/lifecycle/server_test.go | Update assertions to match lambda_microvm_id tag naming. |
| cmd/serverless-init/lifecycle/heartbeat.go | Emit lambda_microvm_id:<id> tag instead of microvm_id:<id>. |
| cmd/serverless-init/lifecycle/heartbeat_test.go | Update tests to expect lambda_microvm_id tagging. |
| cmd/serverless-init/cloudservice/service.go | Extend CloudService interface with ShouldForceFlushAllOnForceFlushToSerializer. |
| cmd/serverless-init/cloudservice/service_test.go | Add test pinning MicroVM as the only service returning true. |
| cmd/serverless-init/cloudservice/microvm.go | Return true for MicroVM to force-flush open buckets on lifecycle flushes. |
| cmd/serverless-init/cloudservice/containerapp.go | Implement new CloudService method (returns false). |
| cmd/serverless-init/cloudservice/cloudrun.go | Implement new CloudService method (returns false). |
| cmd/serverless-init/cloudservice/cloudrun_jobs.go | Implement new CloudService method (returns false). |
| cmd/serverless-init/cloudservice/appservice.go | Implement new CloudService method (returns false). |
Comments suppressed due to low confidence (1)
pkg/serverless/metrics/metric_test.go:300
- Same determinism concern as the positive control: using time.Now() for the sample timestamp means this negative-control test can intermittently flush the sample if a bucket boundary is crossed during the loop, even though the intent is to keep the sample in an always-open bucket. Using a future timestamp makes the test stable and ensures it actually validates open-bucket skipping behavior.
now := float64(time.Now().UnixNano()) / float64(time.Second)
agent.AddEnhancedMetric("test.metric", 1.0, pkgmetrics.MetricSourceServerless, now)
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
86a9eb5 to
ae0ba40
Compare
|
🎯 Code Coverage (details) 🔗 Commit SHA: d98a128 | Docs | Datadog PR Page | Give us feedback! |
8e63683 to
faa91f3
Compare
Files inventory check summaryFile checks results against ancestor 8d8cc340: Results for datadog-agent_7.83.0~devel.git.375.d98a128.pipeline.126931455-1_amd64.deb:No change detected |
Static quality checks✅ Please find below the results from static quality gates Successful checksInfo
21 successful checks with minimal change (< 2 KiB)
|
Regression DetectorRegression Detector ResultsMetrics dashboard Baseline: 8d8cc34 Optimization Goals: ✅ No significant changes detected
|
| perf | experiment | goal | Δ mean % | Δ mean % CI | trials | links |
|---|---|---|---|---|---|---|
| ➖ | quality_gate_private_action_runner | memory utilization | +0.73 | [+0.61, +0.85] | 1 | Logs bounds checks dashboard |
| ➖ | quality_gate_security_no_fs_load | memory utilization | +0.29 | [+0.20, +0.38] | 1 | Logs bounds checks dashboard |
| ➖ | quality_gate_idle | memory utilization | +0.18 | [+0.12, +0.23] | 1 | Logs bounds checks dashboard |
| ➖ | quality_gate_security_mean_fs_load | memory utilization | -0.10 | [-0.14, -0.06] | 1 | Logs bounds checks dashboard |
| ➖ | quality_gate_security_idle | memory utilization | -0.20 | [-0.26, -0.14] | 1 | Logs bounds checks dashboard |
| ➖ | quality_gate_idle_all_features | memory utilization | -0.22 | [-0.27, -0.18] | 1 | Logs bounds checks dashboard |
| ➖ | quality_gate_metrics_logs | memory utilization | -0.55 | [-0.80, -0.30] | 1 | Logs bounds checks dashboard |
| ➖ | quality_gate_logs | % cpu utilization | -0.84 | [-1.85, +0.16] | 1 | Logs bounds checks dashboard |
Bounds Checks: ✅ Passed
| perf | experiment | bounds_check_name | replicates_passed | observed_value | links |
|---|---|---|---|---|---|
| ✅ | quality_gate_idle | intake_connections | 10/10 | 3 ≤ 4 | bounds checks dashboard |
| ✅ | quality_gate_idle | memory_usage | 10/10 | 148.01MiB ≤ 154MiB | bounds checks dashboard |
| ✅ | quality_gate_idle | total_bytes_received | 10/10 | 735.47KiB ≤ 819.20KiB | bounds checks dashboard |
| ✅ | quality_gate_idle_all_features | intake_connections | 10/10 | 3 ≤ 4 | bounds checks dashboard |
| ✅ | quality_gate_idle_all_features | memory_usage | 10/10 | 498.43MiB ≤ 512MiB | bounds checks dashboard |
| ✅ | quality_gate_idle_all_features | total_bytes_received | 10/10 | 1.13MiB ≤ 1.25MiB | bounds checks dashboard |
| ✅ | quality_gate_logs | intake_connections | 10/10 | 3 ≤ 6 | bounds checks dashboard |
| ✅ | quality_gate_logs | memory_usage | 10/10 | 182.95MiB ≤ 195MiB | bounds checks dashboard |
| ✅ | quality_gate_logs | missed_bytes | 10/10 | 0B = 0B | bounds checks dashboard |
| ✅ | quality_gate_logs | total_bytes_received | 10/10 | 263.77MiB ≤ 292MiB | bounds checks dashboard |
| ✅ | quality_gate_metrics_logs | cpu_usage | 10/10 | 363.09 ≤ 2000 | bounds checks dashboard |
| ✅ | quality_gate_metrics_logs | intake_connections | 10/10 | 3 ≤ 6 | bounds checks dashboard |
| ✅ | quality_gate_metrics_logs | memory_usage | 10/10 | 401.56MiB ≤ 430MiB | bounds checks dashboard |
| ✅ | quality_gate_metrics_logs | missed_bytes | 10/10 | 0B = 0B | bounds checks dashboard |
| ✅ | quality_gate_metrics_logs | total_bytes_received | 10/10 | 0.94GiB ≤ 1.04GiB | bounds checks dashboard |
| ✅ | quality_gate_private_action_runner | memory_usage | 10/10 | 71.41MiB ≤ 75MiB | bounds checks dashboard |
| ✅ | quality_gate_security_idle | cpu_usage | 10/10 | 30.99 ≤ 100 | bounds checks dashboard |
| ✅ | quality_gate_security_idle | memory_usage | 10/10 | 301.63MiB ≤ 330MiB | bounds checks dashboard |
| ✅ | quality_gate_security_mean_fs_load | cpu_usage | 10/10 | 83.80 ≤ 200 | bounds checks dashboard |
| ✅ | quality_gate_security_mean_fs_load | memory_usage | 10/10 | 277.98MiB ≤ 310MiB | bounds checks dashboard |
| ✅ | quality_gate_security_no_fs_load | cpu_usage | 10/10 | 25.26 ≤ 100 | bounds checks dashboard |
| ✅ | quality_gate_security_no_fs_load | memory_usage | 10/10 | 288.87MiB ≤ 320MiB | bounds checks dashboard |
Explanation
Confidence level: 90.00%
Effect size tolerance: |Δ mean %| ≥ 5.00%
Performance changes are noted in the perf column of each table:
- ✅ = significantly better comparison variant performance
- ❌ = significantly worse comparison variant performance
- ➖ = no significant change in performance
A regression test is an A/B test of target performance in a repeatable rig, where "performance" is measured as "comparison variant minus baseline variant" for an optimization goal (e.g., ingress throughput). Due to intrinsic variability in measuring that goal, we can only estimate its mean value for each experiment; we report uncertainty in that value as a 90.00% confidence interval denoted "Δ mean % CI".
For each experiment, we decide whether a change in performance is a "regression" -- a change worth investigating further -- if all of the following criteria are true:
-
Its estimated |Δ mean %| ≥ 5.00%, indicating the change is big enough to merit a closer look.
-
Its 90.00% confidence interval "Δ mean % CI" does not contain zero, indicating that if our statistical model is accurate, there is at least a 90.00% chance there is a difference in performance between baseline and comparison variants.
-
Its configuration does not mark it "erratic".
Replicate Execution Details
We run multiple replicates for each experiment/variant. However, we allow replicates to be automatically retried if there are any failures, up to 8 times, at which point the replicate is marked dead and we are unable to run analysis for the entire experiment. We call each of these attempts at running replicates a replicate execution. This section lists all replicate executions that failed due to the target crashing or being oom killed.
Note: In the below tables we bucket failures by experiment, variant, and failure type. For each of these buckets we list out the replicate indexes that failed with an annotation signifying how many times said replicate failed with the given failure mode. In the below example the baseline variant of the experiment named experiment_with_failures had two replicates that failed by oom kills. Replicate 0, which failed 8 executions, and replicate 1 which failed 6 executions, all with the same failure mode.
| Experiment | Variant | Replicates | Failure | Logs | Debug Dashboard |
|---|---|---|---|---|---|
| experiment_with_failures | baseline | 0 (x8) 1 (x6) | Oom killed | Debug Dashboard |
The debug dashboard links will take you to a debugging dashboard specifically designed to investigate replicate execution failures.
❌ Retried Profiling Replicate Execution Failures (ddprof)
Note: Profiling replicas may still be executing. See the debug dashboard for up to date status.
| Experiment | Variant | Replicates | Failure | Debug Dashboard |
|---|---|---|---|---|
| quality_gate_idle_all_features | baseline | 10 | Oom killed | Debug Dashboard |
| quality_gate_idle_all_features | comparison | 10 | Oom killed | Debug Dashboard |
| quality_gate_logs | baseline | 10 | Oom killed | Debug Dashboard |
| quality_gate_logs | comparison | 10 | Oom killed | Debug Dashboard |
| quality_gate_metrics_logs | baseline | 10 | Oom killed | Debug Dashboard |
| quality_gate_metrics_logs | comparison | 10 | Oom killed | Debug Dashboard |
| quality_gate_security_idle | comparison | 10 | Oom killed | Debug Dashboard |
| quality_gate_security_no_fs_load | baseline | 10 | Crashed (exit code: 134) | Debug Dashboard |
| quality_gate_security_no_fs_load | comparison | 10 | Crashed (exit code: 134) | Debug Dashboard |
CI Pass/Fail Decision
✅ Passed. All Quality Gates passed.
- quality_gate_idle, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_idle, bounds check intake_connections: 10/10 replicas passed. Gate passed.
- quality_gate_idle, bounds check total_bytes_received: 10/10 replicas passed. Gate passed.
- quality_gate_idle_all_features, bounds check intake_connections: 10/10 replicas passed. Gate passed.
- quality_gate_idle_all_features, bounds check total_bytes_received: 10/10 replicas passed. Gate passed.
- quality_gate_idle_all_features, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_metrics_logs, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_metrics_logs, bounds check intake_connections: 10/10 replicas passed. Gate passed.
- quality_gate_metrics_logs, bounds check cpu_usage: 10/10 replicas passed. Gate passed.
- quality_gate_metrics_logs, bounds check missed_bytes: 10/10 replicas passed. Gate passed.
- quality_gate_metrics_logs, bounds check total_bytes_received: 10/10 replicas passed. Gate passed.
- quality_gate_private_action_runner, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_security_mean_fs_load, bounds check cpu_usage: 10/10 replicas passed. Gate passed.
- quality_gate_security_mean_fs_load, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_security_no_fs_load, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_security_no_fs_load, bounds check cpu_usage: 10/10 replicas passed. Gate passed.
- quality_gate_security_idle, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_security_idle, bounds check cpu_usage: 10/10 replicas passed. Gate passed.
- quality_gate_logs, bounds check total_bytes_received: 10/10 replicas passed. Gate passed.
- quality_gate_logs, bounds check missed_bytes: 10/10 replicas passed. Gate passed.
- quality_gate_logs, bounds check memory_usage: 10/10 replicas passed. Gate passed.
- quality_gate_logs, bounds check intake_connections: 10/10 replicas passed. Gate passed.
/suspend and /terminate call Flush() to make telemetry durable before a snapshot or teardown, but ForceFlushToSerializer hardcoded forceFlushAll to false, so a metric emitted right before the call could sit in the still-open bucket with no retry once the VM is gone. Thread a real forceFlushAll through Demultiplexer.ForceFlushToSerializer. ServerlessMetricAgent exposes it as two methods — Flush (unforced) and FlushAll (forced) — and lifecycle/server.go picks between them based on the hook, not the cloud service: /terminate calls FlushAll since there's no resume to catch a sample left in an open bucket; /suspend must keep using plain Flush, since forcing the bucket there would let a later flush after resume send a second, partial point for that same timestamp, overwriting the pre-suspend data in the backend instead of merging with it. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
faa91f3 to
d98a128
Compare
Driven by Codex's comment on #54111: #54111 (comment)
/suspendand/terminatecallFlush()to make MicroVM telemetry durable before the VM is snapshotted or torn down.ForceFlushToSerializerignored that intent and always skipped the current, still-open time bucket — a metric emitted right before the call, including the/terminatemetric itself, could get dropped with nothing to retry once the VM is gone. That lines up withaws.lambda.enhanced.microvm.terminatehaving zero data points anywhere in the org for the last 30 days.This threads a real
forceFlushAllflag from MicroVM down throughServerlessMetricAgent.Flush()intoDemultiplexer.ForceFlushToSerializer. Every other cloud service opts out — none of them callFlush()outside the normal periodic/shutdown cycle, so nothing else changes.The other issue from that comment —
SampleDrainernever being wired up — needs new drain plumbing inpkg/aggregator'stimeSamplerWorker; the existingshutdown()/waitForShutdown()is one-shot and would panic on a second/suspend. That's a separate PR.Test plan
forceFlushAllis true, never delivered when false.dda inv teston every touched package (pkg/aggregator,pkg/serverless,cmd/serverless-init,pkg/clusteragent/.../kubernetesadmissionevents) — all green.🤖 Generated with Claude Code