FlowDSL Node Catalog

JSON Stream Source

Ingest arbitrary JSON messages from an external stream — HTTP webhook, NDJSON file, or message bus.

Overview

PropertyValue
Node IDredelay/json-stream-source
Kindsource
Moduleevents-flowdsl
Repogo-flowdsl
Module refevents
Handlergithub.com/redelay/go-flowdsl/nodes.JSONStreamSource
Iconlucide:file-json
Color#0ea5e9
Tagsjson, ingest, source, webhook, stream

Description

Subscribes to an external JSON message stream and emits each decoded document as a packet. The node does NOT interpret the payload; pair it with redelay/json-to-event to map incoming fields to a Redelay event.

Supported transports:

  • http-webhook — exposes an HTTP POST endpoint at the configured path
  • ndjson — tails an NDJSON (newline-delimited) file or URL
  • kafka-raw — reads a Kafka topic where the value is raw JSON (no envelope)
  • nats — subscribes to a NATS subject (core or JetStream)
  • redis-stream — reads entries from a Redis Stream

Each incoming JSON document is emitted once on the message output port. Use settings to scope the stream, require headers/tokens, or batch.

Outputs

  • message — The decoded JSON document, emitted as a generic object. Downstream nodes (e.g. redelay/json-to-event) apply field mapping rules.

Settings

PropertyTypeRequiredDefaultUI GroupDescription
authHeaderstring——SecurityFor http-webhook: name of the header that must carry the shared secret (e.g. X-Webhook-Token).
authSecretstring——SecurityExpected value of the auth header. Rejects requests that do not match.
groupIDstring——DeliveryConsumer group ID used when reading from Kafka or Redis streams.
maxBodyBytesinteger—1048576LimitsReject messages whose JSON body exceeds this many bytes. 0 = unlimited.
natsQueuestring——SourceFor nats transport: optional queue-group name for load-balanced delivery across workers.
pathstring——SourceFor http-webhook transport: the path to mount the receiver on (e.g. /hooks/stripe). Must start with /.
topicstring——SourceFor kafka-raw, nats, and redis-stream: the Kafka topic, NATS subject, or Redis stream key to consume.
transportstring (http-webhook|ndjson|kafka-raw|nats|…)yeshttp-webhookSourceWhere the JSON stream comes from.
urlstring——SourceFor ndjson transport: URL of the NDJSON stream (http/https or file://).

Example

yaml
# Use this node in a flow:
nodes:
  - id: my_step
    name: My Step
    kind: source
    action_ref: redelay/json-stream-source
    config:
      authHeader: "value"
      authSecret: "value"
      groupID: "value"
      maxBodyBytes: 1048576
      natsQueue: "value"
      path: "value"
      topic: "value"
      transport: http-webhook
      url: "value"