Imported from goldsky-io/streamling (
AGENTS.md). Install upstream withnpx skills add goldsky-io/streamling. Copyright stays with the author.
Streamling Development Guide
Streamling is a Rust-based data processing engine built on Arrow and DataFusion. It reads from sources (Kafka), applies transforms (SQL, HTTP, WASM), and writes to sinks (Postgres, ClickHouse, Kafka, SQS, webhook, etc.).
Project Layout
crates/
├── streamling/ # Main binary
├── streamling-core/ # Core topology, operators, pipeline engine
├── streamling-config/ # Pipeline YAML parsing and configuration
├── streamling-connectors/ # Table providers (sources/sinks)
├── streamling-plugin/ # FFI-based plugin system
├── streamling-state/ # State backend (checkpointing)
├── streamling-flink-compat/ # Flink state compatibility
└── streamling-e2e/ # End-to-end test framework (see crates/streamling-e2e/AGENT.md)
plugin_examples/ # Example plugins (basic, low_level)
scripts/ # k3s setup/teardown scripts
Command Runner (just)
All common tasks use just. Run just --list to see everything.
| Command | Purpose |
|---|---|
just build |
Build the workspace |
just fix |
Auto-fix cargo issues |
just lint |
Format + clippy for the main workspace |
just test |
Run unit tests |
just env-setup |
Spin up k3s with Postgres, Kafka, ClickHouse, etc. |
just env-status |
Check k3s cluster health |
just env-teardown |
Tear down k3s cluster |
just e2e-test |
Build + run all e2e tests |
just e2e-test-debug |
E2e tests with full output |
just e2e-list |
List available e2e tests |
Development Workflow
After every task
- Run
just fixto auto-fix warnings and issues. - Run
just lintto verify formatting and clippy pass.
Fixing a bug
- Write a unit test that reproduces the issue first.
- Watch it fail.
- Implement the fix — keep the change as small and easy to understand as possible.
- Verify the test passes.
- Run
just fix && just lint.
Adding a feature
- Prefer unit tests for quick iteration during development.
- If the feature involves pipeline-level behavior (source -> transform -> sink), add an e2e test once the implementation is complete.
- Run
just fix && just lint.
Before pushing
Check that CI/CD has passed for the PR:
gh pr checks 494 # use the relevant PR number
Testing Strategy
Unit tests (preferred for iteration)
just test # Run all unit tests
cargo test -p streamling-core # Run tests for a specific crate
Unit tests are fast and should be the first choice for verifying logic changes.
E2E tests (always run locally)
E2E tests exercise the full streamling binary as a black box against real infrastructure (Kafka, Postgres, ClickHouse, SQS via ElasticMQ) running in a local k3s cluster.
E2E tests must always be created and run locally before pushing. They are the final verification that a pipeline-level change works end-to-end.
# First time: set up the local k3s environment
just env-setup
# Run all e2e tests
just e2e-test
# Run a specific test
just e2e-test test_basic_postgres_sink
# Run with full debug output
just e2e-test-debug
# Use debug build (faster compile, slower runtime)
PROFILE=debug just e2e-test
See crates/streamling-e2e/AGENT.md for detailed guidance on writing e2e tests (test patterns, resource isolation, available helpers).
Pipeline YAML
Pipelines are defined in YAML with three sections: sources, transforms, and sinks. The transforms key is always required, even if empty (transforms: {}).
sources:
my_source:
type: kafka
topic: my_topic
starting_offsets: earliest
primary_key: id
telemetry:
event_time:
column: block_timestamp
unit: seconds
labels:
tier: critical
dataset: v2.evm.blocks
transforms: {}
sinks:
my_sink:
type: postgres
from: my_source
table: output_table
schema: public
primary_key: id
on_conflict: update
telemetry block (optional)
Every source, transform, and sink accepts an optional telemetry block:
event_time— column that carries the record's event time; drivesstreamling_event_time_watermark_millisecondsandstreamling_event_time_lag_millisecondsseries.unitrequired for integer columns (seconds,milliseconds,microseconds).labels— map ofkey: valueidentity labels attached to every metric this node emits. Use for dataset identity, tier/team tagging, destination tagging.
Label constraints (enforced at config load):
| Rule | Reason |
|---|---|
| Max 20 labels per node | Each label is a Prometheus dimension; unbounded maps blow up storage. |
Keys match ^[a-zA-Z_][a-zA-Z0-9_]*$ |
Prometheus label-name grammar. |
Keys cannot start with __ |
Prometheus reserves that prefix for internal use. |
Keys cannot be id, topology_node_type, operator_type, service_instance_id |
These collide with built-in metric dimensions. |
| Keys cannot shadow a node kind's per-type tag | topic on Kafka, table on ClickHouse/Postgres, url on webhooks, type on plugins, language on script transforms. Hybrid sources reserve both table and topic. Overriding would silently replace the real identity on every emitted metric. |
| Values max 256 bytes, no control characters (except tab) | A raw newline in a value can corrupt Prometheus scrape output. |
Precedence when labels collide with plugin-declared labels: plugin wins, and a WARN log names the colliding key and both values. Plugins are authoritative about their own identity (chain_slug, topic, etc.); YAML cannot override them, but the WARN surfaces misconfigurations rather than letting them go silent.
Hybrid sources: telemetry.labels on a hybrid source automatically propagate to both the bounded and unbounded phase-child metric series — declare them once on the parent.
Key Architecture Concepts
- Arrow RecordBatches flow between operators as the internal data format.
- DataFusion provides the SQL engine for transforms and table providers for sources/sinks.
- Checkpointing persists Kafka offsets and operator state for exactly-once semantics.
- Plugin system uses FFI (
cdylib+abi_stable) to load external sources, sinks, and transforms at runtime.
Common Pitfalls
- Tests hang: Add
batch_size: 1to sink config andenv("STREAMLING__RECORD_BATCH_SIZE", "1")to pipeline opts. - k3s not running: Run
just env-statusto check,just env-setupto start. - Linker issues with nextest on macOS: Set
export DYLD_LIBRARY_PATH="$HOME/.rustup/toolchains/1.89.0-aarch64-apple-darwin/lib/rustlib/aarch64-apple-darwin/lib:$DYLD_LIBRARY_PATH".
Code Review Rules
Each rule has a severity: Block (must fix before merge), Recommend (should fix), Minor (nit).
Convention Rules
| Rule ID | Severity | Rule |
|---|---|---|
| CONV-001 | Block | Bug fix PRs MUST include a failing unit test that reproduces the issue |
| CONV-002 | Block | just fix && just lint must pass (no clippy warnings with -D warnings) |
| CONV-004 | Recommend | Pipeline-level features should have an e2e test |
| CONV-005 | Recommend | Changes should be minimal and easy to understand — prefer small, focused PRs |
| CONV-006 | Minor | Commit messages should use conventional format (feat:, fix:, chore:, etc.) |
Rust Rules
| Rule ID | Severity | Rule |
|---|---|---|
| RUST-001 | Block | No .unwrap() in non-test production code — use ?, .expect("reason"), or handle |
| RUST-002 | Block | unsafe blocks must have a // SAFETY: comment explaining soundness |
| RUST-003 | Recommend | Prefer StreamlingError over anyhow::Error for public APIs |
| RUST-004 | Recommend | New public APIs should have doc comments |
| RUST-005 | Minor | Avoid clone() where a reference or Cow would suffice |
Connector Rules (kafka, clickhouse, postgres)
| Rule ID | Severity | Rule |
|---|---|---|
| CONN-001 | Block | Connection errors must be retriable (StreamlingError::retriable) |
| CONN-002 | Recommend | New sink configs should use #[serde(deny_unknown_fields)] |
| CONN-003 | Recommend | Resource cleanup (connections, cursors) must happen in Drop impls or explicit close |
Pipeline / Topology Rules
| Rule ID | Severity | Rule |
|---|---|---|
| PIPE-001 | Block | Arrow schema changes must be backward-compatible or versioned |
| PIPE-002 | Recommend | New transforms must handle empty RecordBatches gracefully |
| PIPE-003 | Recommend | Checkpoint-related changes must preserve exactly-once semantics |