Transformer Simulation
Overview
SDL supports isolated execution of individual Transformers against supplied sample inputs. This allows users to test supported transformer code before deploying a Pipeline. This is particularly useful for Dynamic Transformers, which allow users to provide custom code.
This page defines two related contracts:
-
the authenticated client API used by the Test Code workflow; and
-
the runtime adapter API implemented by Kubernetes-based Transformer images.
The Data Pipeline Engine owns the boundary between these contracts. Client applications do not interpret Kubernetes or runtime-specific responses directly.
Limitations
The supported simulation connection variants are INTERNAL_KAFKA, INTERNAL_POSTGRES, and INTERNAL_ICEBERG. Simulation is limited to Transformer Templates that define only supported input and output types. The Data Pipeline Engine validates these constraints when a Transformer Template is registered.
Dynamic Transformer Test Code supports JavaScript Kafka JSON and a capability-gated Python Kafka JSON profile. Python MinIO simulation and package installation are not supported. A Python template must not advertise simulation until its configured runtime image implements the Python profile in this document.
Enabling Simulation for a Transformer Template
Transformer Templates that support isolated execution must include the simulation-enabled label in their definition:
{
"uid": "1afc0a55-ca7e-40db-9efc-971f85f48b42",
"labels": [
{"name": "simulation-enabled"}
]
// Rest of definition
}
The label declares a capability only. It does not select simulation mode for a normal deployed Pipeline stage. Only the Data Pipeline Engine’s isolated simulator may set TRANSFORMER_SIMULATION_MODE=true. A live stage must start normally and receive connection-derived runtime configuration.
The Python Kafka JSON profile also requires the simulation-python-kafka-json label. This profile label allows the backend and clients to apply Python-specific validation without inferring behavior from the template UID, display name, editor language, or file extension.
Authenticated Client API
The Test Code workflow calls this endpoint:
POST /v1/pipelines/transformer/uid/<uid>/simulate
Content-Type: application/json
The caller must be authenticated and must have read access to the Transformer Template identified by <uid>. The template must advertise simulation-enabled. The endpoint does not persist changes to a Pipeline or Transformer Template.
Client Request
The request contains template-approved configuration and named sample inputs:
{
"configuration": {
"environment": {},
"file_configuration": {
"script": {
"base64_encoded_text": "Y29uc3QgaGFuZGxlciA9IChpbnB1dCkgPT4gaW5wdXQ7"
}
}
},
"inputs": {
"SOURCE_TOPIC": {
"msgs": [
{
"key": "sample-1",
"headers": {},
"value": "{\"name\":\"example\"}"
}
]
}
}
}
Request rules:
-
configurationandinputsare required objects. -
environmentvalues and file names must be allowed by the Transformer Template. -
base64_encoded_textcontains standard base64-encoded file content. -
replicasis ignored for isolated execution and clients should omit it. -
Each input key must match a template input connection or valid variadic expansion.
-
Each input contains exactly one data variant:
msgsfor Kafka, ortablefor Iceberg/Postgres. -
Kafka message
valueis a string. Supported Dynamic Transformer profiles expect that string to contain valid JSON.
Dynamic Transformer Python Profile
The simulation-python-kafka-json profile uses the existing Python template contract:
-
required
scriptfile at/config/script.py; -
optional
requirementsfile at/config/requirements.txt; -
SOURCE_KAFKA.msgsinput; and -
DEST_KAFKA.msgsoutput.
The script must export callable transform(data). For each message, the runtime decodes the JSON object in value, calls transform, and serializes a returned dictionary. Returning None skips the output message. The runtime preserves supported keys and string headers.
The effective requirements content may be blank or comment-only. Any non-comment entry is rejected before Kubernetes simulator resources are created because the isolated workload has no egress and uses a read-only root filesystem. Imports may use the standard library and packages already present in the image. A Python simulation request containing MinIO or another unsupported input is also rejected before resource creation.
Simulation mode loads transform but does not call user main(), initialize Kafka or MinIO clients, or install packages. Live mode continues to call main(). Parity requires the live main() function to pass the same transform callable to start_transform_loop.
Client Response
When the isolated runtime returns a valid response, the Data Pipeline Engine returns HTTP 200 with a client-facing envelope:
{
"status_code": 200,
"metadata": {
"simulation_time_ms": 12
},
"errors": null,
"logs": [
{"msg": "Processed 1 message"}
],
"outputs": {
"DEST_TOPIC": {
"msgs": [
{
"key": "sample-1",
"headers": {},
"value": "{\"name\":\"example\",\"processed\":true}"
}
]
}
}
}
The outer HTTP status describes client-to-Data-Pipeline-Engine handling and orchestration. Body status_code is the HTTP status returned by the isolated runtime. A structured runtime 400 or 500 is therefore returned in an outer HTTP 200 response so the client can display runtime logs and errors.
| Outer HTTP | Body status_code |
Meaning | Client state |
|---|---|---|---|
|
|
Runtime executed successfully. |
Succeeded |
|
|
Runtime rejected the supplied sample or user configuration. |
Runtime validation failed |
|
|
User script or runtime execution failed after startup. |
Runtime failed |
|
Not present |
The request is malformed or invalid, or the template is not simulation-capable. |
Request rejected or unsupported |
|
Not present |
Authentication is missing/invalid, or template read access is denied. |
Unauthorized or forbidden |
|
Not present |
The Transformer Template does not exist. |
Not found |
|
Not present |
The request or runtime response exceeded the configured size or record limit. |
Resource limit exceeded |
|
Not present |
The configured number of concurrent tests is already running. |
Capacity reached; retry later |
|
Not present |
The simulator did not start or finish within its configured time limit. |
Timed out |
|
Not present |
Simulator startup, cleanup, or another platform operation failed before a valid runtime response was available. |
Platform failed |
Non-200 Data Pipeline Engine responses use the Pipelines API error envelope:
{
"error_code": 400,
"error_messages": ["simulation request is not well-formed"]
}
Runtime Adapter API - Kubernetes-based Transformers
This is the internal adapter contract that the Data Pipeline Engine expects each Kubernetes-based Transformer Template to implement when it declares simulation-enabled.
Note that this contract applies only to Transformer Templates that are specified as containerized images to be run in Kubernetes; other instantiation types will require separate simulation APIs.
Assumptions/Prerequisites
-
Simulation cannot occur with Transformer Templates that have any
inputconnections of typeTRANSFORMER, because the method to supply inputs to those connections is unknown by the Data Pipeline Engine. -
A Transformer Template should support running either as a Transformer directly or in simulation mode without rebuilding or changing the image.
Detecting Simulation
Transformer Templates should read the value of the environment variable named TRANSFORMER_SIMULATION_MODE to determine whether or not to start in Simulation (FaaS) or Transformer mode:
-
If the environment variable exists and has the value of
true, then the Transformer should start up in Simulation mode. -
In any other case, the Transformer should start normally.
Startup
If running in Simulation mode, a Transformer should not run an event loop or listener, nor should it expect that any configuration values to connect to data sources are present. For example, the configuration contract for INTERNAL_KAFKA connections will not be fulfilled in Simulation mode. However, any user-provided configuration values (as specified in the Transformer Template’s configuration block) will be filled in by the Data Pipeline Engine.
Instead, the Transformer should launch an HTTP server on port 8111 that contains handlers for two routes:
-
/health: This route should return200 OKonce the Transformer is finished startup and is ready to accept the simulation request. TheGETHTTP verb should be handled. -
/simulate: This route should handle the simulation request coming from the Data Pipeline Engine via an HTTPPOSTrequest, with the JSON-encoded body containing the simulation request itself.
Handling /simulate
The /simulate route should be handled in a very specific way in order to be as efficient as possible. Transformers should anticipate that only a single request to /simulate will be made, and that after returning an HTTP response, the Transformer should shut down. Note that repeated simulation requests from users for the same Transformer Template will result in additional replicas being instantiated, not that the same Transformer will be used to handle more than one request. While this may impact performance slightly, it ensures that repeated simulation requests are stateless.
Request Body
The request body for /simulate is very specific, and is always encoded as application/json:
{
"inputs": {
// Each item corresponds to the name of an input connection in the Transformer Template
"CONNECTION_1": {
// Note that conn_type is elided here, as it isn't important
// Assume for this example that CONNECTION_1 is of type INTERNAL_KAFKA.
"msgs": [
// Each object here should be fed into the Transformer code as though it were coming from Kafka.
{"key": "key1", "headers": {}, "value": "Hello, world!"},
{"key": "key2", "headers": {"__from__": "value1"}, "value": "Hello, world!"},
{"key": "key3", "headers": {}, "value": "Hello, world!"},
]
},
"CONNECTION_2": {
// Assume for this example that CONNECTION_1 is of type INTERNAL_ICEBERG or INTERNAL_POSTGRES.
"table": {
// This is a JSON representation of a table that should be created in the Transformer as input.
"cols": ["col1", "col2", "col3"],
"rows": [
[1, "one", false],
[2, "two", true],
[3, "three", false]
]
}
}
}
}
Upon receipt of this request body, the Transformer should route these inputs into its core transformation logic as though they were coming directly from the connected data sources.
Response
Transformers should run a simulation request as similarly to normal operation as possible. The core transformation code path should be shared by simulated and live execution.
Note that normal logging to stdout or stderr will not be captured by the Data Pipeline Engine when returning results to the client. Any logs the Transformer wants to return must be encoded in the response body, as outlined below.
Once a result has been obtained, the Transformer should return a response to the Data Pipeline Engine containing one of these HTTP status codes:
200 OK
Simulation succeeded. The Transformer should return the simulation results in a JSON response body matching this format:
{
// Additional arbitrary values or objects may be returned by the Transformer in this block. Keys beginning
// with a double-underscore (__) are private and must be removed before the client response.
"metadata": {
"simulation_time_ms": 26, // Will be returned to the client.
"__trace_info": { // Internal only; must be removed before the client response.
"id": 273465,
}
},
// `logs` contains zero or more log messages, as output by the Transformer.
"logs": [
{"msg": "Hello, world!"},
{"msg": "Second log"},
{"msg": "Third log"}
],
// `outputs` should correspond to the list of output connections defined by the Transformer Template.
// Note that the data model returned in each output (table or msgs) is validated by the Data Pipeline Engine
// against the conn_type of the output; a 500 will be returned if the returned data do not match the expected
// conn_type.
"outputs": {
"ICEBERG_TABLE_OUTPUT_CONN": {
"table": {
"cols": ["output1", "output2"],
"rows": [
[2, "asdf"],
[3, "asdf"],
]
}
},
"KAFKA_OUTPUT_CONN": {
"msgs": [
{"key": "key1", "headers": {}, "value": "Hello, world!"},
{"key": "key2", "headers": {"__from__": "value1"}, "value": "Hello, world!"},
{"key": "key3", "headers": {}, "value": "Hello, world!"},
]
},
}
}
400 Bad Request
The Transformer failed to run the simulation request due to issues with the data. - For example, the input Kafka messages were raw JSON when the Transformer expected raw XML.
The Transformer should return the error in a JSON response body matching this format:
{
// Additional arbitrary values or objects may be returned by the Transformer in this block. Keys beginning
// with a double-underscore (__) are private and must be removed before the client response.
"metadata": {
"simulation_time_ms": 26, // Will be returned to the client.
"__trace_info": { // Internal only; must be removed before the client response.
"id": 273465,
}
},
// `logs` contains zero or more log messages, as output by the Transformer.
"logs": [
{"msg": "Hello, world!"},
{"msg": "Second log"},
{"msg": "Third log"}
],
"errors": [
{"msg": "Connection XYZ must contain XML data"}
]
}
500 Internal Server Error
The Transformer failed to run the simulation request for some internal reason. Note that this will be returned to the client as a server-side error, rather than an error on the client’s part.
The Transformer should return the error in a JSON response body matching this format:
{
// Additional arbitrary values or objects may be returned by the Transformer in this block. Keys beginning
// with a double-underscore (__) are private and must be removed before the client response.
"metadata": {
"simulation_time_ms": 26, // Will be returned to the client.
"__trace_info": { // Internal only; must be removed before the client response.
"id": 273465,
}
},
// `logs` contains zero or more log messages, as output by the Transformer.
"logs": [
{"msg": "Hello, world!"},
{"msg": "Second log"},
{"msg": "Third log"}
],
"errors": [
{"msg": "Connection XYZ must contain XML data"}
]
}
Shutdown
As mentioned above, once a request to /simulate has been handled and responded to, the Transformer is obligated to shut down as quickly as possible, to ensure that simulation requests are as efficient as possible.
Runtime-to-Client Mapping
Runtime and client responses use the canonical field name metadata. The Data Pipeline Engine still maps the runtime DTO to a client DTO so runtime-specific behavior does not leak into the browser contract. During a rolling upgrade, the transport adapters may accept the legacy field name meta, but they must normalize it to metadata and must not expose meta in application models. Runtime images must not return the client envelope or set status_code themselves; the backend adds status_code from the runtime HTTP response.
Runtime metadata keys beginning with __ are reserved for internal use. The safety implementation must remove reserved or sensitive metadata before returning metadata to a browser.
Runtime Parity
Isolated and live Dynamic Transformer execution must load the same user script and invoke the same core transformation function for each Kafka message value. JavaScript uses handler; Python uses transform, and live Python code passes that function to start_transform_loop. Given the same script, approved configuration, and ordered input values, transformed values must be semantically equivalent.
Intentional differences include supplied sample data instead of a live connection, one-request HTTP transport, isolated Kubernetes resources, execution limits, and the absence of live offsets/partitions. Timestamps, generated identifiers, batching, and infrastructure logs may also differ when they are not part of the transformer’s business result.
Safety Boundary
The v1 authorization policy requires an authenticated caller with read access to the Transformer Template. Authorization and capability checks complete before the Data Pipeline Engine creates simulator resources. The endpoint does not save the script, sample, or Pipeline state.
Each test runs in a separate Kubernetes Job. The Job uses the df-transformer-simulator ServiceAccount, does not mount a Kubernetes API token, and has no RoleBinding. The Pod and container require non-root execution and RuntimeDefault seccomp. The container cannot escalate privileges, uses a read-only root filesystem, and drops all Linux capabilities.
A per-test NetworkPolicy selects only that simulator Pod. It permits ingress to TCP 8111 from df-backend Pods and denies egress. The simulation Job receives no connection-derived Kafka, object-storage, catalog, or external-egress labels and receives no live data-service credentials.
Default Limits
| Control | Default | Enforcement |
|---|---|---|
Concurrent tests |
|
Additional requests receive HTTP |
Startup timeout |
|
The backend stops waiting and returns HTTP |
Execution timeout |
|
The backend cancels the request, deletes resources, and returns HTTP |
Job deadline and finished-Job TTL |
|
Kubernetes stops unbounded work and removes a completed Job if normal cleanup is interrupted. |
CPU request / limit |
|
Applied to the simulator container. |
Memory request / limit |
|
Applied to the simulator container. |
Client request and runtime response |
|
Oversized data receives HTTP |
Sample and output records |
|
The runtime rejects oversized samples or outputs. |
Browser-visible runtime logs |
|
The runtime and backend truncate the log stream before returning it. |
These values are set under pipeline.simulation in the df-backend chart and passed through PIPELINE_SIMULATION_K8S_* environment variables.
Redaction Rules
Backend and runtime logs record counts, status codes, generated simulator names, and lifecycle events. They must not record raw requests, scripts, file contents, sample values, message keys, headers, topic names, environment values, credentials, tokens, passwords, API keys, cookies, or authorization headers. Runtime errors are mapped to bounded user actions instead of forwarding V8, HTTP, or Kubernetes error text.
The outputs object is the requested transformation result and may contain values derived from the supplied sample. It is shown only to the authorized caller and is not written to backend/runtime logs or retained as evidence by default. Reviewers must use synthetic sentinel data in screenshots and shared evidence.
Cleanup uses a bounded context that is independent of client cancellation and always attempts to delete the Job, Secret, Service, and NetworkPolicy. Kubernetes deadline and TTL controls are the fallback when immediate deletion fails.
Reviewer Checks
-
Confirm a simulator Job names
df-transformer-simulator, setsautomountServiceAccountToken: false, and has the documented Pod/container security contexts and resources. -
Confirm the Job has
activeDeadlineSeconds: 60,ttlSecondsAfterFinished: 60, andbackoffLimit: 0. -
Confirm the matching NetworkPolicy has no egress rules and admits only
df-backendon TCP8111. -
Run success, runtime-error, timeout, oversize, and capacity cases with synthetic sentinel values; search backend, runtime, browser, and retained evidence for those values.
-
Confirm the Job, Secret, Service, and NetworkPolicy are absent after success, failure, cancellation, and timeout.
Recent Payload Decision
Recent-payload selection is not part of the first contract. Users provide sample JSON manually or copy it from an authorized preview. Reading or retaining recent stream payloads requires separate privacy, retention, authorization, classification, and storage decisions.
Functional Flow
When the Data Pipeline Engine receives a simulation request from a client, it requires the following information to be present:
-
The UID of the Transformer Template to execute the simulation request against. This tells the Data Pipeline Engine how to instantiate the app that will receive the simulation request.
-
Configuration for each
inputconnection on the Transformer Template, which will be passed to the app at runtime. This will usually take the form of discrete data, as opposed to a connection to a Dataset. -
General Transformer configuration (environment variables, file configurations, etc.)-effectively the same as what configuration is passed to a Transformer when it is used in a Pipeline.
Based on these pieces of information, the Data Pipeline Engine performs the following steps:
-
It retrieves the Transformer Template with the matching UID from the persistence layer.
-
It validates the configuration and input blocks from the simulation request against the Transformer Template.
-
If successful, it instantiates the given Transformer Template as a run-once task.
-
Once the Transformer is active, it forwards the simulation request to it and waits for a response.
-
Once it has received a response, it shuts the Transformer down (if it has not shut itself down already).
-
It filters runtime
metadata, adds the runtimestatus_code, and returns the client response using the same canonicalmetadatafield name.