Add ClearDBData functionality for analysis and repository components

- Introduced `ClearDBData` methods in `Analysis`, `Alert`, and `BruteForceProtection` components.
- Implemented `Clear` operations for `AlertGroupRepository` and `BruteForceProtectionGroupRepository` to reset database buckets.
- Updated `Analyzer` to invoke `ClearDBData` for cleanup logic.
This commit is contained in:
2026-02-28 11:37:25 +05:00
parent 6b482a350b
commit a648647e4a
6 changed files with 75 additions and 0 deletions
+16
View File
@@ -14,6 +14,7 @@ import (
type Analyzer interface { type Analyzer interface {
Run(ctx context.Context) Run(ctx context.Context)
ClearDBData() error
Close() error Close() error
} }
@@ -86,6 +87,21 @@ func (a *analyzer) Run(ctx context.Context) {
a.logger.Debug("Analyzer is start") a.logger.Debug("Analyzer is start")
} }
func (a *analyzer) ClearDBData() error {
a.logger.Debug("Clear data")
clearDBErrors, err := a.analysis.ClearDBData()
if err != nil {
for _, err := range clearDBErrors {
a.logger.Error(err.Error())
}
return err
}
return nil
}
func (a *analyzer) processLogs(ctx context.Context) { func (a *analyzer) processLogs(ctx context.Context) {
for { for {
select { select {
+19
View File
@@ -1,6 +1,8 @@
package log package log
import ( import (
"fmt"
analysisServices "git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/analyzer/log/analysis" analysisServices "git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/analyzer/log/analysis"
"git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/analyzer/log/analysis/alert_group" "git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/analyzer/log/analysis/alert_group"
"git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/analyzer/log/analysis/brute_force_protection_group" "git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/analyzer/log/analysis/brute_force_protection_group"
@@ -12,6 +14,7 @@ import (
type Analysis interface { type Analysis interface {
Alert(entry *analysisServices.Entry) Alert(entry *analysisServices.Entry)
BruteForceProtection(entry *analysisServices.Entry) BruteForceProtection(entry *analysisServices.Entry)
ClearDBData() ([]error, error)
} }
type analysis struct { type analysis struct {
@@ -36,3 +39,19 @@ func (a *analysis) Alert(entry *analysisServices.Entry) {
func (a *analysis) BruteForceProtection(entry *analysisServices.Entry) { func (a *analysis) BruteForceProtection(entry *analysisServices.Entry) {
a.bruteForceProtectionService.Analyze(entry) a.bruteForceProtectionService.Analyze(entry)
} }
func (a *analysis) ClearDBData() ([]error, error) {
var errClearDB []error
if err := a.alertService.ClearDBData(); err != nil {
errClearDB = append(errClearDB, err)
}
if err := a.bruteForceProtectionService.ClearDBData(); err != nil {
errClearDB = append(errClearDB, err)
}
if len(errClearDB) > 0 {
return nil, fmt.Errorf("failed to clear database data: %v", errClearDB)
}
return errClearDB, nil
}
@@ -13,6 +13,7 @@ import (
type Alert interface { type Alert interface {
Analyze(entry *Entry) Analyze(entry *Entry)
ClearDBData() error
} }
type alert struct { type alert struct {
@@ -82,6 +83,10 @@ func (a *alert) Analyze(entry *Entry) {
} }
} }
func (a *alert) ClearDBData() error {
return a.alertGroupService.ClearDBData()
}
func (a *alert) analyzeRule(rule *config.AlertRule, message string) alertAnalyzeRuleReturn { func (a *alert) analyzeRule(rule *config.AlertRule, message string) alertAnalyzeRuleReturn {
result := alertAnalyzeRuleReturn{ result := alertAnalyzeRuleReturn{
found: false, found: false,
@@ -15,6 +15,7 @@ import (
type BruteForceProtection interface { type BruteForceProtection interface {
Analyze(entry *Entry) Analyze(entry *Entry)
ClearDBData() error
} }
type BlockIPFunc func(blockIP blocking.BlockIP) error type BlockIPFunc func(blockIP blocking.BlockIP) error
@@ -106,6 +107,10 @@ func (p *bruteForceProtection) Analyze(entry *Entry) {
} }
} }
func (p *bruteForceProtection) ClearDBData() error {
return p.groupService.ClearDBData()
}
func (p *bruteForceProtection) analyzeRule(rule *brute_force_protection.Rule, message string) bruteForceProtectionAnalyzeRuleReturn { func (p *bruteForceProtection) analyzeRule(rule *brute_force_protection.Rule, message string) bruteForceProtectionAnalyzeRuleReturn {
result := bruteForceProtectionAnalyzeRuleReturn{ result := bruteForceProtectionAnalyzeRuleReturn{
found: false, found: false,
@@ -2,14 +2,17 @@ package repository
import ( import (
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/db/entity" "git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/db/entity"
"go.etcd.io/bbolt" "go.etcd.io/bbolt"
bboltErrors "go.etcd.io/bbolt/errors"
) )
type AlertGroupRepository interface { type AlertGroupRepository interface {
Update(name string, f func(*entity.AlertGroup) (*entity.AlertGroup, error)) error Update(name string, f func(*entity.AlertGroup) (*entity.AlertGroup, error)) error
Clear() error
} }
type alertGroupRepository struct { type alertGroupRepository struct {
@@ -55,3 +58,15 @@ func (r *alertGroupRepository) Update(name string, f func(*entity.AlertGroup) (*
return b.Put(key, data) return b.Put(key, data)
}) })
} }
func (r *alertGroupRepository) Clear() error {
return r.db.Update(func(tx *bbolt.Tx) error {
err := tx.DeleteBucket([]byte(r.bucket))
if errors.Is(err, bboltErrors.ErrBucketNotFound) {
// If the bucket may not exist, ignore ErrBucketNotFound
return nil
}
_, err = tx.CreateBucketIfNotExists([]byte(r.bucket))
return err
})
}
@@ -2,15 +2,18 @@ package repository
import ( import (
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"net" "net"
"git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/db/entity" "git.kor-elf.net/kor-elf-shield/kor-elf-shield/internal/daemon/db/entity"
"go.etcd.io/bbolt" "go.etcd.io/bbolt"
bboltErrors "go.etcd.io/bbolt/errors"
) )
type BruteForceProtectionGroupRepository interface { type BruteForceProtectionGroupRepository interface {
Update(name string, ip net.IP, f func(*entity.BruteForceProtectionGroup) (*entity.BruteForceProtectionGroup, error)) error Update(name string, ip net.IP, f func(*entity.BruteForceProtectionGroup) (*entity.BruteForceProtectionGroup, error)) error
Clear() error
} }
type bruteForceProtectionGroupRepository struct { type bruteForceProtectionGroupRepository struct {
@@ -60,6 +63,18 @@ func (r *bruteForceProtectionGroupRepository) Update(name string, ip net.IP, f f
}) })
} }
func (r *bruteForceProtectionGroupRepository) Clear() error {
return r.db.Update(func(tx *bbolt.Tx) error {
err := tx.DeleteBucket([]byte(r.bucket))
if errors.Is(err, bboltErrors.ErrBucketNotFound) {
// If the bucket may not exist, ignore ErrBucketNotFound
return nil
}
_, err = tx.CreateBucketIfNotExists([]byte(r.bucket))
return err
})
}
func keyGroupIP(groupID string, ip net.IP) ([]byte, error) { func keyGroupIP(groupID string, ip net.IP) ([]byte, error) {
if ip == nil { if ip == nil {
return nil, fmt.Errorf("ip cannot be nil") return nil, fmt.Errorf("ip cannot be nil")