This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / models / logger.go
3.4 kB 135 lines
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}