diff options
Diffstat (limited to 'internal/fmsg')
| -rw-r--r-- | internal/fmsg/defer.go | 98 | ||||
| -rw-r--r-- | internal/fmsg/fmsg.go | 88 | ||||
| -rw-r--r-- | internal/fmsg/verbose.go | 22 |
3 files changed, 77 insertions, 131 deletions
diff --git a/internal/fmsg/defer.go b/internal/fmsg/defer.go deleted file mode 100644 index 677f5ceb..00000000 --- a/internal/fmsg/defer.go +++ /dev/null @@ -1,98 +0,0 @@ -package fmsg - -import ( - "os" - "sync" - "sync/atomic" -) - -var ( - wstate atomic.Bool - dropped atomic.Uint64 - withhold = make(chan struct{}, 1) - msgbuf = make(chan dOp, 64) // these ops are tiny so a large buffer is allocated for withholding output - - dequeueOnce sync.Once - queueSync sync.WaitGroup -) - -func dequeue() { - go func() { - for { - select { - case op := <-msgbuf: - op.Do() - queueSync.Done() - case <-withhold: - <-withhold - } - } - }() -} - -// queue submits ops to msgbuf but drops messages -// when the buffer is full and dequeue is withholding -func queue(op dOp) { - queueSync.Add(1) - - select { - case msgbuf <- op: - default: - // send the op anyway if not withholding - // as dequeue will get to it eventually - if !wstate.Load() { - msgbuf <- op - } else { - queueSync.Done() - // increment dropped message count - dropped.Add(1) - } - } -} - -type dOp interface{ Do() } - -func Exit(code int) { - Resume() // resume here to avoid deadlock - queueSync.Wait() - os.Exit(code) -} - -func Suspend() { - dequeueOnce.Do(dequeue) - if wstate.CompareAndSwap(false, true) { - queueSync.Wait() - withhold <- struct{}{} - } -} - -func Resume() { - dequeueOnce.Do(dequeue) - if wstate.CompareAndSwap(true, false) { - withhold <- struct{}{} - if d := dropped.Swap(0); d != 0 { - Printf("dropped %d messages during withhold", d) - } - } -} - -type dPrint []any - -func (v dPrint) Do() { - std.Print(v...) -} - -type dPrintf struct { - format string - v []any -} - -func (d *dPrintf) Do() { - std.Printf(d.format, d.v...) -} - -type dPrintln []any - -func (v dPrintln) Do() { - std.Println(v...) -} diff --git a/internal/fmsg/fmsg.go b/internal/fmsg/fmsg.go index d4027204..48538380 100644 --- a/internal/fmsg/fmsg.go +++ b/internal/fmsg/fmsg.go @@ -2,39 +2,85 @@ package fmsg import ( + "bytes" + "io" "log" "os" + "sync" + "sync/atomic" + "syscall" ) -var std = log.New(os.Stderr, "fortify: ", 0) +const ( + bufSize = 4 * 1024 + bufSizeMax = 16 * 1024 * 1024 +) -func SetPrefix(prefix string) { - prefix += ": " - std.SetPrefix(prefix) - std.SetPrefix(prefix) -} +var o = &suspendable{w: os.Stderr} + +// Prepare configures the system logger for [Suspend] and [Resume] to take effect. +func Prepare(prefix string) { log.SetPrefix(prefix + ": "); log.SetFlags(0); log.SetOutput(o) } + +type suspendable struct { + w io.Writer + s atomic.Bool -func Print(v ...any) { - dequeueOnce.Do(dequeue) - queue(dPrint(v)) + buf bytes.Buffer + bufOnce sync.Once + bufMu sync.Mutex + dropped int } -func Printf(format string, v ...any) { - dequeueOnce.Do(dequeue) - queue(&dPrintf{format, v}) +func (s *suspendable) Write(p []byte) (n int, err error) { + if !s.s.Load() { + return s.w.Write(p) + } + s.bufOnce.Do(func() { s.prepareBuf() }) + + s.bufMu.Lock() + defer s.bufMu.Unlock() + + if l := len(p); s.buf.Len()+l > bufSizeMax { + s.dropped += l + return 0, syscall.ENOMEM + } + return s.buf.Write(p) } -func Println(v ...any) { - dequeueOnce.Do(dequeue) - queue(dPrintln(v)) +func (s *suspendable) prepareBuf() { s.buf.Grow(bufSize) } +func (s *suspendable) Suspend() bool { return o.s.CompareAndSwap(false, true) } +func (s *suspendable) Resume() (resumed bool, dropped uintptr, n int64, err error) { + if o.s.CompareAndSwap(true, false) { + o.bufMu.Lock() + defer o.bufMu.Unlock() + + resumed = true + dropped = uintptr(o.dropped) + + o.dropped = 0 + n, err = io.Copy(s.w, &s.buf) + s.buf = bytes.Buffer{} + s.prepareBuf() + } + return } -func Fatal(v ...any) { - Print(v...) - Exit(1) +func Suspend() bool { return o.Suspend() } +func Resume() bool { + resumed, dropped, _, err := o.Resume() + if err != nil { + // probably going to result in an error as well, + // so this call is as good as unreachable + log.Printf("cannot dump buffer on resume: %v", err) + } + if resumed && dropped > 0 { + log.Fatalf("dropped %d bytes while output is suspended", dropped) + } + return resumed } -func Fatalf(format string, v ...any) { - Printf(format, v...) - Exit(1) +func BeforeExit() { + if Resume() { + log.Printf("beforeExit reached on suspended output") + } } diff --git a/internal/fmsg/verbose.go b/internal/fmsg/verbose.go index 72a92a69..a8e4c81f 100644 --- a/internal/fmsg/verbose.go +++ b/internal/fmsg/verbose.go @@ -1,25 +1,23 @@ package fmsg -import "sync/atomic" +import ( + "log" + "sync/atomic" +) var verbose = new(atomic.Bool) -func Verbose() bool { - return verbose.Load() -} - -func SetVerbose(v bool) { - verbose.Store(v) -} +func Load() bool { return verbose.Load() } +func Store(v bool) { verbose.Store(v) } -func VPrintf(format string, v ...any) { +func Verbosef(format string, v ...any) { if verbose.Load() { - Printf(format, v...) + log.Printf(format, v...) } } -func VPrintln(v ...any) { +func Verbose(v ...any) { if verbose.Load() { - Println(v...) + log.Println(v...) } } |
