diff --git a/decoders/netflow/ipfix.go b/decoders/netflow/ipfix.go index b2064896..32db7347 100644 --- a/decoders/netflow/ipfix.go +++ b/decoders/netflow/ipfix.go @@ -456,6 +456,9 @@ type IPFIXPacket struct { FlowSets []interface{} `json:"flow-sets"` } +// DecodedPacket marks this packet as a decoded flow payload. +func (*IPFIXPacket) DecodedPacket() {} + // IPFIXOptionsTemplateFlowSet holds IPFIX options template records. type IPFIXOptionsTemplateFlowSet struct { FlowSetHeader diff --git a/decoders/netflow/nfv9.go b/decoders/netflow/nfv9.go index b74c85ba..c7840c3e 100644 --- a/decoders/netflow/nfv9.go +++ b/decoders/netflow/nfv9.go @@ -113,6 +113,9 @@ type NFv9Packet struct { FlowSets []interface{} `json:"flow-sets"` } +// DecodedPacket marks this packet as a decoded flow payload. +func (*NFv9Packet) DecodedPacket() {} + // NFv9OptionsTemplateFlowSet holds v9 options template records. type NFv9OptionsTemplateFlowSet struct { FlowSetHeader diff --git a/decoders/netflowlegacy/packet.go b/decoders/netflowlegacy/packet.go index a6e97a11..fd4ac0f7 100644 --- a/decoders/netflowlegacy/packet.go +++ b/decoders/netflowlegacy/packet.go @@ -18,6 +18,9 @@ type PacketNetFlowV5 struct { Records []RecordsNetFlowV5 `json:"records"` } +// DecodedPacket marks this packet as a decoded flow payload. +func (*PacketNetFlowV5) DecodedPacket() {} + // RecordsNetFlowV5 represents a single NetFlow v5 record entry. type RecordsNetFlowV5 struct { SrcAddr IPAddress `json:"src-addr"` diff --git a/decoders/sflow/packet.go b/decoders/sflow/packet.go index b2b18b35..74cbf759 100644 --- a/decoders/sflow/packet.go +++ b/decoders/sflow/packet.go @@ -14,6 +14,9 @@ type Packet struct { Samples []interface{} `json:"samples"` } +// DecodedPacket marks this packet as a decoded flow payload. +func (*Packet) DecodedPacket() {} + // SampleHeader contains common sample header fields. type SampleHeader struct { Format uint32 `json:"format"` diff --git a/metrics/producer.go b/metrics/producer.go index 498dca37..54250125 100644 --- a/metrics/producer.go +++ b/metrics/producer.go @@ -19,7 +19,7 @@ type PromProducerWrapper struct { } // Produce forwards to the wrapped producer and updates metrics. -func (p *PromProducerWrapper) Produce(msg interface{}, args *producer.ProduceArgs) ([]producer.ProducerMessage, error) { +func (p *PromProducerWrapper) Produce(msg producer.DecodedPacket, args *producer.ProduceArgs) ([]producer.ProducerMessage, error) { flowMessageSet, err := p.wrapped.Produce(msg, args) if err != nil { return flowMessageSet, fmt.Errorf("metrics producer: %w", err) diff --git a/producer/producer.go b/producer/producer.go index 6b293c10..7a411186 100644 --- a/producer/producer.go +++ b/producer/producer.go @@ -9,10 +9,15 @@ import ( // ProducerMessage is the generic type returned by producers. type ProducerMessage interface{} +// DecodedPacket represents a decoded flow packet ready for production. +type DecodedPacket interface { + DecodedPacket() +} + // ProducerInterface converts decoded packets into producer messages. type ProducerInterface interface { // Converts a message into a list of flow samples - Produce(msg interface{}, args *ProduceArgs) ([]ProducerMessage, error) + Produce(msg DecodedPacket, args *ProduceArgs) ([]ProducerMessage, error) // Indicates to the producer the messages returned were processed Commit([]ProducerMessage) Close() diff --git a/producer/proto/proto.go b/producer/proto/proto.go index 60d95547..8c4af85c 100644 --- a/producer/proto/proto.go +++ b/producer/proto/proto.go @@ -44,7 +44,7 @@ func (p *ProtoProducer) getSamplingRateSystem(args *producer.ProduceArgs) Sampli return sampling } -func (p *ProtoProducer) Produce(msg interface{}, args *producer.ProduceArgs) (flowMessageSet []producer.ProducerMessage, err error) { +func (p *ProtoProducer) Produce(msg producer.DecodedPacket, args *producer.ProduceArgs) (flowMessageSet []producer.ProducerMessage, err error) { tr := uint64(args.TimeReceived.UnixNano()) sa, _ := args.SamplerAddress.Unmap().MarshalBinary() switch msgConv := msg.(type) { diff --git a/producer/raw/raw.go b/producer/raw/raw.go index 4510e591..efe6b748 100644 --- a/producer/raw/raw.go +++ b/producer/raw/raw.go @@ -19,7 +19,7 @@ type RawProducer struct { // RawMessage wraps a decoded packet with metadata. type RawMessage struct { - Message interface{} `json:"message"` + Message producer.DecodedPacket `json:"message"` Src netip.AddrPort `json:"src"` TimeReceived time.Time `json:"time_received"` } @@ -40,7 +40,7 @@ func (m RawMessage) MarshalJSON() ([]byte, error) { tmpStruct := struct { Type string `json:"type"` - Message interface{} `json:"message"` + Message producer.DecodedPacket `json:"message"` Src *netip.AddrPort `json:"src"` TimeReceived *time.Time `json:"time_received"` }{ @@ -68,7 +68,7 @@ func (m RawMessage) MarshalText() ([]byte, error) { } // Produce wraps the decoded packet into a RawMessage. -func (p *RawProducer) Produce(msg interface{}, args *producer.ProduceArgs) ([]producer.ProducerMessage, error) { +func (p *RawProducer) Produce(msg producer.DecodedPacket, args *producer.ProduceArgs) ([]producer.ProducerMessage, error) { // should return msg wrapped // []*interface{msg,} return []producer.ProducerMessage{RawMessage{msg, args.Src, args.TimeReceived}}, nil diff --git a/utils/debug/producer.go b/utils/debug/producer.go index 8e5b2026..6c9bda6e 100644 --- a/utils/debug/producer.go +++ b/utils/debug/producer.go @@ -13,7 +13,7 @@ type PanicProducerWrapper struct { } // Produce calls the wrapped producer and converts panics into errors. -func (p *PanicProducerWrapper) Produce(msg interface{}, args *producer.ProduceArgs) (flowMessageSet []producer.ProducerMessage, err error) { +func (p *PanicProducerWrapper) Produce(msg producer.DecodedPacket, args *producer.ProduceArgs) (flowMessageSet []producer.ProducerMessage, err error) { defer func() { if pErr := recover(); pErr != nil {