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.5 kB 151 lines
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}