package logging import ( "errors" "fmt" "io" "io/fs" "log/slog" "os" "path/filepath" "strconv" "strings" "sync" "time" _ "time/tzdata" ) const weeklyLogRetention = 13 type logRotationTimer interface { C() <-chan time.Time Stop() bool } type realLogRotationTimer struct { *time.Timer } func (t realLogRotationTimer) C() <-chan time.Time { return t.Timer.C } type weeklyLogRuntime struct { now func() time.Time newTimer func(time.Duration) logRotationTimer rename func(string, string) error diagnostic io.Writer } func defaultWeeklyLogRuntime() weeklyLogRuntime { return weeklyLogRuntime{ now: time.Now, newTimer: func(duration time.Duration) logRotationTimer { return realLogRotationTimer{Timer: time.NewTimer(duration)} }, rename: os.Rename, diagnostic: os.Stderr, } } type weeklyLogWriter struct { path string location *time.Location now func() time.Time newTimer func(time.Duration) logRotationTimer rename func(string, string) error diagnostic io.Writer mu sync.Mutex file *os.File weekStart time.Time closed bool stop chan struct{} done chan struct{} closeOnce sync.Once closeErr error } func newWeeklyLogWriter(path string, location *time.Location, runtime weeklyLogRuntime) (*weeklyLogWriter, error) { if location == nil { location = time.Local } if runtime.now == nil { runtime.now = time.Now } if runtime.newTimer == nil { runtime.newTimer = defaultWeeklyLogRuntime().newTimer } if runtime.rename == nil { runtime.rename = os.Rename } if runtime.diagnostic == nil { runtime.diagnostic = os.Stderr } now := runtime.now().In(location) weekStart := logWeekStart(now, location) if info, err := os.Stat(path); err == nil { weekStart = logWeekStart(info.ModTime(), location) } else if !errors.Is(err, os.ErrNotExist) { return nil, fmt.Errorf("inspect active log: %w", err) } file, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o666) if err != nil { return nil, err } writer := &weeklyLogWriter{ path: path, location: location, now: runtime.now, newTimer: runtime.newTimer, rename: runtime.rename, diagnostic: runtime.diagnostic, file: file, weekStart: weekStart, stop: make(chan struct{}), done: make(chan struct{}), } writer.rotateIfNeeded(now) go writer.run() return writer, nil } func (w *weeklyLogWriter) Write(data []byte) (int, error) { now := w.now().In(w.location) w.mu.Lock() rotationErr := w.rotateIfNeededLocked(now) if w.closed || w.file == nil { w.mu.Unlock() w.reportRotationError(rotationErr) return 0, os.ErrClosed } n, writeErr := w.file.Write(data) w.mu.Unlock() w.reportRotationError(rotationErr) return n, writeErr } func (w *weeklyLogWriter) Close() error { w.closeOnce.Do(func() { close(w.stop) <-w.done w.mu.Lock() w.closed = true if w.file != nil { w.closeErr = w.file.Close() w.file = nil } w.mu.Unlock() }) return w.closeErr } func (w *weeklyLogWriter) run() { defer close(w.done) for { now := w.now().In(w.location) next := nextLogWeekStart(now, w.location) duration := next.Sub(now) if duration <= 0 { duration = time.Nanosecond } timer := w.newTimer(duration) select { case <-timer.C(): w.rotateIfNeeded(w.now().In(w.location)) case <-w.stop: timer.Stop() return } } } func (w *weeklyLogWriter) rotateIfNeeded(now time.Time) { w.mu.Lock() err := w.rotateIfNeededLocked(now) w.mu.Unlock() w.reportRotationError(err) } func (w *weeklyLogWriter) rotateIfNeededLocked(now time.Time) error { if w.closed { return nil } targetWeek := logWeekStart(now, w.location) if !targetWeek.After(w.weekStart) { return nil } // Record the attempted week even if rotation fails so every write in a // broken environment does not retry and emit another diagnostic. w.weekStart = targetWeek return w.rotateLocked() } func (w *weeklyLogWriter) rotateLocked() error { if w.file == nil { return fmt.Errorf("active log file is unavailable") } if err := w.file.Close(); err != nil { w.file = nil return errors.Join(fmt.Errorf("close active log: %w", err), w.reopenActiveLocked()) } w.file = nil if err := w.cleanupOlderGenerationsLocked(); err != nil { return errors.Join(err, w.reopenActiveLocked()) } if err := removeIfExists(w.generationPath(weeklyLogRetention)); err != nil { return errors.Join(fmt.Errorf("remove oldest weekly log: %w", err), w.reopenActiveLocked()) } for generation := weeklyLogRetention - 1; generation >= 1; generation-- { source := w.generationPath(generation) if _, err := os.Stat(source); err != nil { if errors.Is(err, os.ErrNotExist) { continue } return errors.Join(fmt.Errorf("inspect weekly log generation %d: %w", generation, err), w.reopenActiveLocked()) } if err := w.rename(source, w.generationPath(generation+1)); err != nil { return errors.Join(fmt.Errorf("shift weekly log generation %d: %w", generation, err), w.reopenActiveLocked()) } } activeRenamed := false if _, err := os.Stat(w.path); err == nil { if err := w.rename(w.path, w.generationPath(1)); err != nil { return errors.Join(fmt.Errorf("archive active log: %w", err), w.reopenActiveLocked()) } activeRenamed = true } else if !errors.Is(err, os.ErrNotExist) { return errors.Join(fmt.Errorf("inspect active log: %w", err), w.reopenActiveLocked()) } file, err := os.OpenFile(w.path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o666) if err == nil { w.file = file return nil } rotationErr := fmt.Errorf("open new active log: %w", err) if activeRenamed { if rollbackErr := w.rename(w.generationPath(1), w.path); rollbackErr != nil { rotationErr = errors.Join(rotationErr, fmt.Errorf("restore archived active log: %w", rollbackErr)) } } return errors.Join(rotationErr, w.reopenActiveLocked()) } func (w *weeklyLogWriter) cleanupOlderGenerationsLocked() error { directory := filepath.Dir(w.path) entries, err := os.ReadDir(directory) if err != nil { return fmt.Errorf("list weekly log directory: %w", err) } base := filepath.Base(w.path) extension := filepath.Ext(base) stem := strings.TrimSuffix(base, extension) prefix := stem + "." for _, entry := range entries { if entry.IsDir() { continue } name := entry.Name() if !strings.HasPrefix(name, prefix) || !strings.HasSuffix(name, extension) { continue } generationText := strings.TrimSuffix(strings.TrimPrefix(name, prefix), extension) generation, err := strconv.Atoi(generationText) if err != nil || generation <= weeklyLogRetention { continue } if err := os.Remove(filepath.Join(directory, name)); err != nil && !errors.Is(err, os.ErrNotExist) { return fmt.Errorf("remove weekly log generation %d: %w", generation, err) } } return nil } func (w *weeklyLogWriter) reopenActiveLocked() error { file, err := os.OpenFile(w.path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o666) if err != nil { return fmt.Errorf("reopen active log: %w", err) } w.file = file return nil } func (w *weeklyLogWriter) generationPath(generation int) string { extension := filepath.Ext(w.path) stem := strings.TrimSuffix(w.path, extension) return fmt.Sprintf("%s.%d%s", stem, generation, extension) } func (w *weeklyLogWriter) reportRotationError(err error) { if err == nil { return } slog.New(slog.NewJSONHandler(w.diagnostic, nil)).Warn( "weekly log rotation failed", "module", "logging", "action", "weekly_log_rotation_failed", "path", w.path, "error.message", err.Error(), ) } func logWeekStart(value time.Time, location *time.Location) time.Time { local := value.In(location) daysSinceMonday := (int(local.Weekday()) + 6) % 7 monday := local.AddDate(0, 0, -daysSinceMonday) return time.Date(monday.Year(), monday.Month(), monday.Day(), 0, 0, 0, 0, location) } func nextLogWeekStart(value time.Time, location *time.Location) time.Time { return logWeekStart(value, location).AddDate(0, 0, 7) } func removeIfExists(path string) error { err := os.Remove(path) if errors.Is(err, fs.ErrNotExist) { return nil } return err }