Files
run/spool/log_spool.go
T

926 lines
31 KiB
Go

package spool
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"time"
"browser.local/run/protocol"
)
type LogSpool struct {
dir string
statePath string
mu *sync.Mutex
watermarks map[string]logStreamWatermark
inflight map[string]bool
}
const (
maxAggregatedLogEntries = 128
maxDurableLogSegmentBytes = 8 * 1024 * 1024
)
type logStreamWatermark struct {
Allocated uint64 `json:"allocated"`
Acknowledged uint64 `json:"acknowledged"`
SourceOffset int64 `json:"sourceOffset,omitempty"`
}
type logSpoolState struct {
Streams map[string]logStreamWatermark `json:"streams"`
}
type LogSourceCursor struct {
StartOffset int64 `json:"startOffset"`
EndOffset int64 `json:"endOffset"`
}
type durableLogBatch struct {
protocol.LogBatchIngestRequest
SourceCursor *LogSourceCursor `json:"_sourceCursor,omitempty"`
}
type pendingLogSegment struct {
path string
modTime time.Time
batch durableLogBatch
}
type pendingLogSegmentWatermark struct {
logStreamID string
lastSeq uint64
sourceCursor *LogSourceCursor
}
type pendingSegmentScanOptions struct {
diagnosticPhase string
}
type pendingSegmentScanStats struct {
candidates int
loaded int
skipped int
quarantined int
bytes int64
}
func NewLogSpool(dir string) (LogSpool, error) {
startedAt := time.Now()
if strings.TrimSpace(dir) == "" {
return LogSpool{}, fmt.Errorf("spool directory is required")
}
logDir := filepath.Join(dir, "logs")
if err := os.MkdirAll(logDir, 0o755); err != nil {
return LogSpool{}, fmt.Errorf("create log spool: %w", err)
}
spool := LogSpool{dir: logDir, statePath: filepath.Join(dir, "log-watermarks.json"), mu: &sync.Mutex{}, watermarks: map[string]logStreamWatermark{}, inflight: map[string]bool{}}
if err := spool.loadWatermarks(); err != nil {
return LogSpool{}, err
}
log.Printf("RUN phase=log_spool.init status=watermarks_loaded streams=%d durationMs=%d", len(spool.watermarks), time.Since(startedAt).Milliseconds())
if err := spool.inspectPendingSegmentFiles(); err != nil {
return LogSpool{}, err
}
log.Printf("RUN phase=log_spool.init status=complete streams=%d durationMs=%d", len(spool.watermarks), time.Since(startedAt).Milliseconds())
return spool, nil
}
func (spool LogSpool) inspectPendingSegmentFiles() error {
spool.mu.Lock()
defer spool.mu.Unlock()
entries, err := os.ReadDir(spool.dir)
if err != nil {
return fmt.Errorf("read log spool: %w", err)
}
paths, stats := spool.pendingSegmentPathsLocked(entries, pendingSegmentScanOptions{diagnosticPhase: "startup_inspect"})
log.Printf("RUN phase=log_spool.inspect status=complete candidates=%d accepted=%d skipped=%d quarantined=%d bytes=%d", stats.candidates, len(paths), stats.skipped, stats.quarantined, stats.bytes)
return nil
}
// restorePendingWatermarks supports upgrades from spools created before
// per-stream watermarks existed. Pending durable segments are still an
// allocation fact and must win over a remote progress lookup.
func (spool *LogSpool) restorePendingWatermarks() error {
startedAt := time.Now()
spool.mu.Lock()
defer spool.mu.Unlock()
pending, stats, err := spool.pendingSegmentWatermarksLocked(pendingSegmentScanOptions{diagnosticPhase: "startup_restore"})
if err != nil {
return err
}
changed := false
for _, segment := range pending {
watermark := spool.watermarks[segment.logStreamID]
if segment.lastSeq > watermark.Allocated {
watermark.Allocated = segment.lastSeq
changed = true
}
if segment.sourceCursor != nil && segment.sourceCursor.EndOffset > watermark.SourceOffset {
watermark.SourceOffset = segment.sourceCursor.EndOffset
changed = true
}
spool.watermarks[segment.logStreamID] = watermark
}
if changed {
if err := spool.persistWatermarksLocked(); err != nil {
return err
}
}
log.Printf("RUN phase=log_spool.restore status=complete pending=%d candidates=%d skipped=%d quarantined=%d changed=%t bytes=%d durationMs=%d", stats.loaded, stats.candidates, stats.skipped, stats.quarantined, changed, stats.bytes, time.Since(startedAt).Milliseconds())
return nil
}
func (spool LogSpool) NextSequence(ctx context.Context, streamID string, recover func(context.Context, string) (uint64, error)) (uint64, error) {
spool.mu.Lock()
defer spool.mu.Unlock()
watermark, known := spool.watermarks[streamID]
if !known {
restored, restoredKnown, err := spool.restorePendingWatermarkLocked(streamID)
if err != nil {
return 0, err
}
if restoredKnown {
watermark = restored
known = true
}
}
if !known && recover != nil {
latest, err := recover(ctx, streamID)
if err != nil {
return 0, err
}
watermark = logStreamWatermark{Allocated: latest, Acknowledged: latest}
}
watermark.Allocated++
spool.watermarks[streamID] = watermark
if err := spool.persistWatermarksLocked(); err != nil {
return 0, err
}
return watermark.Allocated, nil
}
func (spool LogSpool) Enqueue(batch protocol.LogBatchIngestRequest) error {
return spool.enqueue(batch, nil)
}
// EnqueueAggregated extends the newest compatible durable segment so callers
// do not create one upload request for every process-output line.
func (spool LogSpool) EnqueueAggregated(batch protocol.LogBatchIngestRequest, checksum func([]protocol.LogEntry) (string, error)) error {
return spool.enqueue(batch, checksum)
}
// EnqueueNextAggregated allocates the next sequence and commits its batch
// under one spool lock. The durable segment is the source of truth; startup
// restores its watermark if the separate watermark snapshot was interrupted.
func (spool LogSpool) EnqueueNextAggregated(ctx context.Context, batch protocol.LogBatchIngestRequest, cursor *LogSourceCursor, recover func(context.Context, string) (uint64, error), checksum func([]protocol.LogEntry) (string, error)) (uint64, bool, error) {
if len(batch.Entries) != 1 {
return 0, false, fmt.Errorf("next aggregated log batch requires exactly one entry")
}
if cursor != nil && (cursor.StartOffset < 0 || cursor.EndOffset <= cursor.StartOffset) {
return 0, false, fmt.Errorf("source cursor range is invalid")
}
spool.mu.Lock()
defer spool.mu.Unlock()
watermark, known := spool.watermarks[batch.LogStreamID]
if !known {
restored, restoredKnown, err := spool.restorePendingWatermarkLocked(batch.LogStreamID)
if err != nil {
return 0, false, err
}
if restoredKnown {
watermark = restored
known = true
}
}
if !known && recover != nil {
latest, err := recover(ctx, batch.LogStreamID)
if err != nil {
return 0, false, err
}
watermark = logStreamWatermark{Allocated: latest, Acknowledged: latest}
}
if cursor != nil && cursor.EndOffset <= watermark.SourceOffset {
return watermark.Allocated, false, nil
}
sequence := watermark.Allocated + 1
batch.FirstSeq = sequence
batch.LastSeq = sequence
batch.Entries[0].Seq = sequence
var err error
batch.Checksum, err = checksum(batch.Entries)
if err != nil {
return 0, false, err
}
if err := spool.writeAggregatedLocked(durableLogBatch{LogBatchIngestRequest: batch, SourceCursor: cursor}, checksum); err != nil {
return 0, false, err
}
watermark.Allocated = sequence
if cursor != nil {
watermark.SourceOffset = cursor.EndOffset
}
spool.watermarks[batch.LogStreamID] = watermark
// The batch is already synced and durable. A later startup reconstructs
// this watermark from pending segments if this snapshot cannot be written.
_ = spool.persistWatermarksLocked()
return sequence, true, nil
}
func (spool LogSpool) enqueue(batch protocol.LogBatchIngestRequest, checksum func([]protocol.LogEntry) (string, error)) error {
spool.mu.Lock()
defer spool.mu.Unlock()
watermark := spool.watermarks[batch.LogStreamID]
if batch.LastSeq > watermark.Allocated {
watermark.Allocated = batch.LastSeq
spool.watermarks[batch.LogStreamID] = watermark
if err := spool.persistWatermarksLocked(); err != nil {
return err
}
}
return spool.writeAggregatedLocked(durableLogBatch{LogBatchIngestRequest: batch}, checksum)
}
func (spool LogSpool) writeAggregatedLocked(batch durableLogBatch, checksum func([]protocol.LogEntry) (string, error)) error {
if checksum != nil && len(batch.Entries) == 1 {
if merged, previous, ok, err := spool.mergeLatest(batch, checksum); err != nil {
return err
} else if ok {
return spool.replaceBatch(previous.path, merged)
}
}
return spool.writeBatch(batch)
}
func (spool LogSpool) mergeLatest(batch durableLogBatch, checksum func([]protocol.LogEntry) (string, error)) (durableLogBatch, pendingLogSegment, bool, error) {
pending, err := spool.pendingSegmentsLocked()
if err != nil {
return durableLogBatch{}, pendingLogSegment{}, false, err
}
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) {
continue
}
merged := previous.batch
merged.LastSeq = batch.LastSeq
merged.Entries = append(append([]protocol.LogEntry(nil), previous.batch.Entries...), batch.Entries...)
if merged.SourceCursor != nil {
cursor := *merged.SourceCursor
cursor.EndOffset = batch.SourceCursor.EndOffset
merged.SourceCursor = &cursor
}
merged.Checksum, err = checksum(merged.Entries)
if err != nil {
return durableLogBatch{}, pendingLogSegment{}, false, err
}
return merged, previous, true, nil
}
return durableLogBatch{}, pendingLogSegment{}, false, nil
}
func compatibleLogBatch(previous durableLogBatch, next durableLogBatch) bool {
return previous.LogStreamID == next.LogStreamID && previous.ServerInstanceID == next.ServerInstanceID && previous.StreamKey == next.StreamKey && previous.Source == next.Source && previous.Compression == next.Compression && previous.LogSessionID == next.LogSessionID && previous.SessionStartedAt.Equal(next.SessionStartedAt)
}
func compatibleSourceCursor(previous *LogSourceCursor, next *LogSourceCursor) bool {
if previous == nil || next == nil {
return previous == nil && next == nil
}
return next.StartOffset >= previous.EndOffset
}
func (spool LogSpool) writeBatch(batch durableLogBatch) error {
path := spool.batchPath(batch)
if existing, err := readDurableLogBatch(path); err == nil {
if existing.LogStreamID == batch.LogStreamID && existing.FirstSeq == batch.FirstSeq && existing.LastSeq == batch.LastSeq && existing.Checksum == batch.Checksum {
return nil
}
return fmt.Errorf("log spool segment conflicts with committed batch")
} else if !os.IsNotExist(err) {
return err
}
return spool.writeBatchAt(path, batch)
}
func (spool LogSpool) replaceBatch(path string, batch durableLogBatch) error {
if strings.TrimSpace(path) == "" {
return fmt.Errorf("log spool segment path is required")
}
return spool.writeBatchAt(path, batch)
}
func (spool LogSpool) writeBatchAt(path string, batch durableLogBatch) error {
tmp := path + ".tmp"
file, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600)
if err != nil {
return fmt.Errorf("open log spool segment: %w", err)
}
encodeErr := json.NewEncoder(file).Encode(batch)
closeErr := file.Close()
if encodeErr != nil {
_ = os.Remove(tmp)
return fmt.Errorf("encode log spool segment: %w", encodeErr)
}
if closeErr != nil {
_ = os.Remove(tmp)
return fmt.Errorf("close log spool segment: %w", closeErr)
}
if err := syncFile(tmp); err != nil {
_ = os.Remove(tmp)
return err
}
if err := os.Rename(tmp, path); err != nil {
_ = os.Remove(tmp)
return fmt.Errorf("commit log spool segment: %w", err)
}
return nil
}
func (spool LogSpool) restorePendingWatermarkLocked(streamID string) (logStreamWatermark, bool, error) {
if strings.TrimSpace(streamID) == "" {
return logStreamWatermark{}, false, nil
}
startedAt := time.Now()
pending, stats, err := spool.pendingSegmentWatermarksLocked(pendingSegmentScanOptions{diagnosticPhase: "lazy_restore"})
if err != nil {
return logStreamWatermark{}, false, err
}
watermark := spool.watermarks[streamID]
changed := false
for _, segment := range pending {
if segment.logStreamID != streamID {
continue
}
if segment.lastSeq > watermark.Allocated {
watermark.Allocated = segment.lastSeq
changed = true
}
if segment.sourceCursor != nil && segment.sourceCursor.EndOffset > watermark.SourceOffset {
watermark.SourceOffset = segment.sourceCursor.EndOffset
changed = true
}
}
if changed {
spool.watermarks[streamID] = watermark
if err := spool.persistWatermarksLocked(); err != nil {
return logStreamWatermark{}, false, err
}
}
log.Printf("RUN phase=log_spool.restore status=lazy_complete stream=%s found=%t candidates=%d loaded=%d skipped=%d quarantined=%d durationMs=%d", safeSpoolStreamID(streamID), changed, stats.candidates, stats.loaded, stats.skipped, stats.quarantined, time.Since(startedAt).Milliseconds())
return watermark, changed, nil
}
func (spool LogSpool) Pending() ([]protocol.LogBatchIngestRequest, error) {
spool.mu.Lock()
defer spool.mu.Unlock()
segments, err := spool.pendingSegmentsLocked()
if err != nil {
return nil, err
}
batches := make([]protocol.LogBatchIngestRequest, 0, len(segments))
for _, segment := range segments {
batches = append(batches, segment.batch.LogBatchIngestRequest)
}
return batches, nil
}
func (spool LogSpool) pendingSegmentsLocked() ([]pendingLogSegment, error) {
segments, _, err := spool.pendingSegmentsLockedWithOptions(pendingSegmentScanOptions{})
return segments, err
}
func (spool LogSpool) pendingSegmentsLockedWithOptions(options pendingSegmentScanOptions) ([]pendingLogSegment, pendingSegmentScanStats, error) {
entries, err := os.ReadDir(spool.dir)
if err != nil {
return nil, pendingSegmentScanStats{}, fmt.Errorf("read log spool: %w", err)
}
paths, stats := spool.pendingSegmentPathsLocked(entries, options)
segments := make([]pendingLogSegment, 0, len(paths))
for _, path := range paths {
batch, err := readDurableLogBatch(path)
if err != nil {
if os.IsNotExist(err) {
stats.skipped++
continue
}
if spool.quarantineLogSegmentLocked(path, "unreadable") == nil {
stats.quarantined++
log.Printf("RUN phase=log_spool.scan status=segment_quarantined reason=unreadable file=%s error=%s", safeSpoolFileName(path), logSpoolErrorKind(err))
continue
}
stats.skipped++
log.Printf("RUN phase=log_spool.scan status=segment_skipped reason=unreadable file=%s error=%s", safeSpoolFileName(path), logSpoolErrorKind(err))
continue
}
var modTime time.Time
if info, statErr := os.Stat(path); statErr == nil {
modTime = info.ModTime()
}
segments = append(segments, pendingLogSegment{path: path, modTime: modTime, batch: batch})
stats.loaded++
}
sort.SliceStable(segments, func(i, j int) bool { return pendingLogSegmentBefore(segments[i], segments[j]) })
return segments, stats, nil
}
func (spool LogSpool) pendingSegmentWatermarksLocked(options pendingSegmentScanOptions) ([]pendingLogSegmentWatermark, pendingSegmentScanStats, error) {
entries, err := os.ReadDir(spool.dir)
if err != nil {
return nil, pendingSegmentScanStats{}, fmt.Errorf("read log spool: %w", err)
}
paths, stats := spool.pendingSegmentPathsLocked(entries, options)
segments := make([]pendingLogSegmentWatermark, 0, len(paths))
for _, path := range paths {
watermark, err := readDurableLogBatchWatermark(path)
if err != nil {
if os.IsNotExist(err) {
stats.skipped++
continue
}
if spool.quarantineLogSegmentLocked(path, "unreadable") == nil {
stats.quarantined++
log.Printf("RUN phase=log_spool.scan status=segment_quarantined phase=%s reason=unreadable file=%s error=%s", safeDiagnosticPhase(options.diagnosticPhase), safeSpoolFileName(path), logSpoolErrorKind(err))
continue
}
stats.skipped++
log.Printf("RUN phase=log_spool.scan status=segment_skipped phase=%s reason=unreadable file=%s error=%s", safeDiagnosticPhase(options.diagnosticPhase), safeSpoolFileName(path), logSpoolErrorKind(err))
continue
}
if strings.TrimSpace(watermark.LogStreamID) == "" {
if spool.quarantineLogSegmentLocked(path, "missing_stream") == nil {
stats.quarantined++
log.Printf("RUN phase=log_spool.scan status=segment_quarantined phase=%s reason=missing_stream file=%s", safeDiagnosticPhase(options.diagnosticPhase), safeSpoolFileName(path))
continue
}
stats.skipped++
continue
}
segments = append(segments, pendingLogSegmentWatermark{logStreamID: watermark.LogStreamID, lastSeq: watermark.LastSeq, sourceCursor: watermark.SourceCursor})
stats.loaded++
}
return segments, stats, nil
}
func (spool LogSpool) pendingSegmentPathsLocked(entries []os.DirEntry, options pendingSegmentScanOptions) ([]string, pendingSegmentScanStats) {
startedAt := time.Now()
paths := make([]string, 0, len(entries))
stats := pendingSegmentScanStats{}
if options.diagnosticPhase != "" {
log.Printf("RUN phase=log_spool.scan status=starting phase=%s entries=%d", safeDiagnosticPhase(options.diagnosticPhase), len(entries))
}
for _, entry := range entries {
if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".json") {
continue
}
stats.candidates++
path := filepath.Join(spool.dir, entry.Name())
info, err := entry.Info()
if err != nil {
if spool.quarantineLogSegmentLocked(path, "stat_failed") == nil {
stats.quarantined++
log.Printf("RUN phase=log_spool.scan status=segment_quarantined phase=%s reason=stat_failed file=%s error=%s", safeDiagnosticPhase(options.diagnosticPhase), safeSpoolFileName(path), logSpoolErrorKind(err))
continue
}
stats.skipped++
log.Printf("RUN phase=log_spool.scan status=segment_skipped phase=%s reason=stat_failed file=%s error=%s", safeDiagnosticPhase(options.diagnosticPhase), safeSpoolFileName(path), logSpoolErrorKind(err))
continue
}
if !info.Mode().IsRegular() {
if spool.quarantineLogSegmentLocked(path, "non_regular") == nil {
stats.quarantined++
log.Printf("RUN phase=log_spool.scan status=segment_quarantined phase=%s reason=non_regular file=%s", safeDiagnosticPhase(options.diagnosticPhase), safeSpoolFileName(path))
continue
}
stats.skipped++
continue
}
if info.Size() > maxDurableLogSegmentBytes {
if spool.quarantineLogSegmentLocked(path, "oversized") == nil {
stats.quarantined++
log.Printf("RUN phase=log_spool.scan status=segment_quarantined phase=%s reason=oversized file=%s bytes=%d limit=%d", safeDiagnosticPhase(options.diagnosticPhase), safeSpoolFileName(path), info.Size(), maxDurableLogSegmentBytes)
continue
}
stats.skipped++
continue
}
stats.bytes += info.Size()
paths = append(paths, path)
if options.diagnosticPhase != "" && stats.candidates%100 == 0 {
log.Printf("RUN phase=log_spool.scan status=progress phase=%s candidates=%d accepted=%d quarantined=%d skipped=%d durationMs=%d", safeDiagnosticPhase(options.diagnosticPhase), stats.candidates, len(paths), stats.quarantined, stats.skipped, time.Since(startedAt).Milliseconds())
}
}
sort.Strings(paths)
if options.diagnosticPhase != "" {
log.Printf("RUN phase=log_spool.scan status=ready phase=%s candidates=%d accepted=%d quarantined=%d skipped=%d bytes=%d durationMs=%d", safeDiagnosticPhase(options.diagnosticPhase), stats.candidates, len(paths), stats.quarantined, stats.skipped, stats.bytes, time.Since(startedAt).Milliseconds())
}
return paths, stats
}
func pendingLogSegmentBefore(left pendingLogSegment, right pendingLogSegment) bool {
leftBatch := left.batch.LogBatchIngestRequest
rightBatch := right.batch.LogBatchIngestRequest
if leftBatch.LogStreamID == rightBatch.LogStreamID {
if leftBatch.FirstSeq != rightBatch.FirstSeq {
return leftBatch.FirstSeq < rightBatch.FirstSeq
}
return leftBatch.LastSeq < rightBatch.LastSeq
}
leftTime := pendingLogSegmentPriorityTime(left)
rightTime := pendingLogSegmentPriorityTime(right)
if !leftTime.Equal(rightTime) {
return leftTime.After(rightTime)
}
return left.path < right.path
}
func pendingLogSegmentPriorityTime(segment pendingLogSegment) time.Time {
batch := segment.batch.LogBatchIngestRequest
if batch.Source == "process" && strings.TrimSpace(batch.LogSessionID) != "" && !batch.SessionStartedAt.IsZero() {
return batch.SessionStartedAt
}
return segment.modTime
}
func (spool LogSpool) Ack(response protocol.LogBatchIngestResponse) error {
if response.AcceptedFrom == 0 || response.AcceptedTo < response.AcceptedFrom {
return nil
}
spool.mu.Lock()
defer spool.mu.Unlock()
watermark := spool.watermarks[response.LogStreamID]
if response.AcceptedTo > watermark.Acknowledged {
watermark.Acknowledged = response.AcceptedTo
}
if watermark.Allocated < watermark.Acknowledged {
watermark.Allocated = watermark.Acknowledged
}
spool.watermarks[response.LogStreamID] = watermark
if err := spool.persistWatermarksLocked(); err != nil {
return err
}
segments, err := spool.pendingSegmentsLocked()
if err != nil {
return err
}
for _, segment := range segments {
batch := segment.batch
if batch.LogStreamID == response.LogStreamID && batch.FirstSeq >= response.AcceptedFrom && batch.LastSeq <= response.AcceptedTo {
if err := os.Remove(segment.path); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove acknowledged log spool segment: %w", err)
}
}
}
return nil
}
func (spool *LogSpool) loadWatermarks() error {
body, err := os.ReadFile(spool.statePath)
if os.IsNotExist(err) {
return nil
}
if err != nil {
return fmt.Errorf("read log watermarks: %w", err)
}
var state logSpoolState
if err := json.Unmarshal(body, &state); err != nil {
return fmt.Errorf("decode log watermarks: %w", err)
}
if state.Streams != nil {
spool.watermarks = state.Streams
}
return nil
}
func (spool LogSpool) persistWatermarksLocked() error {
body, err := json.Marshal(logSpoolState{Streams: spool.watermarks})
if err != nil {
return err
}
temporary := spool.statePath + ".tmp"
if err := os.WriteFile(temporary, body, 0o600); err != nil {
return fmt.Errorf("write log watermarks: %w", err)
}
if err := syncFile(temporary); err != nil {
_ = os.Remove(temporary)
return err
}
if err := os.Rename(temporary, spool.statePath); err != nil {
_ = os.Remove(temporary)
return fmt.Errorf("commit log watermarks: %w", err)
}
return nil
}
type LogBatchClient interface {
IngestLogBatch(context.Context, protocol.LogBatchIngestRequest) (protocol.LogBatchIngestResponse, error)
}
type LogStreamProgressClient interface {
LogStreamLatestSeq(context.Context, protocol.LogBatchIngestRequest) (uint64, error)
}
type PermanentLogBatchError struct {
Reason string
Err error
}
func (err PermanentLogBatchError) Error() string {
if err.Err == nil {
return "permanent log batch rejection: " + err.Reason
}
return "permanent log batch rejection: " + err.Reason + ": " + err.Err.Error()
}
func (err PermanentLogBatchError) Unwrap() error {
return err.Err
}
func PermanentLogBatchRejection(reason string, err error) error {
return PermanentLogBatchError{Reason: sanitizeSegmentName(reason), Err: err}
}
func (spool LogSpool) Flush(ctx context.Context, client LogBatchClient) (int, error) {
if client == nil {
return 0, fmt.Errorf("log batch client is required")
}
spool.mu.Lock()
pending, err := spool.pendingSegmentsLocked()
if err != nil {
spool.mu.Unlock()
return 0, err
}
for _, segment := range pending {
spool.inflight[segment.path] = true
}
spool.mu.Unlock()
defer func() {
spool.mu.Lock()
defer spool.mu.Unlock()
for _, segment := range pending {
delete(spool.inflight, segment.path)
}
}()
acknowledged := 0
for _, segment := range pending {
batch := segment.batch.LogBatchIngestRequest
if err := ctx.Err(); err != nil {
return acknowledged, err
}
response, err := client.IngestLogBatch(ctx, batch)
if err != nil {
var permanent PermanentLogBatchError
if errors.As(err, &permanent) {
if permanent.Reason == "platform_sequence_gap" {
latestSeq, progressErr := recoverLogStreamLatestSeq(ctx, client, batch)
if progressErr != nil {
return acknowledged, err
}
rejected, rejectErr := spool.RejectStreamAfter(batch.LogStreamID, latestSeq, permanent.Reason)
if rejectErr != nil {
return acknowledged, rejectErr
}
if resetErr := spool.ResetStreamWatermark(batch.LogStreamID, latestSeq); resetErr != nil {
return acknowledged, resetErr
}
acknowledged += rejected
return acknowledged, nil
}
if rejectErr := spool.Reject(batch, permanent.Reason); rejectErr != nil {
return acknowledged, rejectErr
}
acknowledged++
continue
}
return acknowledged, err
}
if !response.Accepted || response.LogStreamID != batch.LogStreamID || response.AcceptedFrom > batch.FirstSeq || response.AcceptedTo < batch.LastSeq {
return acknowledged, fmt.Errorf("platform log acknowledgement does not cover pending batch")
}
if err := spool.Ack(response); err != nil {
return acknowledged, err
}
acknowledged++
}
return acknowledged, nil
}
func recoverLogStreamLatestSeq(ctx context.Context, client LogBatchClient, batch protocol.LogBatchIngestRequest) (uint64, error) {
progressClient, ok := client.(LogStreamProgressClient)
if !ok {
return 0, fmt.Errorf("log stream progress client is required after sequence gap")
}
return progressClient.LogStreamLatestSeq(ctx, batch)
}
func (spool LogSpool) ResetStreamWatermark(logStreamID string, latestSeq uint64) error {
spool.mu.Lock()
defer spool.mu.Unlock()
watermark := spool.watermarks[logStreamID]
watermark.Allocated = latestSeq
watermark.Acknowledged = latestSeq
spool.watermarks[logStreamID] = watermark
return spool.persistWatermarksLocked()
}
func (spool LogSpool) RejectStreamAfter(logStreamID string, latestSeq uint64, reason string) (int, error) {
spool.mu.Lock()
defer spool.mu.Unlock()
segments, err := spool.pendingSegmentsLocked()
if err != nil {
return 0, err
}
rejected := 0
for _, segment := range segments {
batch := segment.batch.LogBatchIngestRequest
if batch.LogStreamID != logStreamID || batch.FirstSeq <= latestSeq {
continue
}
if err := spool.moveLogSegmentToRejectedLocked(segment.path, reason); err != nil {
return rejected, err
}
rejected++
}
return rejected, nil
}
func (spool LogSpool) Reject(batch protocol.LogBatchIngestRequest, reason string) error {
spool.mu.Lock()
defer spool.mu.Unlock()
segments, err := spool.pendingSegmentsLocked()
if err != nil {
return err
}
path := ""
for _, segment := range segments {
current := segment.batch
if current.LogStreamID == batch.LogStreamID && current.FirstSeq == batch.FirstSeq && current.LastSeq == batch.LastSeq && current.Checksum == batch.Checksum {
path = segment.path
break
}
}
if path == "" {
return nil
}
return spool.moveLogSegmentToRejectedLocked(path, reason)
}
func readDurableLogBatch(path string) (durableLogBatch, error) {
if err := validateDurableLogSegment(path); err != nil {
return durableLogBatch{}, err
}
file, err := os.Open(path)
if err != nil {
return durableLogBatch{}, err
}
defer file.Close()
var batch durableLogBatch
if err := json.NewDecoder(file).Decode(&batch); err != nil {
return durableLogBatch{}, fmt.Errorf("decode log spool segment: %w", err)
}
return batch, nil
}
func readDurableLogBatchWatermark(path string) (struct {
LogStreamID string `json:"logStreamId"`
LastSeq uint64 `json:"lastSeq"`
SourceCursor *LogSourceCursor `json:"_sourceCursor,omitempty"`
}, error) {
var watermark struct {
LogStreamID string `json:"logStreamId"`
LastSeq uint64 `json:"lastSeq"`
SourceCursor *LogSourceCursor `json:"_sourceCursor,omitempty"`
}
if err := validateDurableLogSegment(path); err != nil {
return watermark, err
}
file, err := os.Open(path)
if err != nil {
return watermark, err
}
defer file.Close()
if err := json.NewDecoder(file).Decode(&watermark); err != nil {
return watermark, fmt.Errorf("decode log spool segment watermark: %w", err)
}
return watermark, nil
}
func validateDurableLogSegment(path string) error {
info, err := os.Stat(path)
if err != nil {
return err
}
if !info.Mode().IsRegular() {
return fmt.Errorf("log spool segment is not a regular file")
}
if info.Size() > maxDurableLogSegmentBytes {
return fmt.Errorf("log spool segment exceeds %d bytes", maxDurableLogSegmentBytes)
}
return nil
}
func (spool LogSpool) quarantineLogSegmentLocked(path string, reason string) error {
return spool.moveLogSegmentToRejectedLocked(path, reason)
}
func (spool LogSpool) moveLogSegmentToRejectedLocked(path string, reason string) error {
rejectedDir := filepath.Join(filepath.Dir(spool.dir), "logs-rejected")
if err := os.MkdirAll(rejectedDir, 0o755); err != nil {
return fmt.Errorf("create rejected log spool directory: %w", err)
}
rejectedPath := filepath.Join(rejectedDir, fmt.Sprintf("%s.%s.%d", filepath.Base(path), sanitizeSegmentName(reason), time.Now().UTC().UnixNano()))
if err := os.Rename(path, rejectedPath); err != nil {
return fmt.Errorf("move rejected log spool segment: %w", err)
}
return nil
}
func safeSpoolFileName(path string) string {
return sanitizeSegmentName(filepath.Base(path))
}
func safeSpoolStreamID(streamID string) string {
return sanitizeSegmentName(streamID)
}
func safeDiagnosticPhase(value string) string {
if strings.TrimSpace(value) == "" {
return "-"
}
return sanitizeSegmentName(value)
}
func logSpoolErrorKind(err error) string {
if err == nil {
return "-"
}
if os.IsNotExist(err) {
return "not_found"
}
if os.IsPermission(err) {
return "permission"
}
var pathErr *os.PathError
if errors.As(err, &pathErr) {
return sanitizeSegmentName(pathErr.Op + "_" + pathErr.Err.Error())
}
return sanitizeSegmentName(err.Error())
}
func syncFile(path string) error {
file, err := os.OpenFile(path, os.O_RDWR, 0)
if err != nil {
return fmt.Errorf("open spool file for sync: %w", err)
}
defer file.Close()
if err := file.Sync(); err != nil {
return fmt.Errorf("sync spool file: %w", err)
}
return nil
}
func (spool LogSpool) batchPath(batch durableLogBatch) string {
streamID := sanitizeSegmentName(batch.LogStreamID)
return filepath.Join(spool.dir, fmt.Sprintf("%s-%020d-%020d.json", streamID, batch.FirstSeq, batch.LastSeq))
}
func (spool LogSpool) watermarkPath(streamID string) string {
return filepath.Join(filepath.Dir(spool.dir), "log-watermarks", sanitizeSegmentName(streamID)+".json")
}
func (spool LogSpool) readWatermark(streamID string) (uint64, error) {
file, err := os.Open(spool.watermarkPath(streamID))
if err != nil {
return 0, err
}
defer file.Close()
var watermark struct {
LatestSeq uint64 `json:"latestSeq"`
}
if err := json.NewDecoder(file).Decode(&watermark); err != nil {
return 0, fmt.Errorf("decode log watermark: %w", err)
}
return watermark.LatestSeq, nil
}
func sanitizeSegmentName(value string) string {
var builder strings.Builder
for _, r := range value {
if r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' || r == '-' || r == '_' || r == '.' {
builder.WriteRune(r)
continue
}
builder.WriteByte('_')
}
if builder.Len() == 0 {
return "stream"
}
return builder.String()
}