Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions cmd/blitz/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import (
"github.com/observiq/blitz/internal/telemetry/traces"
"github.com/observiq/blitz/output"
fileout "github.com/observiq/blitz/output/file"
flowout "github.com/observiq/blitz/output/flow"
hecout "github.com/observiq/blitz/output/hec"
"github.com/observiq/blitz/output/nop"
otlpgrpc "github.com/observiq/blitz/output/otlp_grpc"
Expand Down Expand Up @@ -367,6 +368,13 @@ func run(cmd *cobra.Command, args []string) error {
logger.Error("Failed to create HEC output", zap.Error(err))
return err
}
case config.OutputTypeFlow:
f := cfg.Output.Flow
outputInstance, err = flowout.New(logger, f.Host, strconv.Itoa(f.Port), flowout.Protocol(f.Protocol), f.Vendor, f.AgentIP, tel)
if err != nil {
logger.Error("Failed to create flow output", zap.Error(err))
return err
}
default:
logger.Error("Invalid output type", zap.String("type", string(cfg.Output.Type)))
return fmt.Errorf("invalid output type: %s", cfg.Output.Type)
Expand Down Expand Up @@ -501,6 +509,11 @@ func createGenerator(logger *zap.Logger, genCfg config.Generator, out output.Out
if tw, ok := out.(output.TraceWriter); ok {
consumers.TraceConsumer = output.WriterAsTraceConsumer(tw, tel)
}
// The flow output implements FlowConsumer directly (it UDP-encodes flows),
// so wire it through when the configured output supports flows.
if fc, ok := out.(embed.FlowConsumer); ok {
consumers.FlowConsumer = fc
}
// Pass the embedded library so an embed_library build resolves package
// sources; without the tag FS() is empty and resolution uses disk (PIPE-1445).
mod, err := dispatch.ForEmbed(logger, genCfg, consumers, embeddedlibrary.FS(), env, tel)
Expand Down
1 change: 1 addition & 0 deletions docker/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@ docker compose -f docker/docker-compose.telemetry-generator.yml up -d
| `blitz-hostmetrics-windows` | Host Metrics | Synthetic host metrics (CPU, memory, disk, etc.) for Windows |
| `blitz-traces` | Traces | Synthetic distributed traces (HTTP + DB spans) |
| `blitz-fix` | FIX | FIX protocol messages (4.2 / 4.4 / 5.0 SP2) across 10 asset categories — NewOrderSingle, ExecutionReport, cancel/replace/status |
| `blitz-flow` | Network Flow | Synthetic NetFlow/IPFIX/sFlow records exported over UDP |

## Running Individual Generators

Expand Down
15 changes: 15 additions & 0 deletions docker/docker-compose.telemetry-generator.yml
Original file line number Diff line number Diff line change
Expand Up @@ -230,6 +230,21 @@ services:
BLITZ_OUTPUT_OTLPGRPC_HOST: bdot-collector
BLITZ_OUTPUT_OTLPGRPC_PORT: "4317"

# Network flow records (NetFlow/IPFIX/sFlow) exported over UDP. Point it at a
# collector with a netflow/sflow receiver (e.g. bdot-collector with a
# netflowreceiver, or a standalone nfcapd/sflowtool).
blitz-flow:
<<: *blitz-common
environment:
BLITZ_GENERATOR_TYPE: flow
BLITZ_GENERATOR_FLOW_WORKERS: ${BLITZ_WORKERS:-1}
BLITZ_GENERATOR_FLOW_RATE: ${BLITZ_RATE:-1s}
BLITZ_GENERATOR_FLOW_SCENARIO: ${FLOW_SCENARIO:-default}
BLITZ_OUTPUT_TYPE: flow
BLITZ_OUTPUT_FLOW_HOST: ${FLOW_COLLECTOR_HOST:-bdot-collector}
BLITZ_OUTPUT_FLOW_PORT: ${FLOW_COLLECTOR_PORT:-2055}
BLITZ_OUTPUT_FLOW_PROTOCOL: ${FLOW_PROTOCOL:-netflow-v9}

networks:
telemetry-net:
driver: bridge
Expand Down
142 changes: 142 additions & 0 deletions docs/generator/flow.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,142 @@
# Network Flow Generator

The flow generator produces network flow records — the classic 5-tuple (source
and destination IP/port plus protocol) with byte/packet counters and routing
dimensions — and the flow output exports them over UDP in the wire format a
collector expects. FlowRecord is blitz's fourth signal type, alongside logs,
metrics, and traces.

The generator is protocol-agnostic: it yields FlowRecords, and the flow output
chooses the wire format. So one generator drives every exporter format.

## Protocol coverage

| Wire format | `protocol` value | Notes |
|-------------|------------------|-------|
| NetFlow v5 | `netflow-v5` | Fixed 48-byte records, IPv4 only, ≤30 flows/packet |
| NetFlow v9 | `netflow-v9` | Template flowsets, periodic template resend |
| IPFIX | `ipfix` | RFC 7011; 64-bit counters, absolute ms timestamps |
| sFlow v5 | `sflow` | Flow samples carrying sampled-IPv4 records |

NetFlow v1/v7/v8 are intentionally omitted (superseded; negligible residual
deployment).

## Vendor flavors

A vendor flavor makes the export genuinely vendor-distinct on the wire — not a
metadata tag. It selects a distinct encoding path whose distinguishing element a
decoder (e.g. netsampler/goflow2, which the collector `netflowreceiver` embeds)
reads back. The flavor is a single field, so flavors are mutually exclusive by
construction, and each must be paired with the protocol whose format carries its
signature.

How the distinction is carried depends on the format:

- **IPFIX** is the only format with a decoder-readable enterprise element: a
template field with the enterprise bit (`0x8000`) set, followed by the
vendor's 4-byte IANA Private Enterprise Number (PEN). goflow2 surfaces this as
the decoded field's PEN. AppFlow, jFlow, and cflowd use this.
- **NetFlow v9** has no PEN mechanism, so a v9 flavor is distinguished by a
proprietary field type id in the template (goflow2 decodes the field by its
type). NetStream and rFlow use this.

A vendor may be paired with any format it supports in the real world;
`Validate()` rejects an unsupported pairing. Because NetFlow v5 is a fixed
layout with no vendor mechanism, a v5 export is byte-identical across vendors
(the pairing is allowed for realism but carries no vendor-distinct element).

| `vendor` | Vendor | Supported formats | IPFIX PEN | v9 field type |
|-------------|--------|-------------------|-----------|---------------|
| `jflow` | Juniper | `netflow-v5`, `netflow-v9`, `ipfix` | 2636 | `0x9003` |
| `netstream` | Huawei | `netflow-v5`, `netflow-v9`, `ipfix` | 2011 | `0x9001` |
| `cflowd` | Nokia/Alcatel-Lucent | `netflow-v5`, `netflow-v9`, `ipfix` | 6527 | `0x9004` |
| `appflow` | Citrix | `ipfix` | 5951 | — |
| `rflow` | Redback/Ericsson | `netflow-v5`, `netflow-v9` | — | `0x9002` |

Format support sources: Juniper Flow Monitoring (jFlow v5/v9/IPFIX), Huawei
NetStream configuration guide (v5/v9/IPFIX), Nokia SR OS cflowd (v5/v8/v9/IPFIX;
v8 aggregation-only, excluded), Citrix AppFlow (IPFIX application), Redback/Ericsson
SmartEdge rFlow (v5/v9). Round-trip verified against netsampler/goflow2 v2.2.6,
the decoder the collector `netflowreceiver` embeds.

## OTLP output

Beyond the UDP wire formats, flows can be emitted as OTLP logs: set
`output.type: otlp-grpc` with `generator.type: flow`. Each FlowRecord is
projected to an OTLP `LogRecord` with the same network-flow semantic-convention
attributes (`source.address`, `destination.port`, `network.transport`,
`flow.io.bytes`, …) the collector's `netflowreceiver` emits for a decoded flow,
so a downstream consumer sees the same log a real collector would produce. This
path uses raw OTLP protobuf and adds no `collector/pdata` dependency.

## Transport

Flow export is UDP only — every flow collector listens on UDP. The flow output
dials `host:port` and sends each encoded packet.

## Scenario presets

The `scenario` setting shapes the byte/packet distribution and source/destination
spread:

| `scenario` | Shape |
|---------------|-------|
| `default` | Broad mix of flow sizes |
| `wan-edge` | Few, large, bidirectional flows |
| `datacenter` | Many small flows |
| `ddos-target` | Many sources onto a single destination |

## Configuration

Generator:

| YAML Path | Flag | Env | Default | Description |
|----------------------------|-------------------------------|---------------------------------|-----------|-------------|
| `generator.flow.workers` | `--generator-flow-workers` | `BLITZ_GENERATOR_FLOW_WORKERS` | `1` | Worker goroutines |
| `generator.flow.rate` | `--generator-flow-rate` | `BLITZ_GENERATOR_FLOW_RATE` | `1s` | Interval per worker |
| `generator.flow.scenario` | `--generator-flow-scenario` | `BLITZ_GENERATOR_FLOW_SCENARIO` | `default` | Traffic-shape preset |
| `generator.flow.seed` | `--generator-flow-seed` | `BLITZ_GENERATOR_FLOW_SEED` | `-1` | RNG seed (negative randomizes) |

Output:

| YAML Path | Flag | Env | Default | Description |
|------------------------|-----------------------------|-------------------------------|--------------|-------------|
| `output.flow.host` | `--output-flow-host` | `BLITZ_OUTPUT_FLOW_HOST` | `127.0.0.1` | Collector host |
| `output.flow.port` | `--output-flow-port` | `BLITZ_OUTPUT_FLOW_PORT` | `2055` | Collector UDP port |
| `output.flow.protocol` | `--output-flow-protocol` | `BLITZ_OUTPUT_FLOW_PROTOCOL` | `netflow-v9` | Wire format |
| `output.flow.vendor` | `--output-flow-vendor` | `BLITZ_OUTPUT_FLOW_VENDOR` | `""` | Optional vendor flavor |
| `output.flow.agentIP` | `--output-flow-agentip` | `BLITZ_OUTPUT_FLOW_AGENTIP` | `""` | sFlow exporter agent IP |

## Example

```yaml
generator:
type: flow
flow:
workers: 2
rate: 500ms
scenario: datacenter
output:
type: flow
flow:
host: collector.example.com
port: 2055
protocol: netflow-v9
```

## Embedding

A host consumes flows in-process by implementing `embed.FlowConsumer`:

```go
type myConsumer struct{}

func (myConsumer) ConsumeFlows(ctx context.Context, recs []embed.FlowRecord) error {
for _, r := range recs {
// r.SrcIP, r.DstIP, r.SrcPort, r.DstPort, r.Protocol, r.Bytes, r.Packets ...
}
return nil
}

host := embed.Host{Flows: myConsumer{}}
```
8 changes: 8 additions & 0 deletions embed/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,3 +23,11 @@ type MetricConsumer interface {
type TraceConsumer interface {
ConsumeTraces(ctx context.Context, spans []Span) error
}

// FlowConsumer consumes batches of flow records produced by blitz flow
// modules (netflow, ipfix, sflow).
//
// Implementations must be safe for concurrent calls.
type FlowConsumer interface {
ConsumeFlows(ctx context.Context, records []FlowRecord) error
}
3 changes: 3 additions & 0 deletions embed/host.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,9 @@ type Host struct {
// Traces is the destination for spans. Nil means spans are dropped.
Traces TraceConsumer

// Flows is the destination for flow records. Nil means flows are dropped.
Flows FlowConsumer

// Resource is the per-session base resource attributes blitz applies
// to every emitted record before module-level overrides merge on top.
//
Expand Down
108 changes: 108 additions & 0 deletions embed/otelpdata/flow.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
// Package otelpdata converts blitz signals to raw OpenTelemetry protobuf
// (go.opentelemetry.io/proto/otlp/*). It deliberately does NOT depend on
// go.opentelemetry.io/collector/pdata — blitz's OTLP path is built on the raw
// proto types, which are already a dependency.
//
// FlowToLogRecord projects a FlowRecord onto an OTLP LogRecord using the same
// attribute keys the collector's netflowreceiver emits for a decoded flow
// (receiver/netflowreceiver/parser.go addMessageAttributes), so a downstream
// consumer sees the same log a real collector would produce from the equivalent
// wire packet.
package otelpdata

import (
"net"

"github.com/observiq/blitz/embed"
commonpb "go.opentelemetry.io/proto/otlp/common/v1"
logspb "go.opentelemetry.io/proto/otlp/logs/v1"
)

// OTel semantic-convention attribute keys shared with netflowreceiver.
const (
attrSourceAddress = "source.address"
attrSourcePort = "source.port"
attrDestinationAddress = "destination.address"
attrDestinationPort = "destination.port"
attrNetworkTransport = "network.transport"
attrNetworkType = "network.type"
)

// FlowToLogRecord converts a FlowRecord to an OTLP LogRecord with network-flow
// semantic-convention attributes matching netflowreceiver's output.
func FlowToLogRecord(f embed.FlowRecord) *logspb.LogRecord {
attrs := []*commonpb.KeyValue{
str(attrSourceAddress, ipString(f.SrcIP)),
i64(attrSourcePort, int64(f.SrcPort)),
str(attrDestinationAddress, ipString(f.DstIP)),
i64(attrDestinationPort, int64(f.DstPort)),
str(attrNetworkTransport, transportName(f.Protocol)),
str(attrNetworkType, "ipv4"),
i64("flow.io.bytes", counter(f.Bytes)),
i64("flow.io.packets", counter(f.Packets)),
i64("flow.tcp_flags", int64(f.TCPFlags)),
i64("flow.ip_tos", int64(f.TOS)),
i64("flow.in_if", int64(f.InputIface)),
i64("flow.out_if", int64(f.OutputIface)),
i64("flow.src_as", int64(f.SrcAS)),
i64("flow.dst_as", int64(f.DstAS)),
i64("flow.sampling_rate", int64(f.SamplingRate)),
}
if len(f.NextHop) > 0 {
attrs = append(attrs, str("flow.next_hop", ipString(f.NextHop)))
}

rec := &logspb.LogRecord{
SeverityNumber: logspb.SeverityNumber_SEVERITY_NUMBER_INFO,
SeverityText: "INFO",
Attributes: attrs,
}
if !f.StartTime.IsZero() {
rec.TimeUnixNano = uint64(f.StartTime.UnixNano()) // #nosec G115 -- wall-clock ns fits uint64 for any realistic time
}
if !f.EndTime.IsZero() {
rec.Attributes = append(rec.Attributes, i64("flow.end", f.EndTime.UnixNano()))
}
if !f.StartTime.IsZero() {
rec.Attributes = append(rec.Attributes, i64("flow.start", f.StartTime.UnixNano()))
}
return rec
}

func str(k, v string) *commonpb.KeyValue {
return &commonpb.KeyValue{Key: k, Value: &commonpb.AnyValue{Value: &commonpb.AnyValue_StringValue{StringValue: v}}}
}

func i64(k string, v int64) *commonpb.KeyValue {
return &commonpb.KeyValue{Key: k, Value: &commonpb.AnyValue{Value: &commonpb.AnyValue_IntValue{IntValue: v}}}
}

// counter narrows a uint64 flow counter to the int64 an OTLP attribute holds,
// matching netflowreceiver, which does the same int64(pm.Bytes) conversion.
func counter(v uint64) int64 {
return int64(v) // #nosec G115 -- semconv counters are int64; matches netflowreceiver
}

func ipString(ip net.IP) string {
if len(ip) == 0 {
return ""
}
return ip.String()
}

// transportName maps an IP protocol number to the OTel network.transport value,
// matching netflowreceiver's getTransportName for the common cases.
func transportName(proto uint8) string {
switch proto {
case 1:
return "icmp"
case 6:
return "tcp"
case 17:
return "udp"
case 58:
return "ipv6-icmp"
default:
return "unknown"
}
}
51 changes: 51 additions & 0 deletions embed/otelpdata/flow_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
package otelpdata

import (
"net"
"testing"
"time"

"github.com/observiq/blitz/embed"
"github.com/stretchr/testify/require"
commonpb "go.opentelemetry.io/proto/otlp/common/v1"
)

func attrs(rec interface{ GetAttributes() []*commonpb.KeyValue }) map[string]*commonpb.AnyValue {
m := map[string]*commonpb.AnyValue{}
for _, kv := range rec.GetAttributes() {
m[kv.Key] = kv.Value
}
return m
}

func TestFlowToLogRecordSemconv(t *testing.T) {
f := embed.FlowRecord{
SrcIP: net.IPv4(10, 0, 0, 1), DstIP: net.IPv4(8, 8, 8, 8),
SrcPort: 54321, DstPort: 443, Protocol: 6, TCPFlags: 0x12,
Bytes: 4200, Packets: 42, InputIface: 2, OutputIface: 3,
SrcAS: 64500, DstAS: 15133, TOS: 0, SamplingRate: 1024,
NextHop: net.IPv4(10, 0, 0, 254),
StartTime: time.Unix(1700000000, 0), EndTime: time.Unix(1700000005, 0),
}
rec := FlowToLogRecord(f)
a := attrs(rec)

require.Equal(t, "10.0.0.1", a[attrSourceAddress].GetStringValue())
require.Equal(t, int64(54321), a[attrSourcePort].GetIntValue())
require.Equal(t, "8.8.8.8", a[attrDestinationAddress].GetStringValue())
require.Equal(t, int64(443), a[attrDestinationPort].GetIntValue())
require.Equal(t, "tcp", a[attrNetworkTransport].GetStringValue())
require.Equal(t, "ipv4", a[attrNetworkType].GetStringValue())
require.Equal(t, int64(4200), a["flow.io.bytes"].GetIntValue())
require.Equal(t, int64(42), a["flow.io.packets"].GetIntValue())
require.Equal(t, int64(64500), a["flow.src_as"].GetIntValue())
require.Equal(t, "10.0.0.254", a["flow.next_hop"].GetStringValue())
require.Equal(t, uint64(time.Unix(1700000000, 0).UnixNano()), rec.TimeUnixNano)
}

func TestTransportName(t *testing.T) {
require.Equal(t, "tcp", transportName(6))
require.Equal(t, "udp", transportName(17))
require.Equal(t, "icmp", transportName(1))
require.Equal(t, "unknown", transportName(99))
}
Loading
Loading