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:

  • configuration and inputs are required objects.

  • environment values and file names must be allowed by the Transformer Template.

  • base64_encoded_text contains standard base64-encoded file content.

  • replicas is 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: msgs for Kafka, or table for Iceberg/Postgres.

  • Kafka message value is 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 script file at /config/script.py;

  • optional requirements file at /config/requirements.txt;

  • SOURCE_KAFKA.msgs input; and

  • DEST_KAFKA.msgs output.

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

200

200

Runtime executed successfully.

Succeeded

200

400

Runtime rejected the supplied sample or user configuration.

Runtime validation failed

200

500

User script or runtime execution failed after startup.

Runtime failed

400

Not present

The request is malformed or invalid, or the template is not simulation-capable.

Request rejected or unsupported

401 or 403

Not present

Authentication is missing/invalid, or template read access is denied.

Unauthorized or forbidden

404

Not present

The Transformer Template does not exist.

Not found

413

Not present

The request or runtime response exceeded the configured size or record limit.

Resource limit exceeded

429

Not present

The configured number of concurrent tests is already running.

Capacity reached; retry later

504

Not present

The simulator did not start or finish within its configured time limit.

Timed out

500

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

  1. Simulation cannot occur with Transformer Templates that have any input connections of type TRANSFORMER, because the method to supply inputs to those connections is unknown by the Data Pipeline Engine.

  2. 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:

  1. /health: This route should return 200 OK once the Transformer is finished startup and is ready to accept the simulation request. The GET HTTP verb should be handled.

  2. /simulate: This route should handle the simulation request coming from the Data Pipeline Engine via an HTTP POST request, 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

4 per backend replica

Additional requests receive HTTP 429 before resource creation.

Startup timeout

30s

The backend stops waiting and returns HTTP 504.

Execution timeout

30s

The backend cancels the request, deletes resources, and returns HTTP 504.

Job deadline and finished-Job TTL

60s / 60s

Kubernetes stops unbounded work and removes a completed Job if normal cleanup is interrupted.

CPU request / limit

100m / 500m

Applied to the simulator container.

Memory request / limit

128Mi / 256Mi

Applied to the simulator container.

Client request and runtime response

1 MiB each

Oversized data receives HTTP 413 and is not forwarded to the browser.

Sample and output records

100 per connection

The runtime rejects oversized samples or outputs.

Browser-visible runtime logs

100 entries, 512 bytes per entry, 32 KiB total

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, sets automountServiceAccountToken: false, and has the documented Pod/container security contexts and resources.

  • Confirm the Job has activeDeadlineSeconds: 60, ttlSecondsAfterFinished: 60, and backoffLimit: 0.

  • Confirm the matching NetworkPolicy has no egress rules and admits only df-backend on TCP 8111.

  • 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:

  1. 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.

  2. Configuration for each input connection 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.

  3. 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:

  1. It retrieves the Transformer Template with the matching UID from the persistence layer.

  2. It validates the configuration and input blocks from the simulation request against the Transformer Template.

  3. If successful, it instantiates the given Transformer Template as a run-once task.

  4. Once the Transformer is active, it forwards the simulation request to it and waits for a response.

  5. Once it has received a response, it shuts the Transformer down (if it has not shut itself down already).

  6. It filters runtime metadata, adds the runtime status_code, and returns the client response using the same canonical metadata field name.