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