Compare commits

..

No commits in common. "development" and "v2.1.1" have entirely different histories.

7 changed files with 12 additions and 923 deletions

View File

@ -33,7 +33,7 @@ import (
) )
const ( const (
buildVersion = "v2.1.3" buildVersion = "v2.1.1"
serviceName = "hardlink" serviceName = "hardlink"
pollingFrequency = 8 * time.Second pollingFrequency = 8 * time.Second
) )

View File

@ -295,60 +295,11 @@ func sendAndReadACK(ctx context.Context, port serialTransport, packet []byte, pr
if err := waitForSequence(ctx, processingDelay); err != nil { if err := waitForSequence(ctx, processingDelay); err != nil {
return err return err
} }
return readACK(ctx, port) response := make([]byte, 3)
} if err := readExact(ctx, port, response); err != nil {
return fmt.Errorf("read ACK: %w", err)
// readACK scans sequentially without reading beyond a complete candidate.
// Wrong-address candidates are consumed in full; no bytes persist across calls.
func readACK(ctx context.Context, port serialTransport) (err error) {
const scanLimit = 64
seen := make([]byte, 0, scanLimit)
defer func() {
if err != nil {
log.Warnf("dispenser ACK acquisition failed: address=% X seen=% X error=%v", Address, seen, err)
} }
}() return checkACK(response)
var candidate [3]byte
used := 0
for len(seen) < scanLimit {
if err := ctx.Err(); err != nil {
return err
}
need := 1
if used > 0 {
need = len(candidate) - used
}
need = min(need, scanLimit-len(seen))
chunk := candidate[used : used+need]
n, readErr := port.Read(chunk)
seen = append(seen, chunk[:n]...)
log.Debugf("dispenser ACK RX: address=% X n=%d bytes=% X error=%v", Address, n, chunk[:n], readErr)
if err := ctx.Err(); err != nil {
return err
}
used += n
if used == 1 && candidate[0] != ACK && candidate[0] != NAK {
log.Debugf("dispenser ACK ignored leading byte: % X", candidate[:1])
used = 0
} else if used == len(candidate) {
if len(Address) >= 2 && candidate[1] == Address[0] && candidate[2] == Address[1] {
log.Debugf("dispenser ACK/NAK accepted: address=% X token=% X", Address, candidate[:])
if candidate[0] == ACK && len(seen) > len(candidate) {
log.Warnf("dispenser ACK resynchronized: address=% X ignored=% X accepted=% X", Address, seen[:len(seen)-len(candidate)], candidate[:])
}
return checkACK(candidate[:])
}
log.Debugf("dispenser ACK ignored wrong-address candidate: address=% X candidate=% X", Address, candidate[:])
used = 0
}
if readErr != nil {
return fmt.Errorf("read ACK: %w", readErr)
}
if n == 0 {
return fmt.Errorf("read ACK: %w", io.ErrNoProgress)
}
}
return fmt.Errorf("read ACK: scan limit of %d bytes exhausted", scanLimit)
} }
// queryStatus accepts only the fixed RF/AP payload sizes, before reading a body. // queryStatus accepts only the fixed RF/AP payload sizes, before reading a body.

View File

@ -6,7 +6,6 @@ import (
"errors" "errors"
"fmt" "fmt"
"io" "io"
"os"
"testing" "testing"
"time" "time"
) )
@ -25,146 +24,6 @@ var vendorFrames = []struct {
{"RS", commandRS, []byte{2, 0x30, 0x30, 0, 2, 0x52, 0x53, 3, 2}}, {"RS", commandRS, []byte{2, 0x30, 0x30, 0, 2, 0x52, 0x53, 3, 2}},
} }
func TestACKScannerFragmentation(t *testing.T) {
transportAddress(t)
for _, address := range []string{"00", "15"} {
Address = []byte(address)
wrong := []byte("15")
if address == "15" {
wrong = []byte("00")
}
for _, prefix := range [][]byte{nil, {0x30}, {0x10}, {address[1]}, {ACK, wrong[0], wrong[1], 0x30}, {NAK, wrong[0], wrong[1]}} {
for _, token := range []byte{ACK, NAK} {
wire := append(append([]byte(nil), prefix...), token, Address[0], Address[1])
for mask := 0; mask < 1<<(len(wire)-1); mask++ {
t.Run(fmt.Sprintf("%s/%X/%d", address, wire, mask), func(t *testing.T) {
p := &scriptedTransport{}
start := 0
for i := 1; i < len(wire); i++ {
if mask&(1<<(i-1)) != 0 {
p.chunks = append(p.chunks, wire[start:i])
start = i
}
}
p.chunks = append(p.chunks, wire[start:])
err := dispatchCommand(context.Background(), p, commandFC7, 0)
if (err != nil) != (token == NAK) {
t.Fatalf("dispatchCommand(% X)=%v", wire, err)
}
want := [][]byte{createPacket(Address, commandFC7)}
if token == ACK {
want = append(want, append([]byte{ENQ}, Address...))
}
if len(p.writes) != len(want) {
t.Fatalf("writes=% X want=% X", p.writes, want)
}
for i := range want {
if !bytes.Equal(p.writes[i], want[i]) {
t.Errorf("write=% X want=% X", p.writes[i], want[i])
}
}
if len(p.chunks) != 0 {
t.Errorf("unread token bytes=% X", p.chunks)
}
})
}
}
}
}
}
func TestACKScannerBoundsAndFollowingTransaction(t *testing.T) {
transportAddress(t)
for _, address := range []string{"00", "15"} {
Address = []byte(address)
ack := append([]byte{ACK}, Address...)
for _, leading := range []int{61, 62, 64} {
wire := append(bytes.Repeat([]byte{0xff}, leading), ack...)
p := &scriptedTransport{chunks: [][]byte{wire}}
err := dispatchCommand(context.Background(), p, commandRS, 0)
if (err == nil) != (leading == 61) {
t.Errorf("leading=%d error=%v", leading, err)
}
if leading != 61 && (len(p.writes) != 1 || len(bytes.Join(p.chunks, nil)) != len(wire)-64) {
t.Fatalf("scan exceeded bound: writes=% X remaining=% X", p.writes, p.chunks)
}
}
wire := append(append(append([]byte{0x10}, ack...), ack...), 0xfe)
p := &scriptedTransport{chunks: [][]byte{wire}}
for i := 0; i < 2; i++ {
if err := dispatchCommand(context.Background(), p, commandFC7, 0); err != nil {
t.Fatal(err)
}
}
if len(p.writes) != 4 || !bytes.Equal(bytes.Join(p.chunks, nil), []byte{0xfe}) {
t.Fatalf("transaction boundary lost: writes=% X remaining=% X", p.writes, p.chunks)
}
}
}
type ackReadErrorTransport struct {
*scriptedTransport
reads, failAt int
err error
}
func (p *ackReadErrorTransport) Read(b []byte) (int, error) {
p.reads++
n, err := p.scriptedTransport.Read(b)
if p.reads == p.failAt {
return n, p.err
}
return n, err
}
func TestACKScannerReadErrors(t *testing.T) {
transportAddress(t)
for _, address := range []string{"00", "15"} {
Address = []byte(address)
for _, token := range []byte{ACK, NAK} {
for _, readErr := range []error{io.EOF, os.ErrDeadlineExceeded, io.ErrClosedPipe} {
for _, failAt := range []int{1, 2} {
p := &ackReadErrorTransport{scriptedTransport: &scriptedTransport{chunks: [][]byte{{token, Address[0], Address[1]}}}, failAt: failAt, err: readErr}
err := dispatchCommand(context.Background(), p, commandFC7, 0)
if failAt == 1 && !errors.Is(err, readErr) {
t.Fatalf("incomplete token error=%v want=%v", err, readErr)
}
if failAt == 2 && ((err == nil) != (token == ACK) || errors.Is(err, readErr)) {
t.Fatalf("complete token %02X with read error returned %v", token, err)
}
wantWrites := 1
if failAt == 2 && token == ACK {
wantWrites = 2
}
if len(p.writes) != wantWrites || p.reads != failAt {
t.Fatalf("writes=%d reads=%d want=%d/%d", len(p.writes), p.reads, wantWrites, failAt)
}
}
}
}
for _, wire := range [][]byte{nil, {0x30, ACK, Address[0]}, {0xff, 0xfe}, {ACK, Address[0]}} {
p := &scriptedTransport{chunks: [][]byte{wire}}
if err := dispatchCommand(context.Background(), p, commandFC7, 0); err == nil || len(p.writes) != 1 {
t.Fatalf("incomplete/garbage % X: err=%v writes=% X", wire, err, p.writes)
}
}
for _, cancelAt := range []int{1, 2} {
ctx, cancel := context.WithCancel(context.Background())
p := &ackReadErrorTransport{scriptedTransport: &scriptedTransport{chunks: [][]byte{{ACK, Address[0], Address[1]}}}, failAt: 2, err: io.EOF}
p.afterRead = func() {
if p.reads == cancelAt {
cancel()
}
}
err := dispatchCommand(ctx, p, commandFC7, 0)
cancel()
if !errors.Is(err, context.Canceled) || len(p.writes) != 1 || p.reads != cancelAt {
t.Fatalf("cancel read %d: err=%v writes=%d reads=%d", cancelAt, err, len(p.writes), p.reads)
}
}
}
}
func TestVendorOutboundFrames(t *testing.T) { func TestVendorOutboundFrames(t *testing.T) {
for _, tc := range vendorFrames { for _, tc := range vendorFrames {
t.Run(tc.name, func(t *testing.T) { t.Run(tc.name, func(t *testing.T) {
@ -357,11 +216,7 @@ func TestTransportCancellation(t *testing.T) {
case "processing wait": case "processing wait":
p.afterWrite = cancel p.afterWrite = cancel
case "after ACK": case "after ACK":
p.afterRead = func() { p.afterRead = cancel
if len(p.chunks) == 1 {
cancel()
}
}
case "after ENQ": case "after ENQ":
wantWrites = 2 wantWrites = 2
p.afterWrite = func() { p.afterWrite = func() {
@ -374,7 +229,7 @@ func TestTransportCancellation(t *testing.T) {
reads := 0 reads := 0
p.afterRead = func() { p.afterRead = func() {
reads++ reads++
if reads == 3 { if reads == 2 {
cancel() cancel()
} }
} }

View File

@ -2,22 +2,17 @@ package logging
import ( import (
"fmt" "fmt"
"io"
"os" "os"
"path/filepath"
"time" "time"
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
) )
// SetupLogging ensures the log directory, opens the rotating log writer, and configures logrus. // setupLogging ensures log directory, opens log file, and configures logrus.
func SetupLogging(logDir, serviceName, buildVersion string) (io.WriteCloser, error) { // Returns the *os.File so caller can defer its Close().
if err := os.MkdirAll(logDir, 0o755); err != nil { func SetupLogging(logDir, serviceName, buildVersion string) (*os.File, error) {
return nil, fmt.Errorf("create log directory: %w", err) fileName := logDir + serviceName + ".log"
} f, err := os.OpenFile(fileName, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0666)
fileName := filepath.Join(logDir, serviceName+".log")
f, err := newWeeklyLogWriter(fileName, time.Local, defaultWeeklyLogRuntime())
if err != nil { if err != nil {
return nil, fmt.Errorf("open log file: %w", err) return nil, fmt.Errorf("open log file: %w", err)
} }

View File

@ -1,319 +0,0 @@
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
}

View File

@ -1,387 +0,0 @@
package logging
import (
"bytes"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
log "github.com/sirupsen/logrus"
)
func TestWeeklyLogRotationShiftsAndRetainsExactGenerations(t *testing.T) {
directory := t.TempDir()
path := filepath.Join(directory, "hardlink.log")
writeLogTestFile(t, path, "active")
for generation := 1; generation <= 14; generation++ {
writeLogTestFile(t, weeklyLogTestGeneration(path, generation), fmt.Sprintf("generation-%d", generation))
}
writeLogTestFile(t, weeklyLogTestGeneration(path, 19), "generation-19")
for name, contents := range map[string]string{
"hardlink.backup.14.log": "related-looking backup",
"another.14.log": "unrelated log",
"hardlink.14.log.bak": "different extension",
} {
writeLogTestFile(t, filepath.Join(directory, name), contents)
}
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
setLogTestModTime(t, path, sunday)
writer := newLogTestWriter(t, path, sunday, weeklyLogRuntime{})
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
if err := writer.Close(); err != nil {
t.Fatal(err)
}
requireLogTestContents(t, path, "")
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "active")
requireLogTestContents(t, weeklyLogTestGeneration(path, 2), "generation-1")
requireLogTestContents(t, weeklyLogTestGeneration(path, 13), "generation-12")
for _, generation := range []int{14, 19} {
if _, err := os.Stat(weeklyLogTestGeneration(path, generation)); !errors.Is(err, os.ErrNotExist) {
t.Errorf("generation %d exists after rotation, want removed: %v", generation, err)
}
}
for name, contents := range map[string]string{
"hardlink.backup.14.log": "related-looking backup",
"another.14.log": "unrelated log",
"hardlink.14.log.bak": "different extension",
} {
requireLogTestContents(t, filepath.Join(directory, name), contents)
}
}
func TestWeeklyLogRotationOccursOncePerWeek(t *testing.T) {
directory := t.TempDir()
path := filepath.Join(directory, "hardlink.log")
writeLogTestFile(t, path, "week-zero\n")
sunday := time.Date(2026, time.January, 4, 18, 0, 0, 0, time.UTC)
setLogTestModTime(t, path, sunday)
writer := newLogTestWriter(t, path, sunday, weeklyLogRuntime{})
writer.rotateIfNeeded(time.Date(2026, time.January, 4, 23, 59, 59, 0, time.UTC))
if _, err := os.Stat(weeklyLogTestGeneration(path, 1)); !errors.Is(err, os.ErrNotExist) {
t.Fatalf("rotation before Monday: %v, want no generation", err)
}
writer.rotateIfNeeded(time.Date(2026, time.January, 5, 0, 0, 0, 0, time.UTC))
if _, err := writer.Write([]byte("week-one\n")); err != nil {
t.Fatal(err)
}
writer.rotateIfNeeded(time.Date(2026, time.January, 6, 12, 0, 0, 0, time.UTC))
if err := writer.Close(); err != nil {
t.Fatal(err)
}
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "week-zero\n")
requireLogTestContents(t, path, "week-one\n")
if _, err := os.Stat(weeklyLogTestGeneration(path, 2)); !errors.Is(err, os.ErrNotExist) {
t.Errorf("same-week rotation created generation 2: %v", err)
}
}
func TestWeeklyLogStartupRotatesOnceAfterSeveralMissedWeeks(t *testing.T) {
directory := t.TempDir()
path := filepath.Join(directory, "hardlink.log")
writeLogTestFile(t, path, "old active")
writeLogTestFile(t, weeklyLogTestGeneration(path, 1), "older history")
setLogTestModTime(t, path, time.Date(2026, time.January, 5, 12, 0, 0, 0, time.UTC))
now := time.Date(2026, time.February, 23, 9, 0, 0, 0, time.UTC)
writer := newLogTestWriter(t, path, now, weeklyLogRuntime{})
if err := writer.Close(); err != nil {
t.Fatal(err)
}
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "old active")
requireLogTestContents(t, weeklyLogTestGeneration(path, 2), "older history")
if _, err := os.Stat(weeklyLogTestGeneration(path, 3)); !errors.Is(err, os.ErrNotExist) {
t.Errorf("missed weeks created generation 3: %v", err)
}
}
func TestWeeklyLogSchedulerUsesLocalMondayBoundary(t *testing.T) {
location, err := time.LoadLocation("Europe/London")
if err != nil {
t.Fatal(err)
}
directory := t.TempDir()
path := filepath.Join(directory, "hardlink.log")
writeLogTestFile(t, path, "sunday")
sunday := time.Date(2026, time.March, 29, 0, 0, 0, 0, location)
setLogTestModTime(t, path, sunday)
clock := &atomic.Pointer[time.Time]{}
clock.Store(&sunday)
created := make(chan *logTestTimer, 2)
rotated := make(chan struct{})
runtime := weeklyLogRuntime{
now: func() time.Time { return *clock.Load() },
newTimer: func(duration time.Duration) logRotationTimer {
timer := &logTestTimer{duration: duration, ch: make(chan time.Time, 1)}
created <- timer
return timer
},
rename: func(oldPath, newPath string) error {
err := os.Rename(oldPath, newPath)
if oldPath == path && err == nil {
close(rotated)
}
return err
},
}
writer, err := newWeeklyLogWriter(path, location, runtime)
if err != nil {
t.Fatal(err)
}
timer := <-created
if timer.duration != 23*time.Hour {
t.Fatalf("timer duration = %v, want 23h across DST boundary", timer.duration)
}
monday := time.Date(2026, time.March, 30, 0, 0, 0, 0, location)
clock.Store(&monday)
timer.ch <- monday
<-rotated
if err := writer.Close(); err != nil {
t.Fatal(err)
}
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "sunday")
}
func TestWeeklyLogConcurrentWritesRemainCompleteAcrossRotation(t *testing.T) {
directory := t.TempDir()
path := filepath.Join(directory, "hardlink.log")
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
writer := newLogTestWriter(t, path, sunday, weeklyLogRuntime{})
const goroutines = 8
const recordsPerGoroutine = 100
writeErrors := make(chan error, goroutines)
var writers sync.WaitGroup
writers.Add(goroutines)
for worker := 0; worker < goroutines; worker++ {
go func() {
defer writers.Done()
for record := 0; record < recordsPerGoroutine; record++ {
if _, err := writer.Write([]byte(fmt.Sprintf("%d:%d\n", worker, record))); err != nil {
writeErrors <- err
return
}
}
}()
}
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
writers.Wait()
close(writeErrors)
for err := range writeErrors {
t.Errorf("weeklyLogWriter.Write() error = %v, want nil", err)
}
if err := writer.Close(); err != nil {
t.Fatal(err)
}
contents := readLogTestFile(t, weeklyLogTestGeneration(path, 1)) + readLogTestFile(t, path)
lines := strings.Split(strings.TrimSpace(contents), "\n")
if len(lines) != goroutines*recordsPerGoroutine {
t.Fatalf("record count = %d, want %d", len(lines), goroutines*recordsPerGoroutine)
}
seen := make(map[string]bool, len(lines))
for _, line := range lines {
if _, _, ok := strings.Cut(line, ":"); !ok {
t.Fatalf("record %q has no separator", line)
}
seen[line] = true
}
if len(seen) != goroutines*recordsPerGoroutine {
t.Errorf("unique record count = %d, want %d", len(seen), goroutines*recordsPerGoroutine)
}
}
func TestWeeklyLogCloseWaitsForRotationAndIsIdempotent(t *testing.T) {
directory := t.TempDir()
path := filepath.Join(directory, "hardlink.log")
writeLogTestFile(t, path, "active")
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
setLogTestModTime(t, path, sunday)
started := make(chan struct{})
release := make(chan struct{})
runtime := weeklyLogRuntime{rename: func(oldPath, newPath string) error {
if oldPath == path {
close(started)
<-release
}
return os.Rename(oldPath, newPath)
}}
writer := newLogTestWriter(t, path, sunday, runtime)
rotationDone := make(chan struct{})
go func() {
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
close(rotationDone)
}()
<-started
closeDone := make(chan error, 1)
go func() { closeDone <- writer.Close() }()
select {
case err := <-closeDone:
t.Fatalf("weeklyLogWriter.Close() during rotation returned %v, want blocked", err)
case <-time.After(20 * time.Millisecond):
}
close(release)
<-rotationDone
if err := <-closeDone; err != nil {
t.Fatal(err)
}
if err := writer.Close(); err != nil {
t.Errorf("second weeklyLogWriter.Close() error = %v, want nil", err)
}
if _, err := writer.Write([]byte("after close")); !errors.Is(err, os.ErrClosed) {
t.Errorf("weeklyLogWriter.Write() after close error = %v, want os.ErrClosed", err)
}
}
func TestWeeklyLogRotationFailureIsNonRecursiveAndPreservesWrites(t *testing.T) {
directory := t.TempDir()
path := filepath.Join(directory, "hardlink.log")
writeLogTestFile(t, path, "before\n")
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
setLogTestModTime(t, path, sunday)
var diagnostics bytes.Buffer
runtime := weeklyLogRuntime{
rename: func(oldPath, newPath string) error {
if oldPath == path {
return errors.New("injected rename failure")
}
return os.Rename(oldPath, newPath)
},
diagnostic: &diagnostics,
}
writer := newLogTestWriter(t, path, sunday, runtime)
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
if _, err := writer.Write([]byte("after\n")); err != nil {
t.Fatalf("weeklyLogWriter.Write() after failed rotation error = %v, want nil", err)
}
if err := writer.Close(); err != nil {
t.Fatal(err)
}
requireLogTestContents(t, path, "before\nafter\n")
if !strings.Contains(diagnostics.String(), "weekly_log_rotation_failed") || !strings.Contains(diagnostics.String(), "injected rename failure") {
t.Errorf("rotation diagnostic = %q, want action and injected error", diagnostics.String())
}
if strings.Contains(readLogTestFile(t, path), "weekly log rotation failed") {
t.Error("rotation diagnostic recursively entered active log")
}
}
func TestSetupLoggingKeepsFilenameAndJSONFormat(t *testing.T) {
logger := log.StandardLogger()
previousOutput := logger.Out
previousFormatter := logger.Formatter
previousLevel := logger.Level
t.Cleanup(func() {
logger.SetOutput(previousOutput)
logger.SetFormatter(previousFormatter)
logger.SetLevel(previousLevel)
})
directory := filepath.Join(t.TempDir(), "nested", "logs")
writer, err := SetupLogging(directory, "hardlink", "v-test")
if err != nil {
t.Fatal(err)
}
log.WithField("probe", "value").Info("test record")
if err := writer.Close(); err != nil {
t.Fatal(err)
}
contents := readLogTestFile(t, filepath.Join(directory, "hardlink.log"))
for _, want := range []string{`"buildVersion":"v-test"`, `"msg":"Logging initialized"`, `"probe":"value"`, `"msg":"test record"`} {
if !strings.Contains(contents, want) {
t.Errorf("hardlink.log contents missing %q: %s", want, contents)
}
}
if _, err := os.Stat(filepath.Join(directory, "hardlink.1.log")); !errors.Is(err, os.ErrNotExist) {
t.Errorf("SetupLogging() created archive during current week: %v", err)
}
}
type logTestTimer struct {
duration time.Duration
ch chan time.Time
stopped atomic.Bool
}
func (t *logTestTimer) C() <-chan time.Time {
return t.ch
}
func (t *logTestTimer) Stop() bool {
return !t.stopped.Swap(true)
}
func newLogTestWriter(t *testing.T, path string, now time.Time, overrides weeklyLogRuntime) *weeklyLogWriter {
t.Helper()
runtime := defaultWeeklyLogRuntime()
runtime.now = func() time.Time { return now }
if overrides.now != nil {
runtime.now = overrides.now
}
if overrides.newTimer != nil {
runtime.newTimer = overrides.newTimer
}
if overrides.rename != nil {
runtime.rename = overrides.rename
}
if overrides.diagnostic != nil {
runtime.diagnostic = overrides.diagnostic
}
writer, err := newWeeklyLogWriter(path, now.Location(), runtime)
if err != nil {
t.Fatal(err)
}
return writer
}
func weeklyLogTestGeneration(path string, generation int) string {
extension := filepath.Ext(path)
return fmt.Sprintf("%s.%d%s", strings.TrimSuffix(path, extension), generation, extension)
}
func writeLogTestFile(t *testing.T, path, contents string) {
t.Helper()
if err := os.WriteFile(path, []byte(contents), 0o600); err != nil {
t.Fatal(err)
}
}
func setLogTestModTime(t *testing.T, path string, value time.Time) {
t.Helper()
if err := os.Chtimes(path, value, value); err != nil {
t.Fatal(err)
}
}
func readLogTestFile(t *testing.T, path string) string {
t.Helper()
contents, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
return string(contents)
}
func requireLogTestContents(t *testing.T, path, want string) {
t.Helper()
if got := readLogTestFile(t, path); got != want {
t.Errorf("%s contents = %q, want %q", filepath.Base(path), got, want)
}
}
var _ io.WriteCloser = (*weeklyLogWriter)(nil)

View File

@ -2,12 +2,6 @@
builtVersion is a const in main.go builtVersion is a const in main.go
#### v2.1.3 - 30 September 2026
feat(logging): add weekly log rotation and retention
#### v2.1.2 - 30 September 2026
fix(dispenser): scan ACK responses without assuming 3-byte reads
#### v2.1.1 - 28 September 2026 #### v2.1.1 - 28 September 2026
feat(dispenser): add worker-owned idle card prestaging feat(dispenser): add worker-owned idle card prestaging