package checkpoint import ( "encoding/base64" "encoding/json" "fmt" "os" "strings" "gitea.maximumdirect.net/eric/notarius/internal/core/source" coreworkspace "gitea.maximumdirect.net/eric/notarius/internal/core/workspace" "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" "gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline" ) type WorkspaceLoader struct { root string identityDigest string } func NewWorkspaceLoader(settings coreworkspace.Settings, identity coreworkspace.CheckpointIdentity) (pipeline.CheckpointLoader, error) { root, err := settings.CheckpointDirectory(identity) if err != nil { return nil, err } if strings.TrimSpace(root) == "" { return pipeline.NoopCheckpointLoader(), nil } return &WorkspaceLoader{root: root, identityDigest: identity.Digest}, nil } func (l *WorkspaceLoader) Enabled() bool { return l != nil && strings.TrimSpace(l.root) != "" } func (l *WorkspaceLoader) Source(moduleKey string) (pipeline.SourceCheckpoint, pipeline.CheckpointDecision) { var manifest coreworkspace.SourceManifest if decision := l.readJSON("source/manifest.json", &manifest); !decision.Reused { return pipeline.SourceCheckpoint{}, decision } if decision := l.validateManifest(manifest.StageManifest, coreworkspace.StageSource, "", moduleKey, coreworkspace.StatusSucceeded, nil); !decision.Reused { return pipeline.SourceCheckpoint{}, decision } var payload sourceDocumentEnvelope if decision := l.readJSON("source/source-document.json", &payload); !decision.Reused { return pipeline.SourceCheckpoint{}, decision } doc := cloneSourceDocument(payload.Document) if err := source.ValidateDocument(&doc); err != nil { return pipeline.SourceCheckpoint{}, invalidDecision("source checkpoint document is invalid: %v", err) } if strings.TrimSpace(manifest.SourceID) != "" && manifest.SourceID != doc.ID { return pipeline.SourceCheckpoint{}, invalidDecision("source checkpoint source id does not match payload") } if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.OutputDigests), digestFingerprints("source_document", doc.Digest)) { return pipeline.SourceCheckpoint{}, invalidDecision("source checkpoint output digest does not match payload") } return pipeline.SourceCheckpoint{Document: &doc}, reusedDecision() } func (l *WorkspaceLoader) Chunk(moduleKey string, sourceDigest string) (pipeline.ChunkCheckpoint, pipeline.CheckpointDecision) { expectedDependencies := digestFingerprints("source_document", sourceDigest) var manifest coreworkspace.ChunkManifest if decision := l.readJSON("chunk/manifest.json", &manifest); !decision.Reused { return pipeline.ChunkCheckpoint{}, decision } if decision := l.validateManifest(manifest.StageManifest, coreworkspace.StageChunk, "", moduleKey, coreworkspace.StatusSucceeded, expectedDependencies); !decision.Reused { return pipeline.ChunkCheckpoint{}, decision } var payload chunksEnvelope if decision := l.readJSON("chunk/chunks.json", &payload); !decision.Reused { return pipeline.ChunkCheckpoint{}, decision } chunks, err := sourceChunksFromEnvelope(payload.Chunks) if err != nil { return pipeline.ChunkCheckpoint{}, invalidDecision("chunk checkpoint payload is invalid: %v", err) } if len(chunks) == 0 { return pipeline.ChunkCheckpoint{}, invalidDecision("chunk checkpoint payload has no chunks") } if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.OutputDigests), chunkOutputDigests(chunks)) { return pipeline.ChunkCheckpoint{}, invalidDecision("chunk checkpoint output digests do not match payload") } return pipeline.ChunkCheckpoint{Chunks: chunks, Warnings: cloneWarnings(payload.Warnings)}, reusedDecision() } func (l *WorkspaceLoader) Extract(laneID string, moduleKey string, dependencies []pipeline.CheckpointFingerprint) (pipeline.ExtractCheckpoint, pipeline.CheckpointDecision) { var manifest coreworkspace.ExtractLaneManifest if decision := l.readJSON(laneManifestPath("extract", laneID), &manifest); !decision.Reused { return pipeline.ExtractCheckpoint{}, decision } if decision := l.validateLaneManifest(manifest.StageManifest, coreworkspace.StageExtract, laneID, moduleKey, dependencies, coreworkspace.StatusSucceeded, coreworkspace.StatusSucceededWithRejections); !decision.Reused { return pipeline.ExtractCheckpoint{}, decision } var payload extractOutputsEnvelope if decision := l.readJSON(lanePayloadPath("extract", laneID, "outputs.json"), &payload); !decision.Reused { return pipeline.ExtractCheckpoint{}, decision } outputs, err := extractOutputsFromEnvelope(payload.Outputs) if err != nil { return pipeline.ExtractCheckpoint{}, invalidDecision("extract checkpoint payload is invalid: %v", err) } if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.OutputDigests), rawOutputDigests(extractPayloads(outputs))) { return pipeline.ExtractCheckpoint{}, invalidDecision("extract checkpoint output digests do not match payload") } return pipeline.ExtractCheckpoint{ Outputs: outputs, Rejected: cloneRejectedOutputs(payload.Rejected), Warnings: cloneWarnings(payload.Warnings), }, reusedDecision() } func (l *WorkspaceLoader) Merge(laneID string, moduleKey string, dependencies []pipeline.CheckpointFingerprint) (pipeline.MergeCheckpoint, pipeline.CheckpointDecision) { var manifest coreworkspace.MergeLaneManifest if decision := l.readJSON(laneManifestPath("merge", laneID), &manifest); !decision.Reused { return pipeline.MergeCheckpoint{}, decision } if decision := l.validateLaneManifest(manifest.StageManifest, coreworkspace.StageMerge, laneID, moduleKey, dependencies, coreworkspace.StatusSucceeded); !decision.Reused { return pipeline.MergeCheckpoint{}, decision } var payload mergeOutputEnvelope if decision := l.readJSON(lanePayloadPath("merge", laneID, "output.json"), &payload); !decision.Reused { return pipeline.MergeCheckpoint{}, decision } output, err := mergeOutputFromEnvelope(payload.Output) if err != nil { return pipeline.MergeCheckpoint{}, invalidDecision("merge checkpoint payload is invalid: %v", err) } if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.OutputDigests), rawOutputDigests([]contracts.RawPayload{output.Payload})) { return pipeline.MergeCheckpoint{}, invalidDecision("merge checkpoint output digest does not match payload") } return pipeline.MergeCheckpoint{Output: output, Warnings: cloneWarnings(payload.Warnings)}, reusedDecision() } func (l *WorkspaceLoader) Normalize(laneID string, moduleKey string, dependencies []pipeline.CheckpointFingerprint) (pipeline.NormalizeCheckpoint, pipeline.CheckpointDecision) { var manifest coreworkspace.NormalizeLaneManifest if decision := l.readJSON(laneManifestPath("normalize", laneID), &manifest); !decision.Reused { return pipeline.NormalizeCheckpoint{}, decision } if decision := l.validateLaneManifest(manifest.StageManifest, coreworkspace.StageNormalize, laneID, moduleKey, dependencies, coreworkspace.StatusSucceeded); !decision.Reused { return pipeline.NormalizeCheckpoint{}, decision } var payload normalizeOutputEnvelope if decision := l.readJSON(lanePayloadPath("normalize", laneID, "output.json"), &payload); !decision.Reused { return pipeline.NormalizeCheckpoint{}, decision } output, err := normalizeOutputFromEnvelope(payload.Output) if err != nil { return pipeline.NormalizeCheckpoint{}, invalidDecision("normalize checkpoint payload is invalid: %v", err) } if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.OutputDigests), rawOutputDigests([]contracts.RawPayload{output.Payload})) { return pipeline.NormalizeCheckpoint{}, invalidDecision("normalize checkpoint output digest does not match payload") } return pipeline.NormalizeCheckpoint{Output: output, Warnings: cloneWarnings(payload.Warnings)}, reusedDecision() } func (l *WorkspaceLoader) readJSON(name string, out any) pipeline.CheckpointDecision { if !l.Enabled() { return pipeline.CheckpointDecision{Reason: "checkpoint loading disabled"} } target, err := coreworkspace.SafePath(l.root, name) if err != nil { return invalidDecision("checkpoint path is invalid: %v", err) } data, err := os.ReadFile(target) if err != nil { if os.IsNotExist(err) { return pipeline.CheckpointDecision{Reason: "checkpoint artifact is missing"} } return invalidDecision("read checkpoint artifact: %v", err) } if err := json.Unmarshal(data, out); err != nil { return invalidDecision("decode checkpoint artifact: %v", err) } return reusedDecision() } func (l *WorkspaceLoader) validateManifest(manifest coreworkspace.StageManifest, stage coreworkspace.StageName, laneID string, moduleKey string, status coreworkspace.StageStatus, dependencies []pipeline.CheckpointFingerprint) pipeline.CheckpointDecision { return l.validateLaneManifest(manifest, stage, laneID, moduleKey, dependencies, status) } func (l *WorkspaceLoader) validateLaneManifest(manifest coreworkspace.StageManifest, stage coreworkspace.StageName, laneID string, moduleKey string, dependencies []pipeline.CheckpointFingerprint, statuses ...coreworkspace.StageStatus) pipeline.CheckpointDecision { if manifest.WorkspaceSchemaVersion != coreworkspace.WorkspaceSchemaVersion { return invalidDecision("checkpoint workspace schema version %q is not supported", manifest.WorkspaceSchemaVersion) } if strings.TrimSpace(l.identityDigest) != "" && manifest.Metadata["checkpoint_identity_digest"] != l.identityDigest { return invalidDecision("checkpoint identity digest does not match current invocation") } if manifest.Stage != stage { return invalidDecision("checkpoint stage %q does not match %q", manifest.Stage, stage) } if strings.TrimSpace(laneID) != "" && manifest.LaneID != laneID { return invalidDecision("checkpoint lane %q does not match %q", manifest.LaneID, laneID) } if strings.TrimSpace(moduleKey) != "" && manifest.ModuleKey != moduleKey { return invalidDecision("checkpoint module %q does not match %q", manifest.ModuleKey, moduleKey) } statusOK := false for _, status := range statuses { if manifest.Status == status { statusOK = true break } } if !statusOK { return invalidDecision("checkpoint status %q cannot be reused", manifest.Status) } if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.DependencyFingerprints), dependencies) { return invalidDecision("checkpoint dependency fingerprints do not match") } return reusedDecision() } func sourceChunksFromEnvelope(values []chunkEnvelope) ([]contracts.SourceChunk, error) { if len(values) == 0 { return nil, nil } out := make([]contracts.SourceChunk, 0, len(values)) for _, value := range values { content, err := contentFromEnvelope(value.Content) if err != nil { return nil, err } out = append(out, contracts.SourceChunk{ ID: value.ID, SourceID: value.SourceID, Index: value.Index, StartUnitID: value.StartUnitID, EndUnitID: value.EndUnitID, Content: content, MediaType: value.Content.MediaType, Units: cloneSourceUnits(value.Units), Metadata: cloneMetadata(value.Metadata), }) } return out, nil } func extractOutputsFromEnvelope(values []extractOutputEnvelope) ([]contracts.ExtractOutput, error) { if len(values) == 0 { return nil, nil } out := make([]contracts.ExtractOutput, 0, len(values)) for _, value := range values { payload, err := rawPayloadFromEnvelope(value.Payload) if err != nil { return nil, err } out = append(out, contracts.ExtractOutput{ LaneID: value.LaneID, ExtractorKey: value.ExtractorKey, SourceID: value.SourceID, ChunkID: value.ChunkID, ChunkIndex: value.ChunkIndex, Schema: value.Schema, Payload: payload, }) } return out, nil } func mergeOutputFromEnvelope(value mergeOutputPayload) (contracts.MergeOutput, error) { payload, err := rawPayloadFromEnvelope(value.Payload) if err != nil { return contracts.MergeOutput{}, err } return contracts.MergeOutput{ LaneID: value.LaneID, MergerKey: value.MergerKey, SourceID: value.SourceID, Schema: value.Schema, Payload: payload, }, nil } func normalizeOutputFromEnvelope(value normalizeOutputPayload) (contracts.NormalizeOutput, error) { payload, err := rawPayloadFromEnvelope(value.Payload) if err != nil { return contracts.NormalizeOutput{}, err } return contracts.NormalizeOutput{ LaneID: value.LaneID, NormalizerKey: value.NormalizerKey, SourceID: value.SourceID, Schema: value.Schema, Payload: payload, }, nil } func rawPayloadFromEnvelope(value binaryEnvelope) (contracts.RawPayload, error) { content, err := contentFromEnvelope(value) if err != nil { return contracts.RawPayload{}, err } return contracts.RawPayload{ Content: content, MediaType: value.MediaType, Metadata: cloneMetadata(value.Metadata), Warnings: cloneWarnings(value.Warnings), }, nil } func contentFromEnvelope(value binaryEnvelope) ([]byte, error) { content, err := base64.StdEncoding.DecodeString(value.ContentBase64) if err != nil { return nil, fmt.Errorf("decode content_base64: %w", err) } if digest := strings.TrimSpace(value.ContentDigest); digest != "" && digest != contentDigest(content) { return nil, fmt.Errorf("content digest mismatch") } return content, nil } func coreworkspaceToPipelineFingerprints(values []coreworkspace.Fingerprint) []pipeline.CheckpointFingerprint { if len(values) == 0 { return nil } out := make([]pipeline.CheckpointFingerprint, 0, len(values)) for _, value := range values { out = append(out, pipeline.CheckpointFingerprint{Name: value.Name, Value: value.Value}) } return normalizeFingerprints(out) } func fingerprintsEqual(a []pipeline.CheckpointFingerprint, b []pipeline.CheckpointFingerprint) bool { a = normalizeFingerprints(a) b = normalizeFingerprints(b) if len(a) != len(b) { return false } for i := range a { if a[i] != b[i] { return false } } return true } func reusedDecision() pipeline.CheckpointDecision { return pipeline.CheckpointDecision{Reused: true, Reason: "checkpoint is valid"} } func invalidDecision(format string, args ...any) pipeline.CheckpointDecision { return pipeline.CheckpointDecision{Reason: fmt.Sprintf(format, args...)} }