Skip to content

About

Heavily Rust-optimised transform pipeline for Elastic Stack data (Beats + Elastic Agent). Runs standalone or as part of the HyperI Data Fusion Engine. Equivalent in JSON output to that of elastic beats/agent + server ingest painless scripting, only fast and CPU lean.

Topics

Resources

Contributing

Security policy

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

158 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

dfe-transform-elastic

Beats and Elastic Agent JSON in, DFE-normalised events out.

Overview

dfe-transform-elastic is a batch transform service. It consumes JSON event batches produced by filebeat, winlogbeat, metricbeat, auditbeat, heartbeat, packetbeat or Elastic Agent, from Kafka or over a direct gRPC push, applies the Elastic ingest-pipeline logic for that source, and produces normalised events downstream the same way.

The transform logic is compiled in. There is no interpreter, no scripting VM and no plugin system: each supported source is a Rust module, and the service resolves one of them by name at startup. That is the whole design decision, and everything else follows from it -- the container is just the binary, the hot path has no dynamic dispatch per event, and adding a source means a release rather than a config change.

It is built on the scalo data-plane runtime: config cascade, logging, metrics, the Kafka and gRPC transports, health probes, memory guard and scaling pressure.

Sibling services. dfe-transform-vrl runs user-supplied VRL programs and is the right choice when the transform must be changed without a release. This service is the right choice when the transform is a known Elastic pipeline and throughput matters.

Quick start

cargo build --release

./target/release/dfe-transform-elastic --config config.yaml

List the sources a build can transform:

./target/release/dfe-transform-elastic sources

Configuration

Loaded from an explicit --config path, or from scalo's config cascade under the DFE_TRANSFORM_ELASTIC environment prefix.

source:
  name: filebeat.okta.default    # one of `sources` above
  brokers: ["kafka:9092"]
  topics: ["raw_events"]
  group_id: dfe-transform-elastic-my-pipeline
  batch_size: 20000

sink:
  topic: normalised_events
  brokers: ["kafka:9092"]        # defaults to the source brokers
  max_message_bytes: 15728640    # ceiling on one outbound record (15 MiB)

That is the shape, not the whole surface. config.example.yaml is the COMPLETE set of defaults with every key commented, generated from the deployment contract and pinned against drift by a test -- read it rather than this snippet when you need a key that is not here. dfe-transform-elastic emit-config reprints it.

Every transformed event goes out as its OWN record, because dfe-loader parses one JSON document per message. max_message_bytes bounds each one, so keep it below your broker's message.max.bytes.

A config naming a source this build does not carry is rejected at startup, not discovered at the first batch.

Events arrive as NDJSON, one JSON object per line, wrapped by one of three producers: Beats and Elastic Agent (beats), dfe-receiver (receiver) or dfe-fetcher (fetcher). source.envelope defaults to auto, which reads the family off each event, and naming one pins it. The producer's own field names are stripped once they have been lifted onto ECS, except the one dfe-loader routes on -- _source from the receiver and _source_fetcher from the fetcher come through unchanged. What each family carries, and what a payload with no marker does, are in docs/architecture.md.

JSON is the only payload format. MessagePack, supported in DFE/XDR 2.0 and 2.1, is deprecated in DFE 2.2 and no longer accepted: the JSON path (SIMD parsing with sonic-rs, zstd on the wire) is fast enough that MessagePack gave no CPU saving.

Each side also carries a transport: bus, the default, is Kafka, and direct is a scalo Push listener inbound and a gRPC push outbound. Where the values come from, which env spelling reaches a --config file, and which sections that file cannot carry are in docs/configuration.md.

Behaviour under bad input

The service reads whatever is on the topic, so its failure modes are stated rather than assumed:

  • Invalid UTF-8 is decoded with U+FFFD replacements, matching what Beats itself substitutes for a file it cannot decode. The payload survives; the substitution is counted on lossy_payloads_total.
  • Input that is not JSON is refused and dead-lettered with the reason payload is not JSON, and counted on parse_errors_total. A record no line of which parses -- MessagePack, or any other binary -- is refused whole, with the bytes the producer sent. In a record that is otherwise JSON, only the bad line is refused and the events beside it still go. This service configures no DLQ, so the dead letter is dropped and counted on pipeline_dead_letters_dropped_total{reason="dead_letter"}, and its record's source is released.
  • An event whose transform errors is counted on events_errored_total and left out of the output. The batch continues.
  • An event too large for one Kafka record is dropped and counted on events_oversize_total. No broker would accept it, and retrying it forever would block the partition behind it.
  • A sink that refuses for now -- a full producer queue, a broker outage -- is retried for as long as it refuses, counted on send_backpressure_total. The Kafka offset commit, or the answer to a push, is held until the sink takes the block, so an outage is waited out and nothing after the block is committed past it.
  • A send no retry can fix STOPS the service with the block unreleased, and the restarted consumer reads it again. Delivery is at-least-once, so a downstream consumer must be idempotent.
  • A record the sink transport would refuse is dropped before the send and counted on pipeline_dead_letters_dropped_total, never counted delivered.

Non-English text is a tested case, not an edge case. tests/unicode.rs runs every registered transform against seventeen scripts and a set of degenerate inputs.

Deployment

The container image, Helm chart, compose fragment and KEDA scaler are all generated from one deployment contract in src/deployment.rs:

dfe-transform-elastic emit-dockerfile > Dockerfile
dfe-transform-elastic emit-chart chart/dfe-transform-elastic
dfe-transform-elastic emit-compose
dfe-transform-elastic generate-artefacts --output-dir docs

generate-artefacts writes into docs/: metrics-manifest.json, deployment-contract.json, container-manifest.json, Dockerfile.runtime, and the reflectable config pair config-schema.* and capability-catalog.*.

It writes an argocd-application.yaml only when deployment.argocd.repo_url is set in the config cascade, and this repo sets none. scalo writes that Application's source as path: chart, the chart here is chart/dfe-transform-elastic, and scalo has no setting for the path, so the Application would sync nothing. The file is gitignored so a local cascade that sets the URL cannot commit one.

Not every artefact is pinned against a fresh regen. committed_config_artefacts_do_not_drift in src/deployment.rs compares config-schema.{json,yaml} and capability-catalog.{json,yaml}, and its sibling tests do the same for the committed deployment-contract.json, Dockerfile, config.example.yaml and the chart's config: block. Nothing compares container-manifest.json or Dockerfile.runtime, so those can and do fall behind -- re-run the command when src/deployment.rs changes rather than assuming a test caught it.

metrics-manifest.json is compared by metric NAME SET in src/metrics.rs, not byte for byte, because it also records the version and commit of the build that wrote it.

The image expects the release binary in the build context:

cargo build --release
cp target/release/dfe-transform-elastic .
docker build -t dfe-transform-elastic .

The source catalogue, sources.yaml, is compiled into the binary. dfe-transform-elastic emit-catalogue prints it, so a deployment takes the catalogue dfe-engine offers from the image it already pulls, with no credential for this repo.

Observability

  • /livez, /readyz, /metrics and /metrics/manifest, all served from the one listener on port 9090. There is no second health port.
  • /scaling/pressure serves a single weighted figure to KEDA: consumer lag at 0.70, batch saturation at 0.30, with memory as a hard gate that forces pressure to 100 before an OOM.
  • dfe-transform-elastic metrics-manifest prints the full metric catalogue without starting the service. The committed copy is docs/metrics-manifest.json.

Development

cargo test --workspace --all-features
hyperi-ci check

librdkafka 2.12.1 or later is needed for the kafka feature, which is on by default. Building without it still works:

cargo build --no-default-features

docs/ carries everything deeper than this page -- architecture.md for the code map, parity.md for what the service promises against Elastic's own output, and compat.md for working a source towards it.

Licence

BUSL-1.1. See LICENSE, and COMMERCIAL.md for commercial terms.

About

Heavily Rust-optimised transform pipeline for Elastic Stack data (Beats + Elastic Agent). Runs standalone or as part of the HyperI Data Fusion Engine. Equivalent in JSON output to that of elastic beats/agent + server ingest painless scripting, only fast and CPU lean.

Topics

Resources

Contributing

Security policy

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Used by

Contributors

Languages