source ip data processor

This commit is contained in:
lion 2025-09-04 14:59:20 +08:00
parent 980c9cfdd4
commit 2735819b4f
2 changed files with 196 additions and 0 deletions

View File

@ -28,6 +28,7 @@ func printHelp() {
fmt.Printf(" search binary xdb search test\n")
fmt.Printf(" bench binary xdb bench test\n")
fmt.Printf(" edit edit the source ip data\n")
fmt.Printf(" process process the source ip data\n")
}
// Iterate the cli flags
@ -589,6 +590,82 @@ func edit() {
}
}
func process() {
var err error
var srcFile, dstFile = "", ""
var fieldList, logLevel = "", ""
var fErr = iterateFlags(func(key string, val string) error {
switch key {
case "src":
srcFile = val
case "dst":
dstFile = val
case "field-list":
fieldList = val
case "log-level":
logLevel = val
default:
return fmt.Errorf("undefined option '%s=%s'\n", key, val)
}
return nil
})
if fErr != nil {
fmt.Printf("failed to parse flags: %s", fErr)
return
}
if srcFile == "" || dstFile == "" {
fmt.Printf("%s process [command options]\n", os.Args[0])
fmt.Printf("options:\n")
fmt.Printf(" --src string source ip text file path\n")
fmt.Printf(" --dst string target ip text file path\n")
fmt.Printf(" --field-list string field index list imploded with ',' eg: 0,1,2,3-6,7\n")
fmt.Printf(" --log-level string set the log level, options: debug/info/warn/error\n")
return
}
// check and apply the log level
err = applyLogLevel(logLevel)
if err != nil {
slog.Error("failed to apply log level", "error", err)
return
}
fields, err := getFilterFields(fieldList)
if err != nil {
slog.Error("failed to get filter fields", "error", err)
return
}
// make the binary file
tStart := time.Now()
processor, err := xdb.NewProcessor(srcFile, dstFile, fields)
if err != nil {
fmt.Printf("failed to create %s\n", err)
return
}
err = processor.Init()
if err != nil {
fmt.Printf("failed Init: %s\n", err)
return
}
slog.Info("Processing", "src", srcFile, "dst", dstFile, "logLevel", logLevel)
err = processor.Start()
if err != nil {
fmt.Printf("failed Start: %s\n", err)
return
}
err = processor.End()
if err != nil {
fmt.Printf("failed End: %s\n", err)
}
slog.Info("processor done", "elapsed", time.Since(tStart))
}
func main() {
if len(os.Args) < 2 {
printHelp()
@ -605,6 +682,8 @@ func main() {
testBench()
case "edit":
edit()
case "process":
process()
default:
printHelp()
}

View File

@ -0,0 +1,117 @@
// Copyright 2022 The Ip2Region Authors. All rights reserved.
// Use of this source code is governed by a Apache2.0-style
// license that can be found in the LICENSE file.
// original source ip processor
package xdb
import (
"fmt"
"log/slog"
"os"
"sort"
"time"
)
type Processor struct {
srcHandle *os.File
dstHandle *os.File
fields []int
segments []*Segment
}
func NewProcessor(srcFile string, dstFile string, fields []int) (*Processor, error) {
// open the source file with READONLY mode
srcHandle, err := os.OpenFile(srcFile, os.O_RDONLY, 0600)
if err != nil {
return nil, fmt.Errorf("open source file `%s`: %w", srcFile, err)
}
// open the destination file with Read/Write mode
dstHandle, err := os.OpenFile(dstFile, os.O_RDWR|os.O_CREATE|os.O_TRUNC, 0666)
if err != nil {
return nil, fmt.Errorf("open target file `%s`: %w", dstFile, err)
}
return &Processor{
srcHandle: srcHandle,
dstHandle: dstHandle,
segments: []*Segment{},
}, nil
}
func (p *Processor) loadSegments() error {
slog.Info("try to load the segments ... ")
var tStart = time.Now()
var iErr = IterateSegments(p.srcHandle, func(l string) {
slog.Debug("loaded", "segment", l)
}, func(seg *Segment) error {
// check the continuity of the data segment
// if err := seg.AfterCheck(last); err != nil {
// return err
// }
// apply the field filter
region, err := RegionFiltering(seg.Region, p.fields)
if err != nil {
return err
}
// slog.Info("filtered", "region", region)
seg.Region = region
p.segments = append(p.segments, seg)
return nil
})
if iErr != nil {
return fmt.Errorf("failed to load segments: %s", iErr)
}
slog.Info("all segments loaded", "length", len(p.segments), "elapsed", time.Since(tStart))
return nil
}
// Init the db binary file
func (p *Processor) Init() error {
// load all the segments
err := p.loadSegments()
if err != nil {
return fmt.Errorf("load segments: %w", err)
}
return nil
}
func (p *Processor) Start() error {
slog.Info("try to sort all the segments based on its start ip ...")
sort.Slice(p.segments, func(i, j int) bool {
return IPCompare(p.segments[i].StartIP, p.segments[j].StartIP) < 0
})
slog.Info("try to write all segments to target file ...")
for _, seg := range p.segments {
_, err := fmt.Fprintln(p.dstHandle, seg.String())
if err != nil {
return fmt.Errorf("write segment index for '%s': %w", seg.String(), err)
}
}
slog.Info("process done", "segments", len(p.segments))
return nil
}
func (p *Processor) End() error {
err := p.dstHandle.Close()
if err != nil {
return err
}
err = p.srcHandle.Close()
if err != nil {
return err
}
return nil
}