Skip to content

Commit 963d19a

Browse files
authored
Fix incorrect double-dequeue after restarting (#12)
1 parent dba729e commit 963d19a

File tree

2 files changed

+30
-1
lines changed

2 files changed

+30
-1
lines changed

queue_test.go

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -400,3 +400,32 @@ func TestQueueCorruptedWritingFile(t *testing.T) {
400400
var e entry.Entry
401401
require.False(t, q.Dequeue(&e))
402402
}
403+
404+
func TestQueueReopen(t *testing.T) {
405+
dataDir := filepath.Join(tmpDir, "pqueue_reopen")
406+
_ = os.RemoveAll(dataDir)
407+
err := os.MkdirAll(dataDir, 0o777)
408+
require.NoError(t, err)
409+
defer func() {
410+
_ = os.RemoveAll(dataDir)
411+
}()
412+
413+
q, err := New(dataDir, 3)
414+
require.NoError(t, err)
415+
416+
require.NoError(t, q.Enqueue([]byte{1, 2, 3}))
417+
418+
var e entry.Entry
419+
require.True(t, q.Dequeue(&e))
420+
require.EqualValues(t, e, []byte{1, 2, 3})
421+
422+
err = q.Close()
423+
require.NoError(t, err)
424+
425+
q, err = New(dataDir, 3)
426+
require.NoError(t, err)
427+
428+
// only one value was enqueued and it was already dequeued, so there shouldn't be anything else
429+
require.False(t, q.Dequeue(&e))
430+
q.Close()
431+
}

utils.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ func loadFileInfos(dir string, infoExtractor func(os.DirEntry) (os.FileInfo, err
5959

6060
files := make([]file, 0, len(fileList))
6161
for i := range fileList {
62-
if strings.HasPrefix(fileList[i].Name(), segPrefix) {
62+
if strings.HasPrefix(fileList[i].Name(), segPrefix) && !strings.HasSuffix(fileList[i].Name(), segOffsetFileSuffix) {
6363
info, e := infoExtractor(fileList[i])
6464
if e != nil {
6565
return nil, e

0 commit comments

Comments
 (0)