From 2735819b4f468685b42383bf6c301d575c0d1e61 Mon Sep 17 00:00:00 2001 From: lion Date: Thu, 4 Sep 2025 14:59:20 +0800 Subject: [PATCH] source ip data processor --- maker/golang/main.go | 79 +++++++++++++++++++++++ maker/golang/xdb/processor.go | 117 ++++++++++++++++++++++++++++++++++ 2 files changed, 196 insertions(+) create mode 100644 maker/golang/xdb/processor.go diff --git a/maker/golang/main.go b/maker/golang/main.go index 0a71198..73cd839 100644 --- a/maker/golang/main.go +++ b/maker/golang/main.go @@ -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() } diff --git a/maker/golang/xdb/processor.go b/maker/golang/xdb/processor.go new file mode 100644 index 0000000..6754a8d --- /dev/null +++ b/maker/golang/xdb/processor.go @@ -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 +}