From b214a37ebc2fde0e467d83a71f28499db9bf6c54 Mon Sep 17 00:00:00 2001 From: barryz Date: Mon, 8 Mar 2021 12:39:16 +0800 Subject: [PATCH] log rotation: write node log to file async --- log/async_file_writer.go | 213 ++++++++++++++++++++++++++++++++++ log/async_file_writer_test.go | 30 +++++ log/handler.go | 19 +++ log/logger.go | 16 +++ node/config.go | 14 +++ node/node.go | 14 +++ 6 files changed, 306 insertions(+) create mode 100644 log/async_file_writer.go create mode 100644 log/async_file_writer_test.go diff --git a/log/async_file_writer.go b/log/async_file_writer.go new file mode 100644 index 0000000000..77ee382904 --- /dev/null +++ b/log/async_file_writer.go @@ -0,0 +1,213 @@ +package log + +import ( + "errors" + "fmt" + "os" + "path/filepath" + "sync" + "sync/atomic" + "time" +) + +type HourTicker struct { + stop chan struct{} + C <-chan time.Time +} + +func NewHourTicker() *HourTicker { + ht := &HourTicker{ + stop: make(chan struct{}), + } + ht.C = ht.Ticker() + return ht +} + +func (ht *HourTicker) Stop() { + close(ht.stop) +} + +func (ht *HourTicker) Ticker() <-chan time.Time { + ch := make(chan time.Time) + go func() { + hour := time.Now().Hour() + ticker := time.NewTicker(time.Second) + defer ticker.Stop() + for { + select { + case t := <-ticker.C: + if t.Hour() != hour { + ch <- t + hour = t.Hour() + } + case <-ht.stop: + return + } + } + }() + return ch +} + +type AsyncFileWriter struct { + filePath string + fd *os.File + + wg sync.WaitGroup + started int32 + buf chan []byte + stop chan struct{} + hourTicker *HourTicker +} + +func NewAsyncFileWriter(filePath string, bufSize int64) (*AsyncFileWriter, error) { + absFilePath, err := filepath.Abs(filePath) + if err != nil { + return nil, fmt.Errorf("get file path of logger error. filePath=%s, err=%s", filePath, err) + } + + return &AsyncFileWriter{ + filePath: absFilePath, + buf: make(chan []byte, bufSize), + stop: make(chan struct{}), + hourTicker: NewHourTicker(), + }, nil +} + +func (w *AsyncFileWriter) initLogFile() error { + var ( + fd *os.File + err error + ) + + realFilePath := w.timeFilePath(w.filePath) + fd, err = os.OpenFile(realFilePath, os.O_CREATE|os.O_APPEND|os.O_RDWR, 0644) + if err != nil { + return err + } + + w.fd = fd + _, err = os.Lstat(w.filePath) + if err == nil || os.IsExist(err) { + err = os.Remove(w.filePath) + if err != nil { + return err + } + } + + err = os.Symlink(realFilePath, w.filePath) + if err != nil { + return err + } + + return nil +} + +func (w *AsyncFileWriter) Start() error { + if !atomic.CompareAndSwapInt32(&w.started, 0, 1) { + return errors.New("logger has already been started") + } + + err := w.initLogFile() + if err != nil { + return err + } + + w.wg.Add(1) + go func() { + defer func() { + atomic.StoreInt32(&w.started, 0) + + w.flushBuffer() + w.flushAndClose() + + w.wg.Done() + }() + + for { + select { + case msg, ok := <-w.buf: + if !ok { + fmt.Fprintln(os.Stderr, "log_writer: buf channel has been closed.") + return + } + w.SyncWrite(msg) + case <-w.stop: + return + } + } + }() + return nil +} + +func (w *AsyncFileWriter) flushBuffer() { + for { + select { + case msg := <-w.buf: + w.SyncWrite(msg) + default: + return + } + } +} + +func (w *AsyncFileWriter) SyncWrite(msg []byte) { + w.rotateFile() + if w.fd != nil { + w.fd.Write(msg) + } +} + +func (w *AsyncFileWriter) rotateFile() { + select { + case <-w.hourTicker.C: + if err := w.flushAndClose(); err != nil { + fmt.Fprintf(os.Stderr, "log_writer: flush and close file error. err=%s", err) + } + if err := w.initLogFile(); err != nil { + fmt.Fprintf(os.Stderr, "log_writer: init log file error. err=%s", err) + } + default: + } +} + +func (w *AsyncFileWriter) Stop() { + w.stop <- struct{}{} + w.wg.Wait() + + w.hourTicker.Stop() +} + +func (w *AsyncFileWriter) Write(msg []byte) (n int, err error) { + buf := make([]byte, len(msg)) + copy(buf, msg) + + select { + case w.buf <- buf: + default: + } + return 0, nil +} + +func (w *AsyncFileWriter) Flush() error { + if w.fd == nil { + return nil + } + return w.fd.Sync() +} + +func (w *AsyncFileWriter) flushAndClose() error { + if w.fd == nil { + return nil + } + + err := w.fd.Sync() + if err != nil { + return err + } + + return w.fd.Close() +} + +func (w *AsyncFileWriter) timeFilePath(filePath string) string { + return filePath + "." + time.Now().Format("2006-01-02_15") +} diff --git a/log/async_file_writer_test.go b/log/async_file_writer_test.go new file mode 100644 index 0000000000..81dc96b4c1 --- /dev/null +++ b/log/async_file_writer_test.go @@ -0,0 +1,30 @@ +package log + +import ( + "io/ioutil" + "os" + "strings" + "testing" +) + +func TestWriter(t *testing.T) { + w, err := NewAsyncFileWriter("./hello.log", 100) + if err != nil { + t.Fatalf(err.Error()) + } + + w.Start() + w.Write([]byte("hello\n")) + w.Write([]byte("world\n")) + w.Stop() + files, _ := ioutil.ReadDir("./") + for _, f := range files { + fn := f.Name() + if strings.HasPrefix(fn, "hello") { + t.Log(fn) + content, _ := ioutil.ReadFile(fn) + t.Log(content) + os.Remove(fn) + } + } +} diff --git a/log/handler.go b/log/handler.go index 4ad433334e..81f250bc08 100644 --- a/log/handler.go +++ b/log/handler.go @@ -5,6 +5,7 @@ import ( "io" "net" "os" + "path" "reflect" "sync" @@ -70,6 +71,24 @@ func FileHandler(path string, fmtr Format) (Handler, error) { return closingHandler{f, StreamHandler(f, fmtr)}, nil } +// RotatingFileHandler returns a handler which writes log records to file chunks +// at the given path. When a file's size reaches the limit, the handler creates +// a new file named after the timestamp of the first log record it will contain. +func RotatingFileHandler(filePath string, limit uint64, formatter Format) (Handler, error) { + if _, err := os.Stat(path.Dir(filePath)); os.IsNotExist(err) { + err := os.MkdirAll(path.Dir(filePath), 0755) + if err != nil { + return nil, fmt.Errorf("could not create directory %s, %v", path.Dir(filePath), err) + } + } + fileWriter, err := NewAsyncFileWriter(filePath, int64(limit)) + if err != nil { + return nil, fmt.Errorf("could not create async file writer, %v", err) + } + fileWriter.Start() + return StreamHandler(fileWriter, formatter), nil +} + // NetHandler opens a socket to the given address and writes records // over the connection. func NetHandler(network, addr string, fmtr Format) (Handler, error) { diff --git a/log/logger.go b/log/logger.go index 276d6969e2..e95ec6e8ac 100644 --- a/log/logger.go +++ b/log/logger.go @@ -243,3 +243,19 @@ func (c Ctx) toArray() []interface{} { return arr } + +func NewFileLvlHandler(logPath string, maxBytesSize uint64, level string) (Handler, error) { + rfh, err := RotatingFileHandler( + logPath, + maxBytesSize, + LogfmtFormat(), + ) + if err != nil { + return nil, err + } + logLevel, err := LvlFromString(level) + if err != nil { + return nil, err + } + return LvlFilterHandler(logLevel, rfh), nil +} diff --git a/node/config.go b/node/config.go index ef1da15d70..f6a3f158cc 100644 --- a/node/config.go +++ b/node/config.go @@ -188,6 +188,9 @@ type Config struct { // Logger is a custom logger to use with the p2p.Server. Logger log.Logger `toml:",omitempty"` + // LogConfig is a configuration of the node asynchronous rotation log. + LogConfig *LogConfig `toml:",omitempty"` + staticNodesWarning bool trustedNodesWarning bool oldGethResourceWarning bool @@ -537,3 +540,14 @@ func (c *Config) warnOnce(w *bool, format string, args ...interface{}) { l.Warn(fmt.Sprintf(format, args...)) *w = true } + +// LogConfig configuration of the asynchronous rotation log. +type LogConfig struct { + // FilePath path of log file, filename must be included, + // If not specified, the datadir will be used as the default filepath + FilePath string + // MaxBytesSize max bytes allowed for log file. + MaxBytesSize uint64 + // Level log level. + Level string +} diff --git a/node/node.go b/node/node.go index 1e65fff1c2..7fc08b78e1 100644 --- a/node/node.go +++ b/node/node.go @@ -21,6 +21,7 @@ import ( "fmt" "net/http" "os" + "path" "path/filepath" "reflect" "strings" @@ -79,6 +80,19 @@ func New(conf *Config) (*Node, error) { } conf.DataDir = absdatadir } + + if conf.LogConfig != nil { + logFilePath := conf.LogConfig.FilePath + if logFilePath == "" { + logFilePath = path.Join(conf.DataDir, "node.log") + } + h, err := log.NewFileLvlHandler(logFilePath, conf.LogConfig.MaxBytesSize, conf.LogConfig.Level) + if err != nil { + return nil, err + } + log.Root().SetHandler(h) + } + if conf.Logger == nil { conf.Logger = log.New() }