package backupengine import ( "context" "crypto/sha256" "encoding/hex" "errors" "fmt" "io" "runtime" "sort" "sync" "github.com/syncova/syncova/packages/repository" ) // defaultQueueDepth ist die Tiefe der Warteschlangen zwischen den Stufen. // // Sie ist der Backpressure-Mechanismus: ist die Schlange voll, blockiert der // Leser, bis ein Arbeiter Kapazität hat. Ohne diese Grenze läse der Chunker // so schnell, wie die Quelle liefert, und füllte den Arbeitsspeicher mit // unverarbeiteten Blöcken (PROMPT.md §80). // // Der Speicherbedarf ist damit nach oben begrenzt: // Tiefe × Höchstblockgröße × (Eingangs- plus Ausgangsschlange). const defaultQueueDepth = 4 // pipelineJob ist ein an einen Arbeiter übergebener Block. type pipelineJob struct { // sequence ist die Position im Datenstrom und stellt die Reihenfolge wieder her. sequence int64 // offset ist die Position des Blocks in den Ursprungsdaten. offset int64 // data ist der Blockinhalt; er gehört ausschließlich diesem Auftrag. data []byte } // pipelineResult ist das Ergebnis eines Arbeiters. type pipelineResult struct { // sequence ist die Position im Datenstrom. sequence int64 // reference beschreibt den abgelegten Block. reference repository.ChunkReference // wasNew meldet, ob der Block neu geschrieben wurde. wasNew bool // processingError ist ein aufgetretener Fehler. processingError error } // chunkSink nimmt die verarbeiteten Blöcke auf. // // Die Schnittstelle entkoppelt die Pipeline vom Repository und macht sie // einzeln prüfbar. type chunkSink interface { // storeChunk legt einen Block ab und liefert dessen Beschreibung. storeChunk(storeContext context.Context, plaintextIdentifier string, storedData []byte, plaintextLength int64) (storedDigest string, wasNew bool, storeError error) // hasChunk meldet, ob ein Block bereits vorliegt. hasChunk(queryContext context.Context, plaintextIdentifier string) (bool, error) // storedDigestOf liefert die Prüfsumme der abgelegten Form eines vorhandenen Blocks. storedDigestOf(queryContext context.Context, plaintextIdentifier string) (string, error) // noteDeduplicated vermerkt einen wiederverwendeten Block in der Statistik. noteDeduplicated(plaintextLength int64) } // repositorySink schreibt in ein Repository. // // Geschrieben wird ueber die **Schreibsession**, gelesen unmittelbar am // Repository. Der Unterschied ist keine Formalie: Nur die Session zaehlt mit, // und ihre Zahlen wandern ins Manifest. Ein Schreibweg am Repository vorbei // hinterliesse ein Manifest mit einer Statistik von null — und damit ein // Repository, das seine eigene Groesse nicht kennt. type repositorySink struct { // backupWriter ist die Schreibsession des laufenden Backups. backupWriter repository.Writer // targetRepository dient den reinen Leseabfragen. targetRepository *repository.LocalRepository } // storeChunk legt einen Block ueber die Schreibsession ab. func (sink *repositorySink) storeChunk(storeContext context.Context, plaintextIdentifier string, storedData []byte, plaintextLength int64) (string, bool, error) { return sink.backupWriter.WriteTransformedChunk(storeContext, plaintextIdentifier, storedData, plaintextLength) } // hasChunk meldet, ob ein Block bereits im Repository liegt. func (sink *repositorySink) hasChunk(queryContext context.Context, plaintextIdentifier string) (bool, error) { return sink.targetRepository.HasChunk(queryContext, plaintextIdentifier) } // storedDigestOf liefert die Prüfsumme der abgelegten Form eines vorhandenen Blocks. func (sink *repositorySink) storedDigestOf(queryContext context.Context, plaintextIdentifier string) (string, error) { return sink.targetRepository.StoredDigestOfChunk(queryContext, plaintextIdentifier) } // noteDeduplicated vermerkt einen wiederverwendeten Block in der Statistik. func (sink *repositorySink) noteDeduplicated(plaintextLength int64) { sink.backupWriter.NoteDeduplicatedChunk(plaintextLength) } // pipeline verarbeitet einen Datenstrom nebenläufig zu abgelegten Blöcken. type pipeline struct { // transformer komprimiert und verschlüsselt die Blöcke. transformer *ChunkTransformer // sink nimmt die verarbeiteten Blöcke auf. sink chunkSink // workerCount ist die Zahl paralleler Arbeiter. workerCount int // queueDepth ist die Tiefe der Warteschlangen. queueDepth int // progressReporter meldet den Fortschritt. progressReporter *ProgressReporter } // process verarbeitet einen Datenstrom vollständig. // // Der Ablauf ist bewusst dreigeteilt: // // - Ein Leser zerlegt den Strom in Blöcke. Er läuft für sich, weil die // Blockgrenzen von den Vorgängerdaten abhängen und sich nicht parallelisieren // lassen. // - Mehrere Arbeiter hashen, deduplizieren, komprimieren, verschlüsseln und // schreiben. Hier liegt die eigentliche Rechenarbeit. // - Ein Sammler ordnet die Ergebnisse wieder nach Position. Die Reihenfolge // im Manifest muss der Reihenfolge der Ursprungsdaten entsprechen, sonst // ergäbe eine Wiederherstellung Unsinn. func (processingPipeline *pipeline) process(processContext context.Context, sourceReader io.Reader, chunkerOptions ChunkerOptions) ([]repository.ChunkReference, error) { // Ein eigener Context erlaubt es, bei einem Fehler alle Beteiligten zu stoppen. pipelineContext, cancelPipeline := context.WithCancel(processContext) defer cancelPipeline() jobQueue := make(chan pipelineJob, processingPipeline.queueDepth) resultQueue := make(chan pipelineResult, processingPipeline.queueDepth) // readError wird vom Leser gesetzt und nach dem Warten ausgewertet. var readError error var readerWaitGroup sync.WaitGroup readerWaitGroup.Add(1) go func() { defer readerWaitGroup.Done() defer close(jobQueue) readError = processingPipeline.runReader(pipelineContext, sourceReader, chunkerOptions, jobQueue) }() var workerWaitGroup sync.WaitGroup for workerIndex := 0; workerIndex < processingPipeline.workerCount; workerIndex++ { workerWaitGroup.Add(1) go func() { defer workerWaitGroup.Done() processingPipeline.runWorker(pipelineContext, jobQueue, resultQueue) }() } // Die Ergebnisschlange wird geschlossen, sobald alle Arbeiter fertig sind. go func() { workerWaitGroup.Wait() close(resultQueue) }() collectedResults := make([]pipelineResult, 0, 64) var firstProcessingError error for pipelineOutcome := range resultQueue { if pipelineOutcome.processingError != nil { // Der erste Fehler stoppt die Pipeline; die übrigen Ergebnisse // werden noch abgeräumt, damit keine Goroutine hängen bleibt. if firstProcessingError == nil { firstProcessingError = pipelineOutcome.processingError cancelPipeline() } continue } collectedResults = append(collectedResults, pipelineOutcome) } readerWaitGroup.Wait() if firstProcessingError != nil { return nil, firstProcessingError } if readError != nil { return nil, readError } if contextError := processContext.Err(); contextError != nil { return nil, contextError } // Die Ergebnisse werden in die Reihenfolge der Ursprungsdaten gebracht. sort.Slice(collectedResults, func(firstIndex int, secondIndex int) bool { return collectedResults[firstIndex].sequence < collectedResults[secondIndex].sequence }) chunkReferences := make([]repository.ChunkReference, 0, len(collectedResults)) for _, pipelineOutcome := range collectedResults { chunkReferences = append(chunkReferences, pipelineOutcome.reference) } return chunkReferences, nil } // runReader zerlegt den Datenstrom und übergibt die Blöcke an die Arbeiter. func (processingPipeline *pipeline) runReader(readContext context.Context, sourceReader io.Reader, chunkerOptions ChunkerOptions, jobQueue chan<- pipelineJob) error { streamChunker := NewChunker(sourceReader, chunkerOptions) for { if contextError := readContext.Err(); contextError != nil { return contextError } nextChunk, chunkError := streamChunker.Next() if errors.Is(chunkError, io.EOF) { return nil } if chunkError != nil { return chunkError } // Der Chunker verwendet seinen Puffer wieder; der Auftrag braucht eine // eigene Kopie, sonst überschriebe der nächste Lesevorgang die Daten, // während ein Arbeiter sie noch verarbeitet. chunkCopy := make([]byte, len(nextChunk.Data)) copy(chunkCopy, nextChunk.Data) // Ist die Schlange voll, blockiert diese Zeile — das ist der // Backpressure-Mechanismus. select { case jobQueue <- pipelineJob{sequence: nextChunk.Sequence, offset: nextChunk.Offset, data: chunkCopy}: case <-readContext.Done(): return readContext.Err() } } } // runWorker verarbeitet Blöcke, bis die Auftragsschlange erschöpft ist. func (processingPipeline *pipeline) runWorker(workerContext context.Context, jobQueue <-chan pipelineJob, resultQueue chan<- pipelineResult) { for pipelineTask := range jobQueue { if contextError := workerContext.Err(); contextError != nil { return } workerOutcome := processingPipeline.processSingleChunk(workerContext, pipelineTask) select { case resultQueue <- workerOutcome: case <-workerContext.Done(): return } } } // processSingleChunk führt Hash, Deduplizierung, Transformation und Ablage aus. func (processingPipeline *pipeline) processSingleChunk(chunkContext context.Context, pipelineTask pipelineJob) pipelineResult { // Der Hash entsteht über dem Klartext. Nur so finden zwei gleiche // Ursprungsblöcke zusammen — verschlüsselte Fassungen wären stets verschieden. plaintextDigest := sha256.Sum256(pipelineTask.data) plaintextIdentifier := hex.EncodeToString(plaintextDigest[:]) chunkReference := repository.ChunkReference{ Identifier: plaintextIdentifier, LogicalOffset: pipelineTask.offset, LogicalLength: int64(len(pipelineTask.data)), } // Liegt der Block bereits vor, entfallen Kompression, Verschlüsselung und // Schreibvorgang vollständig. Genau hier entsteht der Zeitgewinn der // Deduplizierung, nicht erst beim Speicherplatz (PROMPT.md §10). alreadyPresent, presenceError := processingPipeline.sink.hasChunk(chunkContext, plaintextIdentifier) if presenceError != nil { return pipelineResult{sequence: pipelineTask.sequence, processingError: presenceError} } if alreadyPresent { processingPipeline.progressReporter.recordDeduplicatedChunk(int64(len(pipelineTask.data))) // Die Prüfsumme der bereits abgelegten Form gehört ins Manifest, damit // ein Integritätslauf auch diesen Block prüfen kann. storedDigest, digestError := processingPipeline.sink.storedDigestOf(chunkContext, plaintextIdentifier) if digestError != nil { return pipelineResult{sequence: pipelineTask.sequence, processingError: digestError} } chunkReference.StoredDigest = storedDigest // Der Block wurde gelesen und erkannt, nur nicht abgelegt. Ohne diesen // Vermerk bliebe die Statistik eines Laufs ueber unveraenderte Daten // bei null. processingPipeline.sink.noteDeduplicated(int64(len(pipelineTask.data))) return pipelineResult{sequence: pipelineTask.sequence, reference: chunkReference, wasNew: false} } transformedChunk, wasCompressed, transformError := processingPipeline.transformer.Transform(pipelineTask.data) if transformError != nil { return pipelineResult{sequence: pipelineTask.sequence, processingError: transformError} } // Die Klartextlaenge geht mit: Aus der umgewandelten Form laesst sie sich // nicht zurueckrechnen, und ohne sie wuesste die Session nicht, wie viel // Ursprungsdaten sie verarbeitet hat. storedDigest, wasNew, storeError := processingPipeline.sink.storeChunk(chunkContext, plaintextIdentifier, transformedChunk, int64(len(pipelineTask.data))) if storeError != nil { return pipelineResult{sequence: pipelineTask.sequence, processingError: storeError} } chunkReference.StoredLength = int64(len(transformedChunk)) chunkReference.StoredDigest = storedDigest if wasNew { processingPipeline.progressReporter.recordWrittenChunk(int64(len(pipelineTask.data)), int64(len(transformedChunk))) // Der Entropie-Indikator zaehlt nur **neue** Bloecke: Ein // deduplizierter stammt aus einem frueheren Lauf und sagt nichts ueber // die Daten von heute. if !wasCompressed { processingPipeline.progressReporter.recordIncompressibleChunk() } } else { // Ein anderer Arbeiter war schneller: derselbe Block wurde parallel abgelegt. processingPipeline.progressReporter.recordDeduplicatedChunk(int64(len(pipelineTask.data))) } return pipelineResult{sequence: pipelineTask.sequence, reference: chunkReference, wasNew: wasNew} } // defaultWorkerCount bestimmt eine sinnvolle Zahl paralleler Arbeiter. // // Kompression und Verschlüsselung sind rechenintensiv; mehr Arbeiter als Kerne // brächten keinen Gewinn, sondern nur zusätzlichen Speicherbedarf. func defaultWorkerCount() int { availableCores := runtime.NumCPU() // Ein einzelner Kern soll nicht zu null Arbeitern führen. if availableCores < 1 { return 1 } // Eine Obergrenze begrenzt den Speicherbedarf auf großen Maschinen: // jeder Arbeiter hält bis zu einen Block im Speicher. const maximumWorkers = 8 return min(availableCores, maximumWorkers) } // validateWorkerCount prüft eine vorgegebene Arbeiterzahl. func validateWorkerCount(requestedWorkers int) (int, error) { if requestedWorkers == 0 { return defaultWorkerCount(), nil } if requestedWorkers < 0 { return 0, fmt.Errorf("die zahl der arbeiter muss positiv sein, war %d", requestedWorkers) } return requestedWorkers, nil }