Initial commit
This commit is contained in:
@@ -0,0 +1,124 @@
|
||||
package analyzer
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"dip-ids/internal/config"
|
||||
"dip-ids/internal/storage"
|
||||
|
||||
"github.com/google/gopacket"
|
||||
"github.com/google/gopacket/layers"
|
||||
)
|
||||
|
||||
// Engine — основной движок анализа трафика.
|
||||
// Получает пакеты из канала, передаёт каждый всем детекторам,
|
||||
// собирает алерты и передаёт их в хранилище.
|
||||
type Engine struct {
|
||||
db *storage.DB
|
||||
detectors []Detector
|
||||
totalPkts atomic.Int64
|
||||
totalBytes atomic.Int64
|
||||
}
|
||||
|
||||
// New создаёт движок с детекторами, сконфигурированными из cfg
|
||||
func New(db *storage.DB, cfg config.DetectorsConfig) *Engine {
|
||||
return &Engine{
|
||||
db: db,
|
||||
detectors: []Detector{
|
||||
NewPortScanDetector(cfg.PortScan),
|
||||
NewBruteForceDetector(cfg.BruteForce),
|
||||
NewDoSDetector(cfg.DoS),
|
||||
NewARPSpoofDetector(cfg.ARPSpoof),
|
||||
NewDNSAnomalyDetector(cfg.DNSAnomaly),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// Run запускает обработку пакетов из канала pkts до отмены ctx.
|
||||
func (e *Engine) Run(ctx context.Context, pkts <-chan gopacket.Packet) {
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case pkt, ok := <-pkts:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
e.process(pkt)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (e *Engine) process(pkt gopacket.Packet) {
|
||||
e.totalPkts.Add(1)
|
||||
if meta := pkt.Metadata(); meta != nil {
|
||||
e.totalBytes.Add(int64(meta.CaptureLength))
|
||||
}
|
||||
|
||||
// Сохраняем статистику пакета
|
||||
e.recordStat(pkt)
|
||||
|
||||
for _, det := range e.detectors {
|
||||
alerts := det.Analyze(pkt)
|
||||
for i := range alerts {
|
||||
if err := e.db.CreateAlert(&alerts[i]); err != nil {
|
||||
log.Printf("[engine] save alert (%s): %v", det.Name(), err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (e *Engine) recordStat(pkt gopacket.Packet) {
|
||||
stat := &storage.PacketStat{
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
|
||||
if netLayer := pkt.NetworkLayer(); netLayer != nil {
|
||||
stat.SrcIP = netLayer.NetworkFlow().Src().String()
|
||||
stat.DstIP = netLayer.NetworkFlow().Dst().String()
|
||||
}
|
||||
|
||||
// Определяем протокол
|
||||
switch {
|
||||
case pkt.Layer(layers.LayerTypeTCP) != nil:
|
||||
stat.Protocol = "TCP"
|
||||
case pkt.Layer(layers.LayerTypeUDP) != nil:
|
||||
stat.Protocol = "UDP"
|
||||
case pkt.Layer(layers.LayerTypeICMPv4) != nil:
|
||||
stat.Protocol = "ICMP"
|
||||
case pkt.Layer(layers.LayerTypeARP) != nil:
|
||||
stat.Protocol = "ARP"
|
||||
case pkt.Layer(layers.LayerTypeDNS) != nil:
|
||||
stat.Protocol = "DNS"
|
||||
default:
|
||||
stat.Protocol = "OTHER"
|
||||
}
|
||||
|
||||
if meta := pkt.Metadata(); meta != nil {
|
||||
stat.Bytes = int64(meta.CaptureLength)
|
||||
}
|
||||
|
||||
// Сохраняем не каждый пакет (каждый 10-й), чтобы не раздувать БД
|
||||
if e.totalPkts.Load()%10 == 0 {
|
||||
if err := e.db.CreatePacketStat(stat); err != nil {
|
||||
log.Printf("[engine] save stat: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Stats возвращает счётчики текущей сессии
|
||||
func (e *Engine) Stats() (pkts, bytes int64) {
|
||||
return e.totalPkts.Load(), e.totalBytes.Load()
|
||||
}
|
||||
|
||||
// Reset сбрасывает состояние всех детекторов и счётчики
|
||||
func (e *Engine) Reset() {
|
||||
for _, det := range e.detectors {
|
||||
det.Reset()
|
||||
}
|
||||
e.totalPkts.Store(0)
|
||||
e.totalBytes.Store(0)
|
||||
}
|
||||
Reference in New Issue
Block a user