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
3 changes: 3 additions & 0 deletions decoders/netflow/ipfix.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 3 additions & 0 deletions decoders/netflow/nfv9.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 3 additions & 0 deletions decoders/netflowlegacy/packet.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand Down
3 changes: 3 additions & 0 deletions decoders/sflow/packet.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand Down
2 changes: 1 addition & 1 deletion metrics/producer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
7 changes: 6 additions & 1 deletion producer/producer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
2 changes: 1 addition & 1 deletion producer/proto/proto.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
6 changes: 3 additions & 3 deletions producer/raw/raw.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
}
Expand All @@ -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"`
}{
Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion utils/debug/producer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading