diff --git a/docs/adr/adr-0001-retire-sketch-core.md b/docs/adr/adr-0001-retire-sketch-core.md new file mode 100644 index 00000000..2f845f53 --- /dev/null +++ b/docs/adr/adr-0001-retire-sketch-core.md @@ -0,0 +1,137 @@ +# ADR-0001: Retire `sketch-core`; `asap_sketchlib` is the single Rust algorithm crate + +| | | +|---|---| +| Status | **Accepted** (retrospective — work already shipped) | +| Date | 2026-05-01 | +| Deciders | Project ASAP maintainers | +| Supersedes | n/a | +| Superseded by | n/a | + +## Context + +Before this work, the Rust side of the ASAP system had **two** +algorithm crates with overlapping responsibilities: + +- `asap_sketchlib` — the canonical algorithm crate, holding + `DDSketch`, `KLL`, `HLL`, `Count`, `CountMin`, `CMSHeap`, etc. + Already shared with `sketchlib-bench` and external benchmark + harnesses. +- `sketch-core` — a wrapper crate vendored into + `ASAPQuery-backend/asap-common/sketch-core/` (and its sibling + forks in `ASAPQuery/asap-common/sketch-core/` and + `sketchlib-bench/sketch-core/`). It re-exposed the algorithms + with ASAP-specific wire-format types (`DdSketchState`, + `CountSketchState`, etc.), `apply_delta` methods, and the + `Strategy::Legacy` vs `Strategy::Sketchlib` `ImplMode` switch. + +The duplication was a source of: + +- **Cross-language drift risk.** The Go side had one canonical + algorithm crate (`sketchlib-go`); having two on the Rust side + made a third potential drift target. +- **Apply-delta logic in the wrong place.** `apply_delta` + implementations belonged with the algorithms, not in a wrapper + crate. Bug fixes (e.g., DDSketch delta count reconstruction) + had to be written in two places. +- **Maintenance overhead.** Three on-disk forks of the same + crate had to stay in sync; PRs touching the shared types + needed three identical commits. + +The `ImplMode` switch (Legacy vs Sketchlib) was originally +introduced to allow gradual migration off hand-written matrix +implementations onto `asap_sketchlib`-backed implementations. +Once the Sketchlib path was the default and Legacy was no longer +exercised in production paths, the switch became dead weight. + +## Decision + +1. **Retire `sketch-core` entirely.** Move all of its content + into `asap_sketchlib`: + - Wire-format types (`DdSketch`, `CountSketch`, `HllSketch`, + `CountMinSketch`, `KllSketch`, `CountMinSketchWithHeap`) + into the existing `src/sketches/.rs` files alongside + the high-throughput algorithms — single home per sketch + concept. + - `apply_delta` implementations into the same files so the + algorithm + its delta semantics live together. + - New sibling files for sketches that had no existing home in + `asap_sketchlib` (`hydra_kll.rs`, `set_aggregator.rs`, + `delta_set_aggregator.rs`). + - The previous `*_sketchlib.rs` FFI wrapper files inlined into + their main sketch file (e.g., `SketchlibCms` lives directly + inside `countmin.rs`). + - The standalone `asap_runtime` module for the legacy + `ImplMode` configuration was created and then removed (see + decision 3). + +2. **Drop the three on-disk `sketch-core` forks.** All consumers + (`ASAPQuery`, `ASAPQuery-backend`, `sketchlib-bench`) depend + on `asap_sketchlib` directly via git URL. + +3. **Drop the `ImplMode` (Legacy / Sketchlib) dispatch.** Always + use the `asap_sketchlib`-backed implementation. Removed: + - The `asap_runtime` module entirely. + - `KllBackend` enum (Sketchlib + Legacy variants); `KllSketch` + now holds `SketchlibKll` directly. + - The `dsrs` (datasketches-rs) dependency that provided the + legacy KLL backend. + - The `clap` (asap-cli feature) and `ctor` (legacy-mode test + initializer) dependencies. + - The backend's `--sketch-cms-impl` / `--sketch-kll-impl` / + `--sketch-cmwh-impl` CLI args and their `config::configure` + plumbing. + +4. **Rename `CountMinDelta` → `CountMinSketchDelta`** to match + the `Delta` pattern of the other delta types + (`DdSketchDelta`, `CountSketchDelta`, `HllSketchDelta`). + +5. **Rename to avoid wire-format collisions:** + - `octo_delta::HllDelta` (single-register, octo path) keeps + its name; the wire-format multi-register delta becomes + `HllSketchDelta`. + - `common::input::HeapItem` (polymorphic key type) keeps its + name; the wire-format CMSHeap item becomes `CmsHeapItem`. + +## Consequences + +### Positive + +- Single canonical Rust algorithm crate, mirroring `sketchlib-go` + on the Go side. Cross-language drift is now a 1:1 concern, not + 1:N. +- `apply_delta` logic lives next to the algorithm it applies to. +- Three on-disk forks deleted; consumers all track + `asap_sketchlib` directly. +- ~3,000 lines of net code removed once `Legacy` paths and + `dsrs` / `clap` / `ctor` / `asap-cli` deps were dropped. + +### Negative / Tradeoffs + +- The wire-format types now sit alongside the high-throughput + types in the same files (e.g., both `DDSketch` and `DdSketch` + in `src/sketches/ddsketch.rs`). File sizes grow by 30–80%. + Acceptable: the alternative (separate crate or separate + module) recreates the duplication problem this ADR solves. +- Cross-language parity tests between `sketchlib-go` and + `asap_sketchlib` are now the *only* defense against algorithm + drift. The R1 risk in the design doc explicitly calls this + out; mitigation (statistical-output cross-language harness in + `sketchlib-bench`) is tracked for follow-up. + +### Compatibility + +- Wire format unchanged: `SketchEnvelope.payload` bytes are + byte-identical pre/post-retirement. The + `tests::elastic_dsl_query_tests::tests::test_esdsl_time_range_query` + test was relaxed from `assert_eq!(value, 291.0)` to a `±1` + tolerance — `asap_sketchlib`'s KLL gives 290 on this + distribution where `dsrs` gave 291. Both are within KLL's + rank-error bound; the test was previously over-tight. + +## References + +- [ProjectASAP/asap_sketchlib#36](https://github.com/ProjectASAP/asap_sketchlib/pull/36) — sketch-core merged into existing `src/sketches/` layout +- [ProjectASAP/ASAPQuery-backend#73](https://github.com/ProjectASAP/ASAPQuery-backend/pull/73) — backend consumer migration +- [ProjectASAP/ASAPQuery#309](https://github.com/ProjectASAP/ASAPQuery/pull/309) — legacy-fork consumer migration + branch-pin cleanup + elastic DSL test fix +- [ProjectASAP/asap_sketchlib#37](https://github.com/ProjectASAP/asap_sketchlib/pull/37) — wire-format / `apply_delta` semantic alignment with `sketchlib-go` (DDSketch count reconstruction, CountMin/CountSketch field additions, out-of-bounds policy) diff --git a/docs/adr/adr-0002-extract-precompute-runtime.md b/docs/adr/adr-0002-extract-precompute-runtime.md new file mode 100644 index 00000000..9c0fc3f0 --- /dev/null +++ b/docs/adr/adr-0002-extract-precompute-runtime.md @@ -0,0 +1,195 @@ +# ADR-0002: Extract the precompute runtime to `asap-precompute-{go,rs}` + +| | | +|---|---| +| Status | **Proposed** (gates Phase 2 / Phase 3 of the edge-framework migration) | +| Date | 2026-05-02 | +| Deciders | Project ASAP maintainers | +| Supersedes | n/a | +| Superseded by | n/a | + +## Context + +Today the windowing / delta / scheduler state machine lives +inside the Go OTel processors (~3850 LoC across +`opentelemetry-collector-contrib-patch/processor/{ddsketch,kll,hll,countsketch,countminsketch}processor/processor.go`) +and inside the Rust ingest path +(`ASAPQuery-backend/asap-query-engine/src/precompute_operators/*.rs` ++ `drivers/ingest/otel.rs::apply_modified_otlp_delta_bytes`). + +Each of those files conflates four concerns: + +1. **OTel binding** — implementing `processor.Metrics`, accepting + `pmetric.Metrics`, calling `nextConsumer.ConsumeMetrics`. +2. **Data shape adapter** — extracting `(timestamp, attrs, value)` + tuples out of `pmetric.Gauge | Sum | DDSketchDataPoint | …`. +3. **Runtime state machine** — `accumulateIntoWindow`, + `flushWindow`, snapshot caches, label matchers, scheduler. +4. **Output binding** — emitting `pmetric.Metrics` of the right + typed variant, calling `nextConsumer.ConsumeMetrics`. + +(1, 2, 4) are host-specific. (3) is host-neutral. The framework +refactor (see `design-asap-edge-framework.md` §6) extracts (3) so +adapters provide (1, 2, 4) only. + +## Decision + +### What gets extracted + +Two new artifacts, one per language: + +- **`asap-precompute-go`** — Go module living under + `ASAPCollector/asap-precompute-go/` (subdirectory, not a + separate repo for now; promotion to a separate repo is a + Phase-7 / repo-rename concern). +- **`asap-precompute-rs`** — Rust crate living under + `ASAPCollector/asap-precompute-rs/` for build-time deployment + inside an OTel-shaped ingest, and dual-published from the same + source as a normal Rust crate that `ASAPQuery-backend` can + depend on for its backend-side ingest path. Concretely the crate + source lives in this repo and `ASAPQuery-backend` consumes it + via git URL, the same way it depends on `asap_sketchlib` today. + +### Public API + +The trait, types, and config surface are pinned by the design +doc (§5.1, §6.2, §6.3). The summary contract: + +- `Observation` — host-neutral input. Includes timestamp, metric + name, labels, and a `value` (`Float | Hash | Bytes | + Envelope`). +- `Precompute` trait — generic over `SketchT: Sketch`. + - `observe(&Observation) -> Result<(), Overflow>` + - `observe_envelope(SketchEnvelope) -> Result<(), Error>` + - `tick(now_ms) -> Vec` +- `PrecomputeConfig` — exactly today's processor config knobs + (sketch type, window size, matchers, delta thresholds, sketch + params) plus two new fields needed once the runtime is no + longer behind OTel's pipeline backpressure: `max_series` + + `OnOverflow` (Drop | Block | EvictOldest). +- `Sketch` trait family — `Sketch + QuantileSketch + + CardinalitySketch + FrequencySketch` (split). Each in-process + sketch impls only the queries it supports. +- Crash recovery is intentionally **not** in the trait. A future + `PersistentPrecompute: Precompute` extension trait can add + `snapshot()` / `restore()` later. + +### Behavior preservation + +- Wire format unchanged. `SketchEnvelope.payload` bytes are + byte-identical pre/post-extraction. +- Window semantics unchanged. Tumbling-vs-sliding-vs-batch logic + is moved verbatim from each processor file into + `asap-precompute-go::window.go` (Go) / + `asap-precompute-rs::window.rs` (Rust). +- Snapshot cache invariants unchanged. The + `snapshots map[string][]byte` (Go) and + `IngestState.sketch_snapshots` (Rust) maps move into + `snapshot_cache.go` / `snapshot_cache.rs` with the same + per-series-key contract. +- Backwards-compat for the Go OTel processors during Phase 2: + each existing `processor/{ddsketch,kll,hll,countsketch,countminsketch}processor/processor.go` + reduces to a ~50-line shim that delegates to + `asap-precompute-go`. The shim's `ConsumeMetrics` signature, + metric output schema, and config keys remain identical; + config-file changes are not required for existing deployments. + +### Performance contract + +- **Go (Phase 2):** per-observation `Observe` latency p99 must + stay within 10% of the pre-refactor in-line implementation. + Verified via the existing fake-exporter b3-delta benchmark. +- **Rust (Phase 3):** entry point `observe_envelope` stays + bit-identical to today's per-accumulator + `apply_proto_delta_bytes`. No behavior drift on backend + PromQL output. + +### Repo / module layout + +``` +ASAPCollector/ +├── asap-precompute-go/ +│ ├── go.mod // module github.com/ProjectASAP/asap-precompute-go +│ ├── observation.go // Observation type +│ ├── envelope.go // SketchEnvelope type (Go view of the proto) +│ ├── precompute.go // Precompute interface + impl +│ ├── window.go // tumbling / sliding / batch logic +│ ├── snapshot_cache.go // outbound + inbound snapshot caches; ComputeDelta +│ ├── matchers.go // LabelMatcher / aggregate_by / seriesKey +│ ├── config.go // PrecomputeConfig + AggregationMode + OnOverflow +│ ├── adapter.go // Adapter trait + helpers +│ └── controlchannel/ // ControlChannel trait + HttpPollChannel impl +├── asap-precompute-rs/ +│ ├── Cargo.toml // crate name asap-precompute-rs +│ └── src/ +│ ├── lib.rs +│ ├── observation.rs +│ ├── envelope.rs // re-exports asap_sketchlib::proto::sketchlib::SketchEnvelope +│ ├── precompute.rs // Precompute trait + generic impl +│ ├── window.rs +│ ├── snapshot_cache.rs +│ ├── matchers.rs +│ ├── config.rs +│ ├── adapter.rs +│ └── control_channel.rs +└── opentelemetry-collector-contrib-patch/processor/ + ├── ddsketchprocessor/processor.go # Phase 2: ~50 LoC shim delegating to asap-precompute-go + ├── kllprocessor/processor.go # Phase 2: ~50 LoC shim + ├── hllprocessor/processor.go # Phase 2: ~50 LoC shim + ├── countsketchprocessor/processor.go # Phase 2: ~50 LoC shim + └── countminsketchprocessor/processor.go # Phase 2: ~50 LoC shim +``` + +### Sequencing + +- Phase 2 (Go) and Phase 3 (Rust) can proceed in parallel + because they touch different repos. Both must merge before + Phase 4 (Telegraf adapter) starts, since Telegraf reuses + `asap-precompute-go` and the OTel adapter shim must be a + proven shape before generalizing it. + +## Consequences + +### Positive + +- Each existing processor shrinks from ~700–950 LoC to ~50 LoC + shim, removing the same conflated-concern logic five times. +- Telegraf / OTAP / Vector adapters become viable — they reuse + `asap-precompute-{go,rs}` rather than re-implementing window / + snapshot / matcher logic. +- Backend ingest path becomes a thin adapter calling the same + Rust crate the agents would use, eliminating the agent / + backend duplication for delta-apply logic. +- Future Sketch trait additions (e.g., `observe_batch` for + columnar Arrow ingest) become single-crate changes. + +### Negative + +- Cross-repo dependency added: Go OTel patches now pull in + `github.com/ProjectASAP/asap-precompute-go`. The build + pipeline (`build_sketchcollector.sh`) needs to handle two + module sources. +- Test surface doubles temporarily during the migration: each + function moves through a "duplicated, behavior-verified, then + delete the original" sequence to catch divergence. Phase 2 + exit criterion (b3-delta produces same value, p99 within 10%) + is the gate. + +### Compatibility + +- No wire-format changes. +- No config-file changes for existing OTel collector + deployments. +- Backend PromQL output preserved bit-for-bit (R4 mitigation). + +## Phase-2 / Phase-3 execution plan + +See [`docs/phase-2-execution-plan.md`](../phase-2-execution-plan.md) +for the file-by-file extraction map covering all 5 OTel +processors. + +## References + +- [`docs/design-asap-edge-framework.md`](../design-asap-edge-framework.md) §3, §6, §9 (Phases 2–3) +- ADR-0001 (sketch-core retirement — prerequisite that simplified the Rust algorithm crate before this extraction) +- ADR-0003 (adapter trait + control channel — defines what the OTel processor shims look like after extraction) diff --git a/docs/adr/adr-0003-adapter-trait-and-control-channel.md b/docs/adr/adr-0003-adapter-trait-and-control-channel.md new file mode 100644 index 00000000..067661e0 --- /dev/null +++ b/docs/adr/adr-0003-adapter-trait-and-control-channel.md @@ -0,0 +1,179 @@ +# ADR-0003: `Adapter` trait, `ControlChannel` trait, and Strategy A/B encoding + +| | | +|---|---| +| Status | **Proposed** (gates Phase 4 / Phase 5 / Phase 6 of the edge-framework migration) | +| Date | 2026-05-02 | +| Deciders | Project ASAP maintainers | +| Supersedes | n/a | +| Superseded by | n/a | + +## Context + +ASAP's edge precompute runtime needs to ride into multiple host +runtimes (OTel Collector, Telegraf, Vector, OTAP Dataflow). Each +host has its own data model (`pmetric.Metrics` / `telegraf.Metric` +/ `vector::Event` / `arrow.RecordBatch`), config format, and +reload semantics. Without explicit contracts for how the runtime +plugs into each host, every adapter would diverge in subtle ways +and the cross-platform invariants the framework promises (single +sketch source-of-truth, stable wire format, plan-driven +configuration) would erode. + +Three contracts need pinning before adapter implementation +starts: + +1. The **`Adapter` trait** that every Layer-4 shim implements. +2. The **`ControlChannel` trait** that delivers + `PrecomputeConfig` from the controller to the runtime. +3. The **Strategy A vs Strategy B encoding choice** for how + `SketchEnvelope` rides each host's native event model. + +This ADR pins all three. + +## Decision + +### 1. `Adapter` trait + +```rust +pub trait Adapter { + type Event; // pmetric.Metrics, telegraf.Metric, vector.Event, arrow.RecordBatch + + fn decode<'a>(&self, ev: &'a Self::Event) -> Result>, Error>; + fn encode(&self, envelopes: Vec) -> Result; + fn schedule_tick(&self, period: Duration, cb: TickCallback); + fn emit_telemetry(&self, stats: &PrecomputeStats); +} +``` + +Go side has an idiomatic equivalent. + +Adapters emit only raw `Observation`s. `agg_id` resolution +(matchers → agg_id) is the runtime's responsibility, not the +adapter's. This is non-negotiable: the runtime is the single +source of truth for plan binding; adapters do not duplicate +matcher logic. + +### 2. `ControlChannel` trait + +```rust +pub trait ControlChannel: Send { + fn poll(&mut self) -> Option; + fn ack(&mut self, plan_version: u64); +} +``` + +Three implementations land in tree: + +| Impl | Used by | Status | +|---|---|---| +| `OpAmpChannel` | OTel Collector existing deploys | already implemented in `controller/src/opamp/` | +| `HttpPollChannel` | All non-OTel adapters; new `_asap-otelcol_` deployments | new, ~30–50 LoC. Polls controller's existing `GET /api/v1/plan?host_id=` endpoint. | +| `FileWatchChannel` | Sidecar fallback for environments where outbound HTTP isn't allowed | new, last-resort. | + +### 3. Hard rule: `ControlChannel` runs in an internal goroutine / task + +Every adapter, including OTel, runs its `ControlChannel` inside +its own owned goroutine (Go) / tokio task (Rust). The adapter +maintains an `atomic.Pointer[PrecomputeConfig]` (Go) / +`Arc>` (Rust) that the hot path reads +on each `Observe` / `Tick`. + +This rule is forced by an audit finding: every host platform's +native reload mechanism (OTel Collector's SIGHUP + OpAMP +supervisor; Telegraf's `--watch-config`; Vector's +`reload_config_and_respawn`; OTAP's `NodeControlMsg::Config`) +**rebuilds** the plugin instance and clobbers in-memory sketch +state. ASAP cannot route plan pushes through any of those +without losing all sketches. The internal-goroutine pattern is +precedented in OTel's `tailsamplingprocessor`. + +Configuration in the host's config file is **bootstrap-only**: +controller URL, agent ID, auth. Sketch parameters (sketch type, +window size, matchers, max_series, delta thresholds) are +never in host config; they arrive only via `ControlChannel::poll`. + +### 4. Strategy A vs Strategy B encoding + +Every adapter picks one strategy for how `SketchEnvelope.payload` +rides its native event: + +- **Strategy A — Native typed variant.** Extend the host's + schema with a sketch-typed oneOf / variant. Realistic only + for OTel-family platforms. ASAP today uses Strategy A for + OTel Collector via modified-OTLP `pmetric.Metric.data` oneOf. +- **Strategy B — Opaque bytes in the host's nearest binary-clean + carrier.** Use whatever bytes/binary field the native event + already exposes, marked with the standardized + `_asap_envelope` + metadata keys. Required for Telegraf, + Vector, OTAP. + +Standardized well-known keys (project-level standard; every +adapter using Strategy B uses exactly these spellings): + +| Key | Logical type | Required? | Meaning | +|---|---|---|---| +| `_asap_envelope` | bytes | yes | `SketchEnvelope.payload` proto bytes | +| `_asap_sketch_type` | string | yes | `"DDSketch"` / `"KLLSketch"` / etc. | +| `_asap_agg_id` | uint64 | yes | matches controller plan | +| `_asap_schema_version` | uint32 | yes | matches `SketchEnvelope.schema_version` | +| `_asap_window_start_ms` | uint64 | yes | window lower bound | +| `_asap_window_end_ms` | uint64 | yes | window upper bound | +| `_asap_encoding` | string | optional | `PROTO_FULL` (default) / `PROTO_DELTA` / `MSGPACK` | + +Per-platform Strategy-B carrier (verified against current upstream +source): + +| Platform | Carrier for `_asap_envelope` | Notes | +|---|---|---| +| Telegraf | `string`-typed field on `telegraf.Metric` holding 8-bit-clean envelope bytes | `[]byte` is coerced to `string` by `convertField`; bytes survive but the type tag is lost. A custom `asap` Telegraf Serializer is required for HTTP / Kafka / File sinks. InfluxDB line-protocol and Prometheus remote-write are not supported sinks. | +| Vector | `Value::Bytes(Bytes)` field on a `LogEvent` | NOT `MetricValue::Sketch` (hard-typed to `AgentDDSketch`). Required sinks: `vector` native (Protobuf), Kafka with `encoding=native\|raw_message`, S3 with `encoding=native`. | +| OTAP Dataflow | `AttributeValueType::Bytes` on the per-row attribute child batch of `OtapArrowRecords` | NOT a sibling top-level Binary column — OTAP's strict schema validator rejects extension columns. Encoding as an attribute is OTLP round-trippable. | + +### 5. Bandwidth invariant — sketches stay sketches end-to-end + +Every adapter MUST preserve the framework's bandwidth invariant +(see design doc §5.2): `SketchEnvelope.payload` bytes flow +end-to-end without ever being "exploded" back to per-sample +observations. Adapters that fail this invariant break the entire +framework's bandwidth promise. CI signal: encoded byte-count vs +raw-sample byte-count ratio for any adapter must match the +sketch's expected compression ratio. + +## Consequences + +### Positive + +- Adapters get a clear, narrow contract. Changing an adapter + doesn't touch runtime / wire / control-plane code. +- ControlChannel becomes a single abstraction shared across all + four hosts. The same controller can serve all four adapters + via either OpAMP or HTTP-poll without per-host customization. +- Strategy A/B explicit means future platforms have a clear + template. "Take the project's nearest binary-clean field, + attach the well-known keys" is the recipe. + +### Negative + +- Per-platform Strategy-B carrier is non-uniform. The doc + acknowledges this; adapter implementers must read the + per-platform table carefully. +- Internal-goroutine pattern means each adapter ships a copy of + the same controller-poll loop. ~50 LoC duplication is + acceptable; the alternative (a shared library with a `Send + + Sync` poll task) adds Go/Rust async-runtime coupling that's + worse than the duplication. + +### Compatibility + +- OTel adapter today already emits via Strategy A. No + wire-format change. +- New OTel deployments using `HttpPollChannel` instead of OpAMP + are a deploy-time choice; the controller serves both. + +## References + +- [`docs/design-asap-edge-framework.md`](../design-asap-edge-framework.md) §5 (envelope), §6 (Precompute trait), §7 (Adapter + Strategy A/B), §8 (ControlChannel) +- ADR-0001 (sketch-core retirement) — establishes that `SketchEnvelope` is the canonical wire type +- ADR-0002 (extract precompute runtime) — defines the runtime that adapters wrap +- OTel `tailsamplingprocessor` — precedent for internal-goroutine config update pattern diff --git a/docs/phase-2-execution-plan.md b/docs/phase-2-execution-plan.md new file mode 100644 index 00000000..f726730c --- /dev/null +++ b/docs/phase-2-execution-plan.md @@ -0,0 +1,280 @@ +# Phase 2 execution plan — extract `asap-precompute-go` + +_Companion to [ADR-0002](adr/adr-0002-extract-precompute-runtime.md). +File-by-file extraction map for moving the runtime out of the +five Go OTel processors into a shared `asap-precompute-go` +module._ + +## Inventory of source files + +``` +opentelemetry-collector-contrib-patch/processor/ +├── ddsketchprocessor/processor.go 942 LoC +├── kllprocessor/processor.go 720 LoC +├── hllprocessor/processor.go 785 LoC +├── countsketchprocessor/processor.go 663 LoC +└── countminsketchprocessor/processor.go 739 LoC + -------- + 3849 LoC total +``` + +The processors fall into two structural patterns that need to be +harmonized in the extracted runtime: + +- **Pattern A** (DDSketch / KLL / HLL): nested + `resourceWindow → scopeWindow → metricWindow → sketchSeries` + hierarchy. Per-`pmetric.ScopeMetrics` aggregation. Explicit + `accumulateIntoWindow` / `flushWindow` pair. +- **Pattern B** (CountSketch / CountMin): flat + `windowSketch` (map: partitionKey → sketch). Per-metric + aggregation. Timer-driven `startWindowLoop` / + `emitWindowAndReset`. + +Phase 2 unifies both into the generic `Precompute[SketchT]` +shape from ADR-0002. The window manager picks +tumbling/sliding/batch internally based on `PrecomputeConfig`. + +## Target layout + +``` +asap-precompute-go/ +├── go.mod // module github.com/ProjectASAP/asap-precompute-go +├── observation.go // Observation + ObservationValue (~50 LoC) +├── envelope.go // SketchEnvelope view of the sketchlib-go proto (~30 LoC) +├── precompute.go // Precompute interface + generic impl (~250 LoC) +├── window.go // tumbling / sliding / batch logic (~250 LoC) +├── snapshot_cache.go // outbound + inbound caches; ComputeDelta (~200 LoC) +├── matchers.go // LabelMatcher, seriesKey, seriesAttrs (~150 LoC) +├── config.go // PrecomputeConfig, AggregationMode, OnOverflow (~120 LoC) +├── adapter.go // Adapter interface + Decode/Encode helpers (~100 LoC) +├── controlchannel/ +│ ├── channel.go // ControlChannel interface (~30 LoC) +│ ├── http_poll.go // HttpPollChannel impl (~80 LoC) +│ └── opamp.go // OpAmpChannel adapter wrapping existing controller/opamp (~50 LoC) +├── telemetry.go // PrecomputeStats + recordInput/recordOutput (~80 LoC) +└── otel/ // OTel-flavored Adapter helpers (consumed by Phase-2 shims) + ├── decode.go // pmetric.Metrics → []Observation + ├── encode.go // []SketchEnvelope → pmetric.Metrics + └── seriesattrs.go // attribute key construction +``` + +Approximate total: 1500–1800 LoC. Each existing OTel processor +shrinks to ~50–80 LoC shim that constructs a +`Precompute[]` and delegates `ConsumeMetrics`. + +## Function-level extraction map + +For each existing function, where it goes after Phase 2: + +### Pattern A (DDSketch, KLL, HLL) — same map applies to all three + +| Today | Layer | Becomes | +|---|---|---| +| `type resourceWindow / scopeWindow / metricWindow / sketchSeries` | 3 | `precompute.go` — collapse into a single `series` struct keyed by `(agg_id, label_key)` since the resource/scope hierarchy was an OTel-side concern, not an algorithmic concern. | +| `func newProcessor` | 4 | stays in shim; constructs `Precompute[*ddsketch.DDSketch]` | +| `func Start` | 4 | stays in shim; spawns ticker goroutine that calls `Precompute.Tick` and emits via `Adapter.Encode` + `next.ConsumeMetrics` | +| `func Shutdown` | 4 | stays in shim; cancels ticker, calls `Precompute.Shutdown` | +| `func ConsumeMetrics` | 4 | stays in shim; calls `Adapter.Decode(md)` → `for _, o := range obs { p.pc.Observe(o) }` → `next.ConsumeMetrics(ctx, md)` (pass-through) | +| `func processBatch / processScopeMetrics` | 4 | becomes the `otel.Decode` helper; produces `[]Observation` | +| `func consumeDDSketchDataPoints / consumeGaugeDataPoints` | 4 | folded into `otel.Decode`; produces `Observation::Envelope` for sketch-typed inputs and `Observation::Float` for scalar | +| `func decodeDDSketchDataPoint` | 4 | folded into `otel.Decode` envelope path | +| `func cacheInboundSnapshot` | 3 | `snapshot_cache.go::CacheInbound` | +| `func newSketchSeries / updateWindow / merge / ensureSketch` | 3 | `window.go` window-state helpers | +| `func serializeDDSketch` | 1 | already lives in `sketchlib-go`; called via `Sketch.Snapshot()` | +| `func computeDDSketchDelta` | 3 | `snapshot_cache.go::ComputeDelta` | +| `func attributesKey / seriesKey / seriesAttrs / matchesMatchers / newSeriesFrom` | 3 | `matchers.go::SeriesKey / SeriesAttrs / Matches / NewSeries` | +| `func accumulateIntoWindow` | 3 | `window.go::Observe` (merged with `Precompute::Observe`) | +| `func getOrCreateMetricWindow` | 3 | private to `window.go` | +| `func accumulateGaugeMetric / accumulate{DD,KLL,HLL}SketchMetric` | 3+4 | sketch-specific `Sketch::Observe` lives at L1; the routing logic (raw vs envelope) is in `Precompute::Observe` | +| `func flushWindow` | 3 | `Precompute::Tick`; emits `[]SketchEnvelope` | +| `func buildMetric / buildMergedSketchMetric / buildQuantileMetric` | 4 | becomes `otel.Encode`; produces `pmetric.Metrics` from `[]SketchEnvelope` | +| `func enableSelfMonitoring / shutdownMonitor` | 4 | stays in shim | +| `func recordInput / recordOutput / activeSeriesCount` | 3 | `telemetry.go` | + +### Pattern B (CountSketch, CountMin) — same map + +| Today | Layer | Becomes | +|---|---|---| +| `type windowSketch` | 3 | absorbed into `series` struct in `precompute.go` (one map: `(agg_id, label_key) → SketchT`) | +| `func newProcessor / Start / Shutdown / Capabilities / ConsumeMetrics` | 4 | shim | +| `func processMetrics / consumeBatch / ingestMetric / dpValue` | 4 | `otel.Decode`; produces `[]Observation` | +| `func matchesMatchers / encodeKey / seriesKey / seriesAttrs / buildPartitionKey / encodeAttributesAsKey` | 3 | `matchers.go` | +| `func accumulateIntoWindow / updateWindowSketch / mergeWindowSketch` | 3 | `Precompute::Observe + window.go` | +| `func startWindowLoop / emitWindowAndReset / buildWindowMetricsAndReset` | 3+4 | timer goroutine moves to `Precompute` (driven by `Adapter::ScheduleTick`); `buildWindowMetricsAndReset` becomes `otel.Encode` | +| `func inboundDecode{CS,CMS} / mergeWindow{CS,CMS}` | 3 | `snapshot_cache.go::ApplyDelta` (envelope-in path on `Precompute::ObserveEnvelope`) | +| `func serialize{CountSketch,CMS} / deserialize{CMS} / clone{CS,CMS}` | 1 | already `sketchlib-go` API | +| `func newConfiguredCountSketch / nextPowerOfTwo` | 4 | stays in shim (it's `Config` validation) | +| `func recordInput / recordOutput / activeSeriesCount` | 3 | `telemetry.go` | + +### Cross-cutting + +- All five processors have an `enableSelfMonitoring` / + `shutdownMonitor` pair that emits OTel-shaped runtime metrics. + These become two pieces: + - `telemetry.go::PrecomputeStats` (host-neutral counters in + L3) — incremented from inside `Precompute`. + - The shim's `enableSelfMonitoring` reads `PrecomputeStats` + via `Adapter::EmitTelemetry` and constructs the OTel-shaped + metrics. + +## Per-processor shim shape (post-Phase-2) + +Every existing processor file becomes ~50-80 LoC of this shape: + +```go +package ddsketchprocessor + +import ( + precompute "github.com/ProjectASAP/asap-precompute-go" + otelhost "github.com/ProjectASAP/asap-precompute-go/otel" + "github.com/ProjectASAP/asap-precompute-go/controlchannel" + + "go.opentelemetry.io/collector/component" + "go.opentelemetry.io/collector/consumer" + "go.opentelemetry.io/collector/pdata/pmetric" +) + +type ddsketchProcessor struct { + pc precompute.Precompute // generic over the sketch type bound by Config.SketchType + adapter *otelhost.Adapter + cc controlchannel.ControlChannel + cfg *Config + next consumer.Metrics + logger *zap.Logger + monitor *otelhost.SelfMonitor // wraps PrecomputeStats + shutdown chan struct{} +} + +func (p *ddsketchProcessor) Capabilities() consumer.Capabilities { + return consumer.Capabilities{MutatesData: false} +} + +func (p *ddsketchProcessor) Start(ctx context.Context, host component.Host) error { + if err := p.pc.Start(ctx); err != nil { return err } + go p.controlChannelLoop(ctx) + go p.tickLoop(ctx) + return p.monitor.Start(ctx, host) +} + +func (p *ddsketchProcessor) Shutdown(ctx context.Context) error { + close(p.shutdown) + return p.pc.Shutdown(ctx) +} + +func (p *ddsketchProcessor) ConsumeMetrics(ctx context.Context, md pmetric.Metrics) error { + obs, err := p.adapter.Decode(md) + if err != nil { return err } + for _, o := range obs { + if err := p.pc.Observe(&o); err != nil { /* OnOverflow */ } + } + return p.next.ConsumeMetrics(ctx, md) +} + +func (p *ddsketchProcessor) tickLoop(ctx context.Context) { + t := time.NewTicker(p.cfg.WindowSize) + defer t.Stop() + for { + select { + case <-p.shutdown: return + case now := <-t.C: + envelopes := p.pc.Tick(now.UnixMilli()) + md := p.adapter.Encode(envelopes) + if err := p.next.ConsumeMetrics(ctx, md); err != nil { /* log */ } + } + } +} + +func (p *ddsketchProcessor) controlChannelLoop(ctx context.Context) { + t := time.NewTicker(p.cfg.ControlPollInterval) + defer t.Stop() + for { + select { + case <-p.shutdown: return + case <-t.C: + if cs := p.cc.Poll(); cs != nil { + p.pc.UpdateConfig(cs) // atomic swap inside Precompute + } + } + } +} +``` + +The five processors share this scaffolding; the only per-processor +differences are: + +- The generic `Precompute` type parameter (`*ddsketch.DDSketch` + vs `*kll.KllSketch` vs ...). +- The `otel.Adapter` decode path — which `Metric.data` oneOf + variants it recognizes (DDSketch / KLLSketch / HLLSketch / + CountSketch / CountMinSketch). +- Default config values (`alpha`, `k`, `precision`, `width`, + `depth`). + +A future refactor could collapse all five files into a single +generic shim parameterized by sketch type. Phase 2 keeps them +separate for OCB build-config compatibility (`builder-config.yaml` +references each processor's package path). + +## Phase 2 work breakdown + +Sequenced for incremental verification — each step is shippable +and reversible: + +| Step | Scope | Verification | +|---|---|---| +| **2.1** Bootstrap `asap-precompute-go` module + types | `observation.go`, `envelope.go`, `config.go`, `adapter.go` (interface only). No logic, just type definitions matching ADR-0002 contracts. | `go vet ./...` clean. | +| **2.2** Implement `matchers.go` and `snapshot_cache.go` | Move pure-data-structure logic with no host coupling: `LabelMatcher`, `SeriesKey`, snapshot cache, `ComputeDelta`. | Unit tests against fixed input vectors copied from existing processor tests. | +| **2.3** Implement `window.go` and `precompute.go` | Generic Precompute implementation. Drives matchers + snapshot cache. Generic over `Sketch` interface. | Unit tests with mock Sketch (table-driven; verify tumbling, sliding, late-data, max_series, OnOverflow). | +| **2.4** Implement `otel/{decode,encode}.go` | Pmetric decoders for DDSketch + Gauge + Sum. Encoders for `Metric.data = DDSketch{...}`. | Unit tests with fixed `pmetric.Metrics` fixtures. | +| **2.5** Refactor `ddsketchprocessor` to shim | Delete the in-file accumulateIntoWindow/flushWindow/snapshot logic; replace with the shim shape above. Existing tests must still pass. | b3-delta e2e: same `19.49` value at offset −90s. | +| **2.6** Refactor `kllprocessor` | Same as 2.5 but for KLL. | Existing KLL accuracy tests pass. | +| **2.7** Refactor `hllprocessor` | Same as 2.5 but for HLL. | Existing HLL accuracy tests pass. | +| **2.8** Refactor `countsketchprocessor` | Pattern B — verify flat-window collapse to (agg_id, label_key) keying preserves behavior. | CountSketch top-K accuracy reducer (`P8`) matches pre-extraction. | +| **2.9** Refactor `countminsketchprocessor` | Same as 2.8 for CMS. | P8 accuracy reducer matches. | +| **2.10** `controlchannel/http_poll.go` + adapter wiring | First non-OpAMP control channel. Backward-compat: existing OpAMP-driven deploys keep using `OpAmpChannel`; new deploys can opt into `HttpPollChannel`. | b3-delta e2e survives a runtime config push (sketch_type unchanged, window_size changed) without state loss. | +| **2.11** Performance gate | Per-observation latency p99 within 10% of pre-refactor (R2). | Bench against the existing fake-exporter throughput harness. | + +Steps 2.1–2.4 can run in parallel (no inter-dependencies). +Steps 2.5–2.9 are sequential (each builds on the verified shim +shape from the previous). 2.10–2.11 gate the phase exit. + +## Risks during the migration + +- **R-Phase-2-A: Pattern-B → unified series-map collapse hides + partition semantics.** CountSketch / CMS today aggregate by + `partitionKey` (CS) or `aggregationKey` (CMS); these are + string concatenations of `(metric_name, label_subset)`. The + unified `(agg_id, label_key)` schema must preserve the same + string. Mitigation: in 2.8 / 2.9, write a key-equivalence test + before refactoring. +- **R-Phase-2-B: Tick goroutine races with ConsumeMetrics.** Today + each processor uses an internal mutex. The extracted + `Precompute` must keep the same locking discipline (per-series + rwmutex; coarse global mutex around tick swap). Mitigation: + keep mutex shape identical in 2.3; race detector run on + refactored shims. +- **R-Phase-2-C: OCB build manifest drift.** ASAPCollector's + `builder-config.yaml` references each processor's Go package + path. After Phase 2, those paths still work (the processor + packages still exist; they just call a new module). The new + `asap-precompute-go` module needs to be added to OCB's `gomod` + list. Mitigation: 2.5 includes `build_sketchcollector.sh` + smoke run before merge. + +## Phase exit criterion (blocking) + +1. All five OTel processors are ≤80 LoC each (excluding factory + / config-validation boilerplate). +2. b3-delta e2e produces the observed `19.49` value at offset + −90s, identical to pre-extraction. +3. Per-observation `Observe` latency p99 within 10% of + pre-refactor. +4. P8 accuracy reducer per-row error matches pre-extraction + for all five sketch types. +5. `cargo test` / `go test` clean across all touched packages. +6. `controlchannel.HttpPollChannel` smoke test: collector + running with the new control channel survives a controller + plan push that changes window size (without losing sketch + state mid-window). + +If any of (1)–(6) fails, Phase 2 does not merge.