This repository has no description
1package models
2
3import (
4 "bytes"
5 "encoding/json"
6 "fmt"
7 "io"
8 "os"
9 "path/filepath"
10 "strings"
11)
12
13type WorkflowLogger interface {
14 Close() error
15 DataWriter(idx int, stream string) io.Writer
16 ControlWriter(idx int, step Step, stepStatus StepStatus) io.Writer
17}
18
19type NullLogger struct{}
20
21func (l NullLogger) Close() error { return nil }
22func (l NullLogger) DataWriter(idx int, stream string) io.Writer { return io.Discard }
23func (l NullLogger) ControlWriter(idx int, step Step, stepStatus StepStatus) io.Writer {
24 return io.Discard
25}
26
27type FileWorkflowLogger struct {
28 file *os.File
29 encoder *json.Encoder
30 mask *SecretMask
31 dataWriters []*dataWriter
32}
33
34func NewFileWorkflowLogger(baseDir string, wid WorkflowId, secretValues []string) (WorkflowLogger, error) {
35 path := LogFilePath(baseDir, wid)
36 file, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644)
37 if err != nil {
38 return nil, fmt.Errorf("creating log file: %w", err)
39 }
40 return &FileWorkflowLogger{
41 file: file,
42 encoder: json.NewEncoder(file),
43 mask: NewSecretMask(secretValues),
44 }, nil
45}
46
47func LogFilePath(baseDir string, workflowID WorkflowId) string {
48 logFilePath := filepath.Join(baseDir, fmt.Sprintf("%s.log", workflowID.String()))
49 return logFilePath
50}
51
52func (l *FileWorkflowLogger) Close() error {
53 for _, w := range l.dataWriters {
54 if err := w.flush(); err != nil {
55 return err
56 }
57 }
58 return l.file.Close()
59}
60
61func (l *FileWorkflowLogger) DataWriter(idx int, stream string) io.Writer {
62 w := &dataWriter{
63 logger: l,
64 idx: idx,
65 stream: stream,
66 }
67 l.dataWriters = append(l.dataWriters, w)
68 return w
69}
70
71func (l *FileWorkflowLogger) ControlWriter(idx int, step Step, stepStatus StepStatus) io.Writer {
72 return &controlWriter{
73 logger: l,
74 idx: idx,
75 step: step,
76 stepStatus: stepStatus,
77 }
78}
79
80type dataWriter struct {
81 logger *FileWorkflowLogger
82 idx int
83 stream string
84 // trailing bytes held back so a secret split across writes still
85 // matches, flushed on Close or once enough data arrives
86 pending []byte
87}
88
89func (w *dataWriter) Write(p []byte) (int, error) {
90 w.pending = append(w.pending, p...)
91 if err := w.flushCompleteLines(); err != nil {
92 return 0, err
93 }
94 return len(p), nil
95}
96
97func (w *dataWriter) flushCompleteLines() error {
98 limit := len(w.pending) - w.logger.mask.Window()
99 if limit <= 0 {
100 return nil
101 }
102
103 for {
104 lineEnd := bytes.IndexByte(w.pending[:limit], '\n')
105 if lineEnd < 0 {
106 return nil
107 }
108 lineEnd++
109 line := append([]byte(nil), w.pending[:lineEnd]...)
110 w.pending = w.pending[lineEnd:]
111 limit -= lineEnd
112 if err := w.emit(line); err != nil {
113 return err
114 }
115 }
116}
117
118// the writer is done, so a buffered tail can no longer grow into a full
119// secret and goes out as-is
120func (w *dataWriter) flush() error {
121 if len(w.pending) == 0 {
122 return nil
123 }
124 pending := w.pending
125 w.pending = nil
126 return w.emit(pending)
127}
128
129func (w *dataWriter) emit(p []byte) error {
130 line := strings.TrimRight(string(p), "\r\n")
131 if w.logger.mask != nil {
132 line = w.logger.mask.Mask(line)
133 }
134 entry := NewDataLogLine(w.idx, line, w.stream)
135 return w.logger.encoder.Encode(entry)
136}
137
138type controlWriter struct {
139 logger *FileWorkflowLogger
140 idx int
141 step Step
142 stepStatus StepStatus
143}
144
145func (w *controlWriter) Write(_ []byte) (int, error) {
146 entry := NewControlLogLine(w.idx, w.step, w.stepStatus)
147 if err := w.logger.encoder.Encode(entry); err != nil {
148 return 0, err
149 }
150 return len(w.step.Name()), nil
151}