Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0987817aae |
@@ -1,15 +1,7 @@
|
|||||||
package reportbuf
|
package reportbuf
|
||||||
|
|
||||||
import "time"
|
|
||||||
|
|
||||||
// Flush writes the buffered reports to a file now, as the periodic
|
// Flush writes the buffered reports to a file now, as the periodic
|
||||||
// flush does, so tests need not wait a minute for it.
|
// flush does, so tests need not wait a minute for it.
|
||||||
func (b *Buffer) Flush() error {
|
func (b *Buffer) Flush() error {
|
||||||
return b.flushLocked()
|
return b.flushLocked()
|
||||||
}
|
}
|
||||||
|
|
||||||
// StopClock makes every report file the buffer writes from now on
|
|
||||||
// carry the timestamp at, as if all were written in one millisecond.
|
|
||||||
func (b *Buffer) StopClock(at time.Time) {
|
|
||||||
b.now = func() time.Time { return at }
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -58,9 +58,6 @@ type Buffer struct {
|
|||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
maxBytes int64
|
maxBytes int64
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
// now is the clock report files are named by: time.Now, except
|
|
||||||
// in tests that need two flushes to share a timestamp.
|
|
||||||
now func() time.Time
|
|
||||||
// seq numbers the report files, so that two named in the same
|
// seq numbers the report files, so that two named in the same
|
||||||
// millisecond still get different names.
|
// millisecond still get different names.
|
||||||
seq atomic.Uint64
|
seq atomic.Uint64
|
||||||
@@ -87,7 +84,6 @@ func New(
|
|||||||
done: make(chan struct{}),
|
done: make(chan struct{}),
|
||||||
log: params.Logger.Get(),
|
log: params.Logger.Get(),
|
||||||
maxBytes: params.Config.DataDirMaxBytes,
|
maxBytes: params.Config.DataDirMaxBytes,
|
||||||
now: time.Now,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
lc.Append(fx.Hook{
|
lc.Append(fx.Hook{
|
||||||
@@ -222,7 +218,7 @@ func (b *Buffer) drainBuf() []byte {
|
|||||||
func (b *Buffer) writeFile(data []byte) error {
|
func (b *Buffer) writeFile(data []byte) error {
|
||||||
// The timestamp comes first, so the names sort by time; the number
|
// The timestamp comes first, so the names sort by time; the number
|
||||||
// after it tells apart files named in the same millisecond.
|
// after it tells apart files named in the same millisecond.
|
||||||
ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z")
|
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
|
||||||
name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix)
|
name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix)
|
||||||
path := filepath.Join(b.dataDir, name)
|
path := filepath.Join(b.dataDir, name)
|
||||||
|
|
||||||
|
|||||||
@@ -314,23 +314,40 @@ func TestConcurrentAppendsStopAtCap(t *testing.T) {
|
|||||||
// as a flush for size and the final flush at shutdown can: each flush
|
// as a flush for size and the final flush at shutdown can: each flush
|
||||||
// must write a file of its own, and the files must hold every report.
|
// must write a file of its own, and the files must hold every report.
|
||||||
func TestTwoFlushesInOneMillisecond(t *testing.T) {
|
func TestTwoFlushesInOneMillisecond(t *testing.T) {
|
||||||
const flushes = 2
|
// A pair of flushes may straddle a millisecond, which proves
|
||||||
|
// nothing, so pairs are flushed until one falls within one.
|
||||||
|
const maxPairs = 1000
|
||||||
|
|
||||||
dir := t.TempDir()
|
dir := t.TempDir()
|
||||||
t.Setenv("DATA_DIR", dir)
|
t.Setenv("DATA_DIR", dir)
|
||||||
|
|
||||||
buf := startBuffer(t)
|
buf := startBuffer(t)
|
||||||
buf.StopClock(time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC))
|
flushes := 0
|
||||||
|
|
||||||
for id := 1; id <= flushes; id++ {
|
for pair := 1; ; pair++ {
|
||||||
err := buf.Append(map[string]int{"id": id})
|
start := time.Now().Truncate(time.Millisecond)
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("append report %d: %v", id, err)
|
for range 2 {
|
||||||
|
flushes++
|
||||||
|
|
||||||
|
err := buf.Append(map[string]int{"id": flushes})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("append report %d: %v", flushes, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = buf.Flush()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("flush %d: %v", flushes, err)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
err = buf.Flush()
|
if time.Now().Truncate(time.Millisecond).Equal(start) {
|
||||||
if err != nil {
|
break
|
||||||
t.Fatalf("flush %d: %v", id, err)
|
}
|
||||||
|
|
||||||
|
if pair == maxPairs {
|
||||||
|
t.Fatalf("no pair of flushes fell within one millisecond "+
|
||||||
|
"in %d tries", maxPairs)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user