From e73fe765c3b1d87b4e38180cc46e94f6abec5856 Mon Sep 17 00:00:00 2001 From: npc0-hue Date: Wed, 16 Sep 2026 13:34:45 +0800 Subject: [PATCH] Recover log streams stuck on straddling spool segments A segment that the platform already acknowledged could still be extended by the next tailed line, so its checksum covered entries stored under a different batch. The platform then rejected the same body every second while RejectStreamAfter skipped it, because it only quarantined segments that start after the platform latest. Every newly allocated line was quarantined by the following recovery pass, so live log ingest never resumed. Quarantine pending segments that reach beyond the acknowledged range, never extend a segment the platform already stored, and log the gap recovery so a stalled stream is diagnosable from the spool alone. --- spool/log_spool.go | 19 +++++++++++-- spool/log_spool_test.go | 63 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 80 insertions(+), 2 deletions(-) diff --git a/spool/log_spool.go b/spool/log_spool.go index 5c2f9ec..213a16c 100644 --- a/spool/log_spool.go +++ b/spool/log_spool.go @@ -230,9 +230,14 @@ func (spool LogSpool) mergeLatest(batch durableLogBatch, checksum func([]protoco if err != nil { return durableLogBatch{}, pendingLogSegment{}, false, err } + acknowledged := spool.watermarks[batch.LogStreamID].Acknowledged for index := len(pending) - 1; index >= 0; index-- { previous := pending[index] - if spool.inflight[previous.path] || !compatibleLogBatch(previous.batch, batch) || previous.batch.LastSeq+1 != batch.FirstSeq || len(previous.batch.Entries)+len(batch.Entries) > maxAggregatedLogEntries || !compatibleSourceCursor(previous.batch.SourceCursor, batch.SourceCursor) { + // Extending a segment the platform already stored would rewrite its + // checksum for a range that can no longer be uploaded, and the segment + // then becomes permanently unresendable. Leave it for the deduplicating + // retry and start a new segment instead. + if spool.inflight[previous.path] || previous.batch.LastSeq <= acknowledged || !compatibleLogBatch(previous.batch, batch) || previous.batch.LastSeq+1 != batch.FirstSeq || len(previous.batch.Entries)+len(batch.Entries) > maxAggregatedLogEntries || !compatibleSourceCursor(previous.batch.SourceCursor, batch.SourceCursor) { continue } merged := previous.batch @@ -652,10 +657,12 @@ func (spool LogSpool) Flush(ctx context.Context, client LogBatchClient) (int, er if progressErr != nil { return acknowledged, err } + log.Printf("RUN phase=log_spool.flush status=sequence_gap stream=%s platformLatestSeq=%d batchFirstSeq=%d batchLastSeq=%d", safeSpoolStreamID(batch.LogStreamID), latestSeq, batch.FirstSeq, batch.LastSeq) rejected, rejectErr := spool.RejectStreamAfter(batch.LogStreamID, latestSeq, permanent.Reason) if rejectErr != nil { return acknowledged, rejectErr } + log.Printf("RUN phase=log_spool.flush status=sequence_gap_recovered stream=%s platformLatestSeq=%d rejectedSegments=%d", safeSpoolStreamID(batch.LogStreamID), latestSeq, rejected) if resetErr := spool.ResetStreamWatermark(batch.LogStreamID, latestSeq); resetErr != nil { return acknowledged, resetErr } @@ -709,7 +716,15 @@ func (spool LogSpool) RejectStreamAfter(logStreamID string, latestSeq uint64, re rejected := 0 for _, segment := range segments { batch := segment.batch.LogBatchIngestRequest - if batch.LogStreamID != logStreamID || batch.FirstSeq <= latestSeq { + // A pending segment is only resendable when it lies entirely inside the + // platform's acknowledged range (it is then deduplicated by sequence + // checksum). Every segment that reaches beyond the platform latest, or + // that straddles it, can never be accepted again: its checksum covers + // entries the platform already stored under a different batch. Keeping + // such a straddling segment pending made the spool resend the same + // rejected body forever, and every newly allocated line was quarantined + // by the next recovery pass, so the stream never resumed. + if batch.LogStreamID != logStreamID || batch.LastSeq <= latestSeq { continue } if err := spool.moveLogSegmentToRejectedLocked(segment.path, reason); err != nil { diff --git a/spool/log_spool_test.go b/spool/log_spool_test.go index a01302b..b5518f1 100644 --- a/spool/log_spool_test.go +++ b/spool/log_spool_test.go @@ -188,6 +188,69 @@ func TestLogSpoolSequenceGapResetsWatermarkToPlatformLatest(t *testing.T) { } } +func TestLogSpoolSequenceGapRejectsStraddlingAcknowledgedSegment(t *testing.T) { + root := t.TempDir() + logSpool, err := NewLogSpool(root) + if err != nil { + t.Fatalf("new log spool: %v", err) + } + // The acknowledged range ends inside this segment: sequences 10..13 were + // extended after the platform had already stored sequence 10. + if err := logSpool.Enqueue(validSpoolLogBatch(10, 13)); err != nil { + t.Fatalf("enqueue straddling batch: %v", err) + } + if err := logSpool.Enqueue(validSpoolLogBatch(14, 14)); err != nil { + t.Fatalf("enqueue later batch: %v", err) + } + client := &recoveringSequenceGapLogBatchClient{latestSeq: 10} + flushed, err := logSpool.Flush(context.Background(), client) + if err != nil || flushed != 2 { + t.Fatalf("recover sequence gap: flushed=%d err=%v", flushed, err) + } + pending, err := logSpool.Pending() + if err != nil || len(pending) != 0 { + t.Fatalf("expected unresendable segments rejected, pending=%+v err=%v", pending, err) + } + rejected, err := os.ReadDir(filepath.Join(root, "logs-rejected")) + if err != nil || len(rejected) != 2 { + t.Fatalf("expected two rejected segments, rejected=%+v err=%v", rejected, err) + } + checksum := func(entries []protocol.LogEntry) (string, error) { return fmt.Sprintf("sha256:%d", len(entries)), nil } + next := validSpoolLogBatch(0, 0) + sequence, appended, err := logSpool.EnqueueNextAggregated(context.Background(), next, nil, nil, checksum) + if err != nil || !appended || sequence != 11 { + t.Fatalf("expected allocation to resume after platform latest, sequence=%d appended=%t err=%v", sequence, appended, err) + } +} + +func TestLogSpoolDoesNotExtendAcknowledgedSegment(t *testing.T) { + root := t.TempDir() + logSpool, err := NewLogSpool(root) + if err != nil { + t.Fatalf("new log spool: %v", err) + } + checksum := func(entries []protocol.LogEntry) (string, error) { return fmt.Sprintf("sha256:%d", len(entries)), nil } + if err := logSpool.EnqueueAggregated(validSpoolLogBatch(1, 1), checksum); err != nil { + t.Fatalf("enqueue acknowledged segment: %v", err) + } + logSpool.mu.Lock() + watermark := logSpool.watermarks["log-1"] + watermark.Acknowledged = 1 + logSpool.watermarks["log-1"] = watermark + logSpool.mu.Unlock() + if err := logSpool.EnqueueAggregated(validSpoolLogBatch(2, 2), checksum); err != nil { + t.Fatalf("enqueue after acknowledged segment: %v", err) + } + entries, err := os.ReadDir(filepath.Join(root, "logs")) + if err != nil || len(entries) != 2 { + t.Fatalf("expected the acknowledged segment untouched and a new segment, entries=%+v err=%v", entries, err) + } + pending, err := logSpool.Pending() + if err != nil || len(pending) != 2 || pending[0].FirstSeq != 1 || pending[0].LastSeq != 1 || pending[1].FirstSeq != 2 || pending[1].LastSeq != 2 { + t.Fatalf("acknowledged segment was extended: pending=%+v err=%v", pending, err) + } +} + func TestLogSpoolRestoresPendingAllocationWithoutWatermark(t *testing.T) { root := t.TempDir() first, err := NewLogSpool(root)