Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions diskq/diskq.go
Original file line number Diff line number Diff line change
Expand Up @@ -243,6 +243,10 @@ func newDiskq(name string, dataDIR string, maxBytesPerFile int64,
d.logf(ERROR, "DISKQUEUE(%s) failed to retrieveMetaData - %s", d.name, err)
}

// clean up stale temp and bad files left by previous crashed instances
d.cleanupTempFiles()
d.cleanupBadFiles()

go d.ioLoop()
return d
}
Expand Down Expand Up @@ -348,9 +352,46 @@ func (d *DiskQueue) deleteAllFiles() error {
return innerErr
}

d.cleanupTempFiles()
d.cleanupBadFiles()

return err
}

// cleanupTempFiles removes stale .tmp metadata files left by crashed processes
func (d *DiskQueue) cleanupTempFiles() {
pattern := fmt.Sprintf(path.Join(d.dataDIR, "%s.diskqueue.meta.dat.*.tmp"), d.name)
matches, err := filepath.Glob(pattern)
if err != nil {
d.logf(WARN, "DISKQUEUE(%s) failed to glob temp files - %s", d.name, err)
return
}
for _, fn := range matches {
if err = os.Remove(fn); err != nil && !os.IsNotExist(err) {
d.logf(WARN, "DISKQUEUE(%s) failed to remove temp file %s - %s", d.name, fn, err)
} else {
d.logf(INFO, "DISKQUEUE(%s) removed stale temp file %s", d.name, fn)
}
}
}

// cleanupBadFiles removes .bad files that were created when corrupt data files were encountered
func (d *DiskQueue) cleanupBadFiles() {
pattern := fmt.Sprintf(path.Join(d.dataDIR, "%s.diskqueue.*.dat.bad"), d.name)
matches, err := filepath.Glob(pattern)
if err != nil {
d.logf(WARN, "DISKQUEUE(%s) failed to glob bad files - %s", d.name, err)
return
}
for _, fn := range matches {
if err = os.Remove(fn); err != nil && !os.IsNotExist(err) {
d.logf(WARN, "DISKQUEUE(%s) failed to remove bad file %s - %s", d.name, fn, err)
} else {
d.logf(INFO, "DISKQUEUE(%s) removed bad file %s", d.name, fn)
}
}
}

func (d *DiskQueue) skipToNextRWFile() error {
var err error

Expand Down
Loading