-
-
Notifications
You must be signed in to change notification settings - Fork 174
Expand file tree
/
Copy pathline_prefix_writer.go
More file actions
155 lines (132 loc) · 3.71 KB
/
Copy pathline_prefix_writer.go
File metadata and controls
155 lines (132 loc) · 3.71 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
package io
import (
"bytes"
stdio "io"
"sync"
"github.com/cloudposse/atmos/pkg/perf"
)
const (
carriageReturnByte = '\r'
lineFeedByte = '\n'
)
var (
crlfBytes = []byte{carriageReturnByte, lineFeedByte}
crBytes = []byte{carriageReturnByte}
lfBytes = []byte{lineFeedByte}
)
// LinePrefixWriter prefixes complete lines and serializes writes through a
// shared output lock. Partial lines are buffered until Flush or a line ending.
type LinePrefixWriter struct {
mu sync.Mutex
writeMu *sync.Mutex
prefix string
w stdio.Writer
buffer []byte
}
// NewLinePrefixWriter creates a writer that prefixes every rendered line with
// "[prefix] ". A shared writeMu serializes writes across writers targeting the
// same terminal.
func NewLinePrefixWriter(prefix string, w stdio.Writer, writeMu *sync.Mutex) *LinePrefixWriter {
if prefix != "" {
prefix = "[" + prefix + "] "
}
return NewLinePrefixWriterRaw(prefix, w, writeMu)
}
// NewLinePrefixWriterRaw is like NewLinePrefixWriter but uses the prefix
// verbatim (no surrounding brackets or spacing), letting callers supply a
// pre-styled or colored label. A shared writeMu serializes writes across writers
// targeting the same terminal.
func NewLinePrefixWriterRaw(prefix string, w stdio.Writer, writeMu *sync.Mutex) *LinePrefixWriter {
defer perf.Track(nil, "io.NewLinePrefixWriterRaw")()
if writeMu == nil {
writeMu = &sync.Mutex{}
}
return &LinePrefixWriter{
writeMu: writeMu,
prefix: prefix,
w: w,
}
}
func (w *LinePrefixWriter) Write(p []byte) (int, error) {
w.mu.Lock()
defer w.mu.Unlock()
if len(p) == 0 {
return 0, nil
}
w.buffer = append(w.buffer, p...)
// Hold writeMu for the whole flush so that every line produced by this
// single Write call reaches the shared writer as one contiguous block.
// Locking per-line let a concurrent node's writer interleave a line in
// between two lines emitted from the same Write call.
w.writeMu.Lock()
defer w.writeMu.Unlock()
if err := w.flushCompleteLinesLocked(); err != nil {
return 0, err
}
return len(p), nil
}
// Flush writes any trailing partial line.
func (w *LinePrefixWriter) Flush() error {
defer perf.Track(nil, "io.LinePrefixWriter.Flush")()
w.mu.Lock()
defer w.mu.Unlock()
if len(w.buffer) == 0 {
return nil
}
w.writeMu.Lock()
defer w.writeMu.Unlock()
if err := w.flushCompleteLinesLocked(); err != nil {
return err
}
if len(w.buffer) == 0 {
return nil
}
line := append([]byte(nil), w.buffer...)
if err := w.writeLine(line); err != nil {
return err
}
w.buffer = w.buffer[:0]
return nil
}
// flushCompleteLinesLocked writes buffered complete lines while w.mu and w.writeMu are held.
func (w *LinePrefixWriter) flushCompleteLinesLocked() error {
for {
idx := lineEndIndex(w.buffer)
if idx < 0 {
return nil
}
end := idx + 1
if w.buffer[idx] == carriageReturnByte && end < len(w.buffer) && w.buffer[end] == lineFeedByte {
end++
}
line := append([]byte(nil), w.buffer[:end]...)
if err := w.writeLine(line); err != nil {
return err
}
w.buffer = w.buffer[end:]
}
}
// writeLine writes one already-delimited line with the configured prefix.
// Callers must hold w.writeMu.
func (w *LinePrefixWriter) writeLine(line []byte) error {
if w.w == nil {
return nil
}
line = bytes.ReplaceAll(line, crlfBytes, lfBytes)
line = bytes.ReplaceAll(line, crBytes, lfBytes)
if w.prefix == "" {
_, err := w.w.Write(line)
return err
}
_, err := stdio.WriteString(w.w, w.prefix+string(line))
return err
}
// lineEndIndex returns the first complete line-ending byte position or -1 when absent.
func lineEndIndex(p []byte) int {
for i, c := range p {
if c == lineFeedByte || (c == carriageReturnByte && i+1 < len(p)) {
return i
}
}
return -1
}