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)