syncova-backup/packages/backupengine/pipeline.go
Jerrit Fritzsche 610719c316
Some checks failed
CI / Backend (Go) (push) Failing after 3m7s
CI / Frontend (React/TypeScript) (push) Successful in 37s
CI / Sicherheitsprüfungen (push) Successful in 44s
Syncova Backups V1
Enterprise-Backup-, Recovery-, Verification-, Security- und
Monitoring-Plattform fuer Proxmox VE, Windows, Linux und Dateisysteme.

Der Leitsatz, der fast jede Entscheidung erklaert: Ein Backup gilt erst als
vertrauenswuerdig, wenn Integritaet geprueft und Wiederherstellbarkeit
nachgewiesen wurde. Deshalb steigt ein Wiederherstellungspunkt erst nach einem
tatsaechlich durchgefuehrten Restore-Test auf "recoverable", und Unbekanntes
geht in keine Bewertung als "gut" ein.

Umfang (Phasen 0-23):

- Repository Engine: inhaltsadressierte Bloecke, atomares Commit-Protokoll,
  Katalogaufbau allein aus den Manifesten — ohne Datenbank
- Backup Engine: inhaltsabhaengiges Chunking, Deduplizierung trotz
  Verschluesselung, zstd, AES-256-GCM, Streaming mit Gegendruck
- Agenten fuer Windows und Linux mit Auftragsabholung (Pull-Modell)
- Proxmox-Provider mit beiden Zugriffswegen auf die Sicherungsarchive
- Scheduler, Recovery Engine mit Pruefpunkt, Verification, Unveraenderlichkeit
- Weboberflaeche, Kennzahlen, Meldungen, Berichte, Security Center,
  Ransomware-Heuristik (meldet, handelt nie)
- Disaster Recovery, Haertung, Leistungsmessung, Chaos Testing
- Eingefrorene Vertraege fuer API, Migrationen, Backup-Format und Repository
- Auslieferungspaket fuer linux/amd64, linux/arm64 und windows/amd64

Nicht enthalten und als solches gekennzeichnet: Kapazitaetsprognose, Backup
Copy, Changed Block Tracking bei Proxmox, erweiterte Attribute und ACLs.

Gebaut, aber nie auf echter Hardware gefahren: der Windows-Dienst, die
systemd-Einheit und der verpflichtende Proxmox-Meilenstein — ob eine
wiederhergestellte VM startet, ist ungeprueft. Einzelheiten in CHANGELOG.md
und docs/release-candidate.md.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-17 09:10:54 +02:00

366 lines
13 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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
}