Running Vector to reshape events costs you a subprocess, its config surface and its failure modes. This embeds the VRL engine instead, so the transform runs in-process and the wrapper owns memory, backpressure and offsets directly.
Embedded VRL (Vector Remap Language) transform engine with wrapper-controlled Kafka source/sink.
dfe-transform-vrl is a Rust binary that runs VRL transforms on Kafka event streams. It embeds the VRL crate directly, giving full control over memory, backpressure, and wire format - without running Vector as a subprocess. It is built on the scalo data-plane runtime (config cascade, logging, metrics, Kafka transport, health probes, scaling).
Why this exists (vs dfe-transform-vector):
| Aspect | dfe-transform-vrl | dfe-transform-vector |
|---|---|---|
| Transform engine | VRL crate (in-process) | Vector subprocess |
| Memory control | Bounded buffers | Vector unbounded |
| Supported transforms | VRL only | All Vector transforms |
| Container image | ~20 MiB | ~170 MiB |
Use dfe-transform-vrl for VRL-only transforms (the common case). Use
dfe-transform-vector when you need Vector-native transforms like lua,
aggregate, dedupe, or throttle.
# Build
cargo build --release
# Run with config file
./target/release/dfe-transform-vrl --config config.yaml
# Run with environment overrides
DFE_TRANSFORM_SOURCE_BROKERS=kafka:9092 \
DFE_TRANSFORM_SOURCE_TOPICS=raw_events \
DFE_TRANSFORM_SINK_TOPIC=enriched_events \
./target/release/dfe-transform-vrlSee config.example.yaml for full configuration reference.
pipeline:
name: "my-pipeline"
batch_size: 1000
source:
brokers: ["kafka:9092"]
topics: ["raw_events"]
group_id: "dfe-transform-vrl-my-pipeline"
transforms:
dir: "/etc/dfe/transforms/"
sink:
brokers: ["kafka:9092"]
topic: "enriched_events"Transform files contain raw VRL source code. Files are loaded in sorted order by filename and executed sequentially against each event:
# transforms/01_parse.vrl
.parsed = parse_json!(.message)
del(.message)
# transforms/02_enrich.vrl
.environment = get_env_var("ENVIRONMENT") ?? "unknown"
.processed_at = now()
pipelines/filebeat/ ships a pre-canned, pure-VRL port of the DFE 2.1
filebeat-compat templates (Cisco Meraki / IOS / Umbrella) plus its
timezones.csv enrichment table. It is a convenience bundle, not engine
capability: opt in by pointing transforms.dir and enrichment_tables at
those files, exactly like any user-supplied transform.
INTERIM: elastic compatibility is being replaced by
dfe-transform-elastic (Rust-native, in beta). See
pipelines/filebeat/README.md for wiring,
routing behaviour, and known limitations.
The big dials have flat env var overrides for K8s. This is the whole list -- anything not here is not read:
| Env Var | Config Field |
|---|---|
DFE_TRANSFORM_PIPELINE_NAME |
pipeline.name |
DFE_TRANSFORM_KAFKA_SASL_USERNAME |
source.sasl.username + sink.sasl.username |
DFE_TRANSFORM_KAFKA_SASL_PASSWORD |
source.sasl.password + sink.sasl.password |
DFE_TRANSFORM_SOURCE_BROKERS |
source.brokers |
DFE_TRANSFORM_SOURCE_TOPICS |
source.topics |
DFE_TRANSFORM_SOURCE_GROUP_ID |
source.group_id |
DFE_TRANSFORM_SOURCE_SASL_USERNAME |
source.sasl.username |
DFE_TRANSFORM_SOURCE_SASL_PASSWORD |
source.sasl.password |
DFE_TRANSFORM_SINK_BROKERS |
sink.brokers |
DFE_TRANSFORM_SINK_TOPIC |
sink.topic |
DFE_TRANSFORM_SINK_KEY_FIELD |
sink.key_field |
DFE_TRANSFORM_SINK_COMPRESSION |
sink.compression |
DFE_TRANSFORM_SINK_SASL_USERNAME |
sink.sasl.username |
DFE_TRANSFORM_SINK_SASL_PASSWORD |
sink.sasl.password |
DFE_TRANSFORM_TRANSFORMS_DIR |
transforms.dir |
The chart mounts the Kafka Secret into the KAFKA_SASL_* pair, which reaches
both endpoints and beats whatever the config file set for either. The
SOURCE_/SINK_ names override it back, per endpoint -- but chart/ injects
only the shared pair, so a two-cluster deployment has to add them to the chart.
Either half turns SASL on, an enabled block still missing one once the env
layer has run refuses to start, and the password is redacted on every output
path (x-scalo-secret + writeOnly in the emitted schema).
tls.ca_cert_file is a path, and the chart mounts no certificate file. Put the
CA's PEM text in the rendered config instead:
config:
source:
tls:
enabled: true
librdkafka_options:
ssl.ca.pem: |
-----BEGIN CERTIFICATE-----
...
-----END CERTIFICATE-----The same two keys go under sink. A CA is public, so a ConfigMap can hold it.
Mutual TLS is not supported through the chart, which mounts no client key. KEDA
scales on CPU and opens no Kafka connection, so it needs no CA.
metrics, logger, scaling, worker_pool, batch_processing,
self_regulation and version_check are resolved by scalo, from a cascade
that discovers files by fixed base name (settings.yaml, defaults.yaml) and
therefore never reads the mounted config.yaml. Writing one of those sections
into that file parses cleanly and changes nothing; the wrapper warns at startup
when it finds one. Set them through the env layer, where the section nests on a
double underscore:
METRICS_ADDR=0.0.0.0:9090 # or DFE_TRANSFORM_METRICS__ADDRESS
LOG_LEVEL=debug # or DFE_TRANSFORM_LOGGER__LEVEL
LOG_FORMAT=json # or DFE_TRANSFORM_LOGGER__FORMAT
DFE_TRANSFORM_SCALING__MEMORY_GATE_THRESHOLD=0.8
DFE_TRANSFORM_BATCH_PROCESSING__MAX_CHUNK_SIZE=10000A single underscore (DFE_TRANSFORM_METRICS_ADDRESS) produces a flat key that
matches no section and is ignored; that spelling also warns.
These parse, validate, and reach nothing. Setting one logs a warning at
startup naming the replacement. See config::INERT_SETTINGS.
| Setting | Why | Use instead |
|---|---|---|
pipeline.batch_size |
the batch engine is built by the scalo runtime before run_service |
batch_processing.max_chunk_size |
pipeline.batch_timeout_ms |
the governed driver has no partial-batch timer | -- |
sink.key_field |
the producer's key argument carries the destination topic, not a partition key (scalo-rs#37) | -- |
source.commit_interval_ms |
auto-commit is off; the engine commits at the at-least-once barrier | -- |
ConfigReloader still re-reads and re-validates the file on a change or a
SIGHUP, so a bad edit is caught, but every field it carries is on that list --
no pipeline behaviour changes on reload today.
Kubernetes liveness probe. Returns 200 OK when the process is running.
Kubernetes readiness probe. Returns 503 if Kafka connections are unhealthy.
Prometheus metrics endpoint.
flowchart LR
SRC["Kafka source<br/>rdkafka consumer"] -->|"JSON"| VRL["VRL engine<br/>in-process, Value in/out"]
VRL -->|"JSON"| SINK["Kafka sink<br/>rdkafka producer"]
SINK -. "delivery confirmed -> commit source offset (at-least-once)" .-> SRC
The wrapper owns both Kafka connections. Consumer offsets are committed only after producer delivery confirmation (at-least-once guarantee).
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.
A record that is not JSON is dropped, counted in records_error_total{stage="deserialise"}, and its source is released Dropped so it is not redelivered.
# Run tests
cargo nextest run
# Run with debug logging
RUST_LOG=debug cargo run -- --config config.yaml
# Check config without running
cargo run -- config-check --config config.yamlchart/ is generated. cargo run --bin dfe-transform-vrl -- emit-chart chart
rewrites it from src/deployment.rs::contract(), so a hand edit under chart/
is reverted the next time anyone regenerates. Fix the contract, not the output.
test_committed_chart_matches_the_generator runs scalo's assert_no_chart_drift, so any chart file that differs from emit-chart fails the suite. A hand fix the generator cannot yet make is pinned there as a ChartPatch, never exempted.
The ScaledObject scales on CPU alone, because contract() sets KafkaLagTrigger::disabled(). Consumer-group lag rises when a downstream stage breaks, and more replicas cannot fix that. The deployed chart in dfe-infra (helm/library/dfe-common/templates/_keda.tpl) renders its own ScaledObject, CPU plus a scaling-pressure trigger through dfe-keda-shim wherever keda.pressure.enabled is set.
- docs/architecture.md - What the service is, why it is shaped that way, and the invariants
- docs/DESIGN.md - Deeper design detail: memory budget, reload matrix
- config.example.yaml - Configuration reference
This project is licensed under the Business Source License 1.1 (BUSL-1.1). See LICENSE for details.
Copyright (c) 2026 HYPERI PTY LIMITED
For commercial licensing options, see COMMERCIAL.md.
A single Rust binary that runs VRL transforms over records in the DFE data path
-- Kafka in and out on the bus transport, gRPC in and out on the direct
transport. Two boundaries get assumed wrongly. First, it is VRL only: a pipeline
needing lua, aggregate, dedupe, throttle or sample belongs in
dfe-transform-vector, which keeps the Vector subprocess and pays for it. Second,
this crate does not own its own event loop. scalo's BatchEngine::pipeline
drives recv -> process -> send -> release; this crate supplies the process
closure, the produce sink, the config, the VRL compiler and the enrichment
registry. Reading src/pipeline.rs expecting to find the loop is the usual wrong
turn.
| Path | Holds |
|---|---|
src/pipeline.rs |
The process closure and sink handed to scalo's batch engine |
src/engine/ |
compiler.rs builds one program from the transform files, runner.rs runs it per event, budget.rs holds the compile memory floor |
src/config/ |
Config, validation, HotConfig, INERT_SETTINGS, SCALO_CASCADE_SECTIONS |
src/enrichment/ |
Table loading (CSV, JSON, YAML, MMDB, STIX, SQLite), refresh, the custom VRL functions |
src/deployment.rs |
contract() -- the single source the Dockerfile and chart/ are generated from |
chart/, Dockerfile |
Generator output, not hand-authored |
pipelines/filebeat/ |
Opt-in data bundle (212,218 bytes of VRL plus a lookup table), not engine capability |
tests/ |
integration/ runs mostly without Kafka, e2e/ needs a broker, TESTING.md explains the modes |
docs/architecture.md |
Why the service is shaped this way, and the invariants |
make check # hyperi-ci check -- quality + test, the pre-push gate
cargo nextest run --features enrichment-mmdb,enrichment-sqlite # what CI actually runs
cargo nextest run --all-features --run-ignored # adds the broker-dependent e2e tests
cargo run -- config-check --config config.yaml # validate a config without starting
cargo run --bin dfe-transform-vrl -- emit-chart chart # regenerate the chartGreen lies here in three ways, and all three are on by default.
default = [] in Cargo.toml, so a bare cargo nextest run does not compile the MMDB or SQLite enrichment tests in at all. .hyperi-ci.yaml adds enrichment-mmdb and enrichment-sqlite to both the test and the build feature sets -- the build too, because the Dockerfile only copies the binary, so a default build ships a container that rejects type: mmdb and type: sqlite tables at startup.
Kafka-dependent tests call skip_if_no_kafka!() and skip cleanly when no broker
is reachable. They report as not-failed, which is not the same as proven.
The e2e tests are #[ignore] and need --run-ignored plus a broker. Set
TEST_MODE=docker for the dfe-docker infra profile on localhost:19092, or
TEST_MODE=remote for the cluster endpoints.
One more: the push trigger in .github/workflows/ci.yml carries
paths-ignore: docs/**, **.md, so a docs-only push runs no jobs. The
pull_request trigger has no such filter, so the PR is where a docs change gets
checked.
| Don't | Do | Why |
|---|---|---|
Hand-edit chart/ or Dockerfile |
Fix src/deployment.rs::contract() and regenerate |
Both are generator output and a hand edit is reverted by the next regeneration. The chart once mounted the Kafka SASL Secret into env names nothing read, so credentials reached the pod and were ignored -- fixed in the generator so emit-chart keeps it |
Turn the Kafka lag trigger back on in contract() |
Leave KafkaLagTrigger::disabled() and scale on CPU plus scaling pressure |
Lag rises when a downstream stage breaks, so a lag trigger adds pods that wait on the same broken stage. The lag trigger was also the one that kept reading a kafka block this app's values do not have, fixed three times (#37, #65, #70) |
Put a scalo section (metrics, logger, scaling, worker_pool, batch_processing, self_regulation, version_check) in the config file |
Set it through the env layer, on a double underscore | scalo's cascade finds files by fixed base name and can never be pointed at config.yaml, so the section parses and reaches nothing. Full list above under Configuration |
Trust pipeline.batch_size, pipeline.batch_timeout_ms, sink.key_field or source.commit_interval_ms |
Size a chunk with batch_processing.max_chunk_size |
config::INERT_SETTINGS -- accepted, validated, reaching nothing. Table above under Configuration |
Call a bare cargo nextest run green |
Pass --features enrichment-mmdb,enrichment-sqlite |
default = [], so the MMDB and SQLite tests are not compiled in and the run is green without having tested them |
Set sasl.enabled with an empty username or password |
Supply both, or neither | librdkafka's SCRAM check is a NULL check that an empty string passes, so that pod authenticated against nothing and still reported Ready. Refused at startup now |
| Add a serialise path that redacts by field name | Keep the password a scalo::SensitiveString |
Redaction is by type on every path -- Debug, the /config dump, the emitted schema. The figment round-trip has to be wrapped in expose_during or a file-sourced password reaches the broker as the literal ***REDACTED*** |
| Size the container from steady state when the program is large | Leave headroom for the compile | Compilation scales with the VRL, and the bundled filebeat program is OOM-killed under a 32 MiB limit. Below the floor the kernel kills the process mid-compile, seen as an exit-137 restart loop with nothing in the log. engine::budget refuses first and names all three numbers |
Pick the optimisation tier with publish-target |
Use build.skip_optimize in .hyperi-ci.yaml, or the skip-optimize dispatch input for one run |
hyperi-ci ignores publish-target. A build that ships gets PGO and BOLT from scripts/pgo-workload.sh unless optimisation is skipped, and a stable release that skips it is refused unless the dispatch passes release-unoptimized: true |
Inbound -- what this repo depends on:
- scalo-rs (
cargo-dep).Cargo.tomldeclares thescalocrate by range, a second range covering the dev dependency. A scalo release arrives through that range:cargo update -p scaloand rebuild if it admits the version, widen the range first if not. - scalo-rs (
generated-file, lockstep). TheDockerfileand everything underchart/are written by scalo's generators from this crate'sdeployment::contract(). A generator or schema change upstream means regenerating with the command in the file's own header and committing the diff. - dfe-infra (
apps.yaml, deploy-time authority). The suite manifest declares what this app is:multiplicity: per_config(one deployment per source config, never a singleton),scale_deployed: true, both transports, and a source binding deriving{source}_landin,{source}_loadout anddfe-transform-vrl-{source}as the consumer group. Adding a source is a manifest edit, not a change here.
Outbound -- what depends on this repo:
- dfe-infra (
image-pin, lockstep). Pins this repo's ghcr image as a tag plus the digest that makes it immutable. A release here means bumping the tag and re-resolving the digest there, withcheck_versions_drift.pyconfirming the chart's appVersion and the digest mirror agree.
dfe-transform-splack runs this repo's engine image with a rule-config chart of its own. It is out of suite scope and no work here is driven by it.