From 625c982dee21a2471d71923731754cbb1dffd493 Mon Sep 17 00:00:00 2001 From: eric <1048315650@qq.com> Date: Tue, 14 Apr 2026 09:56:22 +0800 Subject: [PATCH] =?UTF-8?q?feat(diskq):=20=E6=B8=85=E7=90=86=E5=9D=8F?= =?UTF-8?q?=E7=9B=98=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- diskq/diskq.go | 41 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/diskq/diskq.go b/diskq/diskq.go index 24c4201..453299b 100644 --- a/diskq/diskq.go +++ b/diskq/diskq.go @@ -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 } @@ -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