package cli import ( "bytes" "io" "strings" "time" "git.unkin.net/unkin/logarchiver/internal/event" ) // lineFilter is a predicate over a single NDJSON event line. type lineFilter func(raw []byte) bool // newLineFilter builds a predicate from optional host/time constraints. A nil // filter (all constraints empty) means "pass everything". func newLineFilter(host string, from, to time.Time) lineFilter { if host == "" && from.IsZero() && to.IsZero() { return nil } return func(raw []byte) bool { meta := event.Extract(raw) if host != "" && !globMatch(host, meta.Host) { return false } if !from.IsZero() || !to.IsZero() { // Events without a parseable timestamp are kept (we cannot exclude // them on time grounds without dropping data). if meta.Ok { if !from.IsZero() && meta.Timestamp.Before(from) { return false } if !to.IsZero() && meta.Timestamp.After(to) { return false } } } return true } } // filterWriter forwards only complete NDJSON lines that satisfy filter. It // buffers a trailing partial line across Writes so streaming decryption can feed // it arbitrary chunks. Flush must be called at end to emit any final unterminated // line. A nil filter forwards bytes verbatim. type filterWriter struct { dst io.Writer filter lineFilter buf bytes.Buffer } func newFilterWriter(dst io.Writer, filter lineFilter) *filterWriter { return &filterWriter{dst: dst, filter: filter} } func (w *filterWriter) Write(p []byte) (int, error) { if w.filter == nil { return w.dst.Write(p) } w.buf.Write(p) for { data := w.buf.Bytes() i := bytes.IndexByte(data, '\n') if i < 0 { break } line := data[:i] if len(bytes.TrimSpace(line)) > 0 && w.filter(line) { if _, err := w.dst.Write(line); err != nil { return 0, err } if _, err := w.dst.Write([]byte{'\n'}); err != nil { return 0, err } } w.buf.Next(i + 1) } return len(p), nil } // Flush emits a trailing line that had no terminating newline. func (w *filterWriter) Flush() error { if w.filter == nil { return nil } line := bytes.TrimRight(w.buf.Bytes(), "\n") w.buf.Reset() if len(bytes.TrimSpace(line)) > 0 && w.filter(line) { if _, err := w.dst.Write(line); err != nil { return err } if _, err := w.dst.Write([]byte{'\n'}); err != nil { return err } } return nil } // globMatch matches pattern against s where '*' matches any run of characters. // With no '*', it is an exact match. func globMatch(pattern, s string) bool { if !strings.Contains(pattern, "*") { return pattern == s } parts := strings.Split(pattern, "*") // Anchor first part. if !strings.HasPrefix(s, parts[0]) { return false } s = s[len(parts[0]):] for _, part := range parts[1 : len(parts)-1] { if part == "" { continue } idx := strings.Index(s, part) if idx < 0 { return false } s = s[idx+len(part):] } // Anchor last part. return strings.HasSuffix(s, parts[len(parts)-1]) }