Skip to content

Latest commit

 

History

History
660 lines (492 loc) · 14.9 KB

File metadata and controls

660 lines (492 loc) · 14.9 KB

Transformer Examples Cookbook

Copy-paste ready examples for common use cases.

Table of Contents


Basic Echo

Simply forward incoming JSON to another endpoint.

File: /etc/howk/transformers/echo.lua

-- Echo transformer - forwards payload unchanged
local json = require("json")

local data = json.decode(incoming)
howk.post("https://httpbin.org/post", data)

Test:

curl -X POST http://localhost:8080/incoming/echo \
  -H "Content-Type: application/json" \
  -d '{"hello": "world"}'

Event Filtering

Only process specific event types.

File: /etc/howk/transformers/filter.lua

-- Only forward "critical" severity events
local json = require("json")

local ok, data = pcall(json.decode, incoming)
if not ok then
    return  -- Silently drop invalid JSON
end

if data.severity ~= "critical" then
    log.info("Skipping non-critical event", {severity = data.severity})
    return  -- 0 webhooks created
end

howk.post("https://pagerduty.com/webhook", data)

Webhook Routing

Route different events to different endpoints.

File: /etc/howk/transformers/router.lua

-- Route events based on 'type' field
local json = require("json")

local ok, event = pcall(json.decode, incoming)
if not ok then
    log.error("Invalid JSON")
    return
end

-- Define routes
local routes = {
    user_created = "https://crm.internal/webhooks/user",
    order_placed = "https://fulfillment.internal/webhooks/order",
    payment_received = "https://accounting.internal/webhooks/payment"
}

local endpoint = routes[event.type]
if endpoint then
    howk.post(endpoint, event)
    log.info("Routed event", {type = event.type, endpoint = endpoint})
else
    log.warn("No route for event type", {type = event.type})
end

Payload Transformation

Transform the payload before forwarding.

File: /etc/howk/transformers/transform.lua

-- Transform GitHub webhook to internal format
local json = require("json")

local ok, github = pcall(json.decode, incoming)
if not ok then
    log.error("Invalid GitHub payload")
    return
end

-- Build internal format
local internal = {
    source = "github",
    event_type = headers["X-GitHub-Event"],
    repository = github.repository and github.repository.full_name,
    actor = github.sender and github.sender.login,
    timestamp = github.repository and github.repository.updated_at,
    metadata = {
        ref = github.ref,
        commit = github.after,
        installation_id = github.installation and github.installation.id
    }
}

howk.post("https://platform.internal/events", internal)

Error Handling Patterns

Pattern 1: Graceful Degradation

local json = require("json")

local ok, data = pcall(json.decode, incoming)
if not ok then
    -- Log error but don't crash
    log.error("Invalid JSON, trying raw forward")
    
    -- Forward raw body anyway
    howk.post("https://logs.internal/raw", {
        raw_body = incoming:sub(1, 10000),  -- Limit size
        error = "json_parse_failed",
        received_at = os.date("!%Y-%m-%dT%H:%M:%SZ")
    })
    return
end

-- Normal processing...

Pattern 2: Validation with Early Return

local json = require("json")

local ok, data = pcall(json.decode, incoming)
if not ok then
    log.error("JSON parse failed", {error = tostring(data)})
    return
end

-- Validate required fields
if not data.user_id then
    log.error("Missing user_id")
    return
end

if not data.action then
    log.error("Missing action")
    return
end

-- All validations passed, process...
howk.post("https://api.internal/process", data)

Pattern 3: Try/Catch for HTTP Calls

local json = require("json")
local http = require("http")

local data = json.decode(incoming)

-- Try to enrich with user data
local ok, user_data, err = pcall(http.get, 
    "https://api.internal/users/" .. data.user_id)

if ok and not err then
    data.user = user_data
else
    log.warn("Failed to enrich user data", {user_id = data.user_id, error = err})
    -- Continue without enrichment
end

howk.post("https://webhook.internal/event", data)

Working with Headers

Extract Specific Headers

local json = require("json")

-- Get authentication info from headers
local auth_header = headers["Authorization"]
local event_type = headers["X-Webhook-Event"]
local signature = headers["X-Signature"]

log.info("Received webhook", {
    event = event_type,
    has_auth = auth_header ~= nil,
    has_signature = signature ~= nil
})

-- Include headers in forwarded payload
howk.post("https://internal.webhook/process", {
    original_body = json.decode(incoming),
    original_event_type = event_type,
    processed_at = os.time()
})

Validate Webhook Signature (HMAC)

local json = require("json")

-- Simple HMAC validation (requires crypto module setup)
local function validate_signature(body, signature, secret)
    -- Note: This is a placeholder - actual HMAC requires specific crypto
    -- For production, implement proper HMAC-SHA256 comparison
    return true  -- Implement actual validation
end

local signature = headers["X-Hub-Signature-256"]
local secret = config.webhook_secret  -- From .json config

if not validate_signature(incoming, signature, secret) {
    log.error("Invalid signature")
    return
end

-- Signature valid, process webhook
local data = json.decode(incoming)
howk.post("https://internal.webhook/process", data)

KV Store Usage

Deduplication

local json = require("json")
local kv = require("kv")

-- Prevent duplicate webhook processing
local data = json.decode(incoming)
local dedup_key = "dedup:" .. (data.id or "")

-- Check if already processed
local existing = kv.get(dedup_key)
if existing then
    log.info("Duplicate webhook ignored", {id = data.id})
    return
end

-- Mark as processed (TTL: 24 hours)
kv.set(dedup_key, "1", 86400)

-- Process the webhook
howk.post("https://api.internal/webhook", data)

Rate Limiting

local json = require("json")
local kv = require("kv")

-- Simple rate limiting per user
local data = json.decode(incoming)
local user_id = data.user_id

if not user_id then
    log.error("Missing user_id")
    return
end

local rate_key = "rate_limit:" .. user_id
local count = kv.get(rate_key) or "0"
count = tonumber(count)

if count >= 100 then
    log.warn("Rate limit exceeded", {user_id = user_id})
    return
end

-- Increment counter (expires in 1 hour)
kv.set(rate_key, tostring(count + 1), 3600)

-- Process webhook
howk.post("https://api.internal/webhook", data)

State Persistence

local json = require("json")
local kv = require("kv")

-- Track sequence numbers for ordered processing
local data = json.decode(incoming)
local stream_id = data.stream_id
local sequence = data.sequence_number

-- Get last processed sequence
local last_seq_key = "seq:" .. stream_id
local last_seq = tonumber(kv.get(last_seq_key) or "0")

if sequence <= last_seq then
    log.info("Out of order event ignored", {
        stream = stream_id,
        received = sequence,
        last_processed = last_seq
    })
    return
end

-- Process event
howk.post("https://api.internal/ordered-events", data)

-- Update sequence
kv.set(last_seq_key, tostring(sequence), 86400)

HTTP Calls from Scripts

Enrich Payload with External Data

local json = require("json")
local http = require("http")

local data = json.decode(incoming)

-- Fetch additional user info
local ok, response = pcall(http.get, 
    "https://api.internal/users/" .. data.user_id,
    {["Authorization"] = "Bearer " .. config.internal_token}
)

if ok and response and response.status_code == 200 then
    local user_info = json.decode(response.body)
    data.user_email = user_info.email
    data.user_tier = user_info.subscription_tier
else
    log.warn("Failed to fetch user info", {user_id = data.user_id})
end

howk.post("https://analytics.internal/track", data)

Conditional Processing Based on External State

local json = require("json")
local http = require("http")

local data = json.decode(incoming)

-- Check if feature flag is enabled
local ok, flag_response = pcall(http.get,
    "https://feature-flags.internal/check/" .. data.feature,
    {},
    {cache_ttl = 60}  -- Cache for 60 seconds
)

local enabled = true  -- Default
if ok and flag_response and flag_response.status_code == 200 then
    local flag = json.decode(flag_response.body)
    enabled = flag.enabled
end

if not enabled then
    log.info("Feature disabled, skipping", {feature = data.feature})
    return
end

-- Feature enabled, process
howk.post("https://api.internal/feature-events", data)

Multi-Destination Fan-Out

Send one event to multiple destinations.

local json = require("json")

local data = json.decode(incoming)

-- Define destinations
local destinations = {
    {url = "https://analytics.internal/track", priority = "high"},
    {url = "https://audit.internal/log", priority = "high"},
    {url = "https://metrics.internal/count", priority = "low"},
    {url = "https://backup.internal/archive", priority = "low"}
}

-- Send to all destinations
for _, dest in ipairs(destinations) do
    -- Customize payload per destination
    local payload = {
        original_event = data,
        destination = dest.url,
        priority = dest.priority,
        forwarded_at = os.date("!%Y-%m-%dT%H:%M:%SZ")
    }
    
    -- Each call returns a webhook ID
    local id = howk.post(dest.url, payload)
    log.debug("Created webhook", {id = id, destination = dest.url})
end

-- All 4 webhooks created

Delivery-Time Secret Injection

When the destination puts secrets in the URL or headers (Google Chat, Slack incoming webhooks, Discord, internal services keyed by query string), keep them out of Kafka by declaring unresolved templates in the transformer's .json. The worker resolves them against its process env at HTTP-send time. See transformers.md → Delivery-Time Overrides for the full contract.

Google Chat space (URL-keyed)

File: /etc/howk/transformers/gondola-google-chat.lua

local json = require("json")

local ok, data = pcall(json.decode, incoming)
if not ok then
    log.error("Invalid JSON")
    return
end

howk.post(
  "https://chat.googleapis.com/v1/spaces/AAQAIYOJwdQ/messages?messageReplyOption=REPLY_MESSAGE_FALLBACK_TO_NEW_THREAD",
  { text = data.text or "(no text)" }
)

File: /etc/howk/transformers/gondola-google-chat.json

{
  "allowed_domains": ["chat.googleapis.com"],
  "_delivery_query_params": {
    "key":   "${GONDOLA_GOOGLE_CHAT_KEY}",
    "token": "${GONDOLA_GOOGLE_CHAT_TOKEN}"
  }
}

Worker env (only on the worker pod):

export GONDOLA_GOOGLE_CHAT_KEY=AIza...
export GONDOLA_GOOGLE_CHAT_TOKEN=abc...

The webhook stored on Kafka carries the bare URL plus the unresolved ${...} templates; the URL on the wire carries &key=...&token=....

Slack incoming webhook (URL-keyed)

Slack incoming webhooks already embed the secret in the path. The same pattern works if you'd rather keep the secret out of Kafka and inject the path segment via header or alternative scheme — here we show the equivalent using a separate token query param against an internal proxy:

File: /etc/howk/transformers/slack-alerts.lua

local json = require("json")

local data = json.decode(incoming)
howk.post("https://slack-proxy.internal/services/T000/B000/post", {
    text       = data.message,
    username   = "howk",
    icon_emoji = ":bell:"
})

File: /etc/howk/transformers/slack-alerts.json

{
  "allowed_domains": ["slack-proxy.internal"],
  "_delivery_query_params": {
    "token": "${SLACK_PROXY_TOKEN}"
  },
  "_delivery_headers": {
    "Authorization": "Bearer ${SLACK_PROXY_BEARER}"
  }
}

Both reserved keys can be combined; either one alone is fine. Empty resolutions (missing env var) are dropped instead of sent as empty values.

Authorization-header pattern (header-keyed)

For destinations that take a Bearer token in Authorization:

File: /etc/howk/transformers/internal-api.json

{
  "allowed_domains": ["api.internal.example.com"],
  "_delivery_headers": {
    "Authorization":   "Bearer ${INTERNAL_API_TOKEN}",
    "X-Tenant-Secret": "${INTERNAL_TENANT_SECRET}"
  }
}

The header value is added to the outbound request headers and overrides any header with the same name set earlier in the pipeline.


Configuration Examples

Basic Config (all domains allowed)

File: /etc/howk/transformers/simple.json

{}

Production Config (restricted domains)

File: /etc/howk/transformers/production.json

{
  "allowed_domains": [
    "api.internal.company.com",
    "*.internal.company.com",
    "hooks.slack.com"
  ],
  "webhook_secret": "${WEBHOOK_SECRET}",
  "internal_token": "${INTERNAL_API_TOKEN}"
}

With Authentication

File: /etc/howk/transformers/secure.passwd

# Bcrypt hashed passwords
admin:$2b$12$xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
service:$2b$12$yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy

Testing Your Scripts

Quick Test

# Test with curl
curl -X POST http://localhost:8080/incoming/my-script \
  -H "Content-Type: application/json" \
  -d '{"test": "data"}' | jq .

With Authentication

curl -X POST http://localhost:8080/incoming/secure-script \
  -u admin:password \
  -H "Content-Type: application/json" \
  -d '{"test": "data"}'

Load Testing

# Using hey (https://github.com/rakyll/hey)
hey -n 1000 -c 10 \
  -m POST \
  -H "Content-Type: application/json" \
  -d '{"test": "data"}' \
  http://localhost:8080/incoming/my-script

Common Pitfalls

1. Not Using pcall for json.decode

local json = require("json")

-- BAD - will crash on invalid JSON
local data = json.decode(incoming)

-- GOOD - graceful handling
local ok, data = pcall(json.decode, incoming)
if not ok then
    log.error("Invalid JSON")
    return
end

2. Ignoring howk.post Return Value

-- GOOD - log the created webhook ID
local id = howk.post("https://api.example.com/webhook", data)
log.info("Created webhook", {id = id})

3. Not Checking Config Exists

-- GOOD - provide default
local timeout = config.timeout or 30
local secret = config.webhook_secret  -- may be nil

4. Logging Sensitive Data

-- BAD - logs sensitive data
log.info("Received", {body = incoming, headers = headers})

-- GOOD - selective logging
log.info("Received webhook", {
    content_type = headers["Content-Type"],
    size = #incoming,
    event_type = data.event_type  -- safe field
})