diff --git a/docs/roadmap/implementation.md b/docs/roadmap/implementation.md index ad4c4b3..1f8a8bc 100644 --- a/docs/roadmap/implementation.md +++ b/docs/roadmap/implementation.md @@ -798,7 +798,7 @@ compatibility tests pass. ## Stage 7: Add cross-run, corruption, and concurrency hardening -**Status:** Not started +**Status:** Complete ### Objective diff --git a/internal/core/source/chunk_plan.go b/internal/core/source/chunk_plan.go index 0d4a89d..5150c1e 100644 --- a/internal/core/source/chunk_plan.go +++ b/internal/core/source/chunk_plan.go @@ -207,7 +207,34 @@ func cloneJSONMap(values map[string]any) map[string]any { } cloned := make(map[string]any, len(values)) for key, value := range values { - cloned[key] = value + cloned[key] = cloneJSONValue(value) } return cloned } + +func cloneJSONValue(value any) any { + switch typed := value.(type) { + case map[string]any: + return cloneJSONMap(typed) + case []any: + cloned := make([]any, len(typed)) + for i := range typed { + cloned[i] = cloneJSONValue(typed[i]) + } + return cloned + case json.RawMessage: + return append(json.RawMessage(nil), typed...) + case []byte: + return append([]byte(nil), typed...) + case []string: + return append([]string(nil), typed...) + case map[string]string: + cloned := make(map[string]string, len(typed)) + for key, item := range typed { + cloned[key] = item + } + return cloned + default: + return value + } +} diff --git a/internal/core/source/chunk_plan_test.go b/internal/core/source/chunk_plan_test.go index 8f0c9c7..beaf600 100644 --- a/internal/core/source/chunk_plan_test.go +++ b/internal/core/source/chunk_plan_test.go @@ -202,6 +202,22 @@ func TestMaterializeChunkPlanExactOutputAndMutationSafety(t *testing.T) { } } +func TestMaterializeChunkPlanDeepClonesUnitMetadata(t *testing.T) { + doc := planDocument() + doc.Units[0].Metadata = map[string]any{"nested": map[string]any{"values": []any{json.RawMessage(`{"ok":true}`)}}} + chunks, err := MaterializeChunkPlan(doc, validChunkPlan(doc)) + if err != nil { + t.Fatal(err) + } + nested := chunks[0].Units[0].Metadata["nested"].(map[string]any) + nested["values"].([]any)[0].(json.RawMessage)[0] = '[' + nested["changed"] = true + original := doc.Units[0].Metadata["nested"].(map[string]any) + if _, exists := original["changed"]; exists || string(original["values"].([]any)[0].(json.RawMessage)) != `{"ok":true}` { + t.Fatalf("source metadata changed through materialized chunk: %#v", doc.Units[0].Metadata) + } +} + func planDocument() *SourceDocument { doc := &SourceDocument{ID: "source-plan", Kind: "test", Format: "application/test", Digest: "sha256:source-plan"} for _, id := range []int{10, 20, 30, 40, 50} { diff --git a/internal/framework/chunkplan/import_boundaries_test.go b/internal/framework/chunkplan/import_boundaries_test.go new file mode 100644 index 0000000..183b77f --- /dev/null +++ b/internal/framework/chunkplan/import_boundaries_test.go @@ -0,0 +1,62 @@ +package chunkplan + +import ( + "go/parser" + "go/token" + "path/filepath" + "runtime" + "strconv" + "strings" + "testing" +) + +const repositoryImportPrefix = "gitea.maximumdirect.net/eric/notarius/internal/" + +func TestPlanStoreAndSourceImportBoundaries(t *testing.T) { + repositoryRoot := repositoryRoot(t) + for _, tc := range []struct { + name string + directory string + forbidden []string + }{ + {name: "source is framework and module independent", directory: "internal/core/source", forbidden: []string{"framework/", "modules/"}}, + {name: "plan store is module independent", directory: "internal/framework/chunkplan", forbidden: []string{"modules/"}}, + } { + t.Run(tc.name, func(t *testing.T) { + files, err := filepath.Glob(filepath.Join(repositoryRoot, tc.directory, "*.go")) + if err != nil { + t.Fatal(err) + } + for _, filename := range files { + if strings.HasSuffix(filename, "_test.go") { + continue + } + parsed, err := parser.ParseFile(token.NewFileSet(), filename, nil, parser.ImportsOnly) + if err != nil { + t.Fatal(err) + } + for _, item := range parsed.Imports { + path, err := strconv.Unquote(item.Path.Value) + if err != nil { + t.Fatal(err) + } + path = strings.TrimPrefix(path, repositoryImportPrefix) + for _, prefix := range tc.forbidden { + if strings.HasPrefix(path, prefix) { + t.Fatalf("%s imports %q, forbidden by %s boundary", filepath.Base(filename), path, tc.name) + } + } + } + } + }) + } +} + +func repositoryRoot(t *testing.T) string { + t.Helper() + _, filename, _, ok := runtime.Caller(0) + if !ok { + t.Fatal("resolve test location") + } + return filepath.Clean(filepath.Join(filepath.Dir(filename), "..", "..", "..")) +} diff --git a/internal/framework/chunkplan/store.go b/internal/framework/chunkplan/store.go index 304904f..97a9697 100644 --- a/internal/framework/chunkplan/store.go +++ b/internal/framework/chunkplan/store.go @@ -17,7 +17,8 @@ import ( const SchemaVersion = pipeline.ChunkPlanSchemaVersion type filesystemStore struct { - root string + root string + write func(string, []byte) error } func NewFilesystemStore(root string) (pipeline.ChunkPlanStore, error) { @@ -28,7 +29,7 @@ func NewFilesystemStore(root string) (pipeline.ChunkPlanStore, error) { if strings.ContainsRune(root, '\x00') { return nil, fmt.Errorf("chunk plan root must not contain NUL") } - return &filesystemStore{root: filepath.Clean(root)}, nil + return &filesystemStore{root: filepath.Clean(root), write: writeAtomic}, nil } func (s *filesystemStore) Load(sourceDigest string) (pipeline.ChunkPlanRecord, pipeline.ChunkPlanDecision, error) { @@ -77,7 +78,11 @@ func (s *filesystemStore) Save(record pipeline.ChunkPlanRecord) error { return fmt.Errorf("encode chunk plan record: %w", err) } data = append(data, '\n') - if err := writeAtomic(target, data); err != nil { + writer := s.write + if writer == nil { + writer = writeAtomic + } + if err := writer(target, data); err != nil { return fmt.Errorf("write chunk plan: %w", err) } return nil @@ -164,7 +169,16 @@ func invalidDecision(reason string) (pipeline.ChunkPlanRecord, pipeline.ChunkPla return pipeline.ChunkPlanRecord{}, pipeline.ChunkPlanDecision{Status: pipeline.ChunkPlanInvalid, Reason: reason}, nil } +type atomicWriteHooks struct { + BeforeCreateTemp func() error + BeforeRename func() error +} + func writeAtomic(target string, data []byte) error { + return writeAtomicWithHooks(target, data, atomicWriteHooks{}) +} + +func writeAtomicWithHooks(target string, data []byte, hooks atomicWriteHooks) error { dir := filepath.Dir(target) if err := os.MkdirAll(dir, 0o700); err != nil { return err @@ -173,6 +187,11 @@ func writeAtomic(target string, data []byte) error { return err } + if hooks.BeforeCreateTemp != nil { + if err := hooks.BeforeCreateTemp(); err != nil { + return err + } + } temp, err := os.CreateTemp(dir, ".plan.json.tmp-*") if err != nil { return err @@ -199,6 +218,11 @@ func writeAtomic(target string, data []byte) error { if err := temp.Close(); err != nil { return err } + if hooks.BeforeRename != nil { + if err := hooks.BeforeRename(); err != nil { + return err + } + } if err := os.Rename(tempPath, target); err != nil { return err } diff --git a/internal/framework/chunkplan/store_test.go b/internal/framework/chunkplan/store_test.go index a4e4646..6634215 100644 --- a/internal/framework/chunkplan/store_test.go +++ b/internal/framework/chunkplan/store_test.go @@ -3,6 +3,8 @@ package chunkplan import ( "bytes" "encoding/json" + "errors" + "fmt" "os" "path/filepath" "reflect" @@ -125,6 +127,7 @@ func TestFilesystemStoreReportsInvalidRecordsAsRecoverable(t *testing.T) { {name: "unknown field", mutate: func(data []byte) []byte { return bytes.Replace(data, []byte(`{"schema_version"`), []byte(`{"unknown":true,"schema_version"`), 1) }, want: "unknown"}, + {name: "truncated JSON", mutate: func(data []byte) []byte { return data[:len(data)/2] }, want: "decode"}, {name: "schema mismatch", mutate: replaceJSON(`notarius.chunk-plan.v1`, `notarius.chunk-plan.v2`), want: "schema_version"}, {name: "source mismatch", mutate: replaceJSON(testSourceDigest, "sha256:"+strings.Repeat("b", 64)), want: "source_digest"}, {name: "plan digest mismatch", mutate: func(data []byte) []byte { @@ -138,6 +141,9 @@ func TestFilesystemStoreReportsInvalidRecordsAsRecoverable(t *testing.T) { {name: "noncanonical annotation", mutate: func(data []byte) []byte { return bytes.Replace(data, []byte(`{"value":1}`), []byte(`{ "value": 1 }`), 1) }, want: "canonical JSON"}, + {name: "bad boundary", mutate: func(data []byte) []byte { + return bytes.Replace(data, []byte(`"start_unit_id":1`), []byte(`"start_unit_id":0`), 1) + }, want: "boundaries must be positive"}, {name: "trailing JSON", mutate: func(data []byte) []byte { return append(data, []byte(` {}`)...) }, want: "trailing"}, } for _, tc := range tests { @@ -226,6 +232,77 @@ func TestFilesystemStoreConcurrentWritersExposeCompleteRecord(t *testing.T) { } } +func TestFilesystemStoreReadersObserveOnlyCompleteRecordsDuringWrites(t *testing.T) { + store := newStore(t, t.TempDir()) + if err := store.Save(testRecord(t, 1)); err != nil { + t.Fatal(err) + } + const writers = 12 + const readers = 12 + errs := make(chan error, writers+readers) + start := make(chan struct{}) + var writersDone sync.WaitGroup + for i := 0; i < writers; i++ { + writersDone.Add(1) + go func(value int) { + defer writersDone.Done() + <-start + errs <- store.Save(testRecord(t, value+2)) + }(i) + } + for i := 0; i < readers; i++ { + go func() { + <-start + for attempt := 0; attempt < 50; attempt++ { + record, decision, err := store.Load(testSourceDigest) + if err != nil || decision.Status != pipeline.ChunkPlanHit { + errs <- fmt.Errorf("Load() decision=%#v error=%v", decision, err) + return + } + if err := validateRecord(record, testSourceDigest); err != nil { + errs <- fmt.Errorf("reader observed invalid record: %w", err) + return + } + } + errs <- nil + }() + } + close(start) + writersDone.Wait() + for i := 0; i < readers+writers; i++ { + if err := <-errs; err != nil { + t.Fatal(err) + } + } +} + +func TestFilesystemStoreInterruptedWritesPreservePreviousRecord(t *testing.T) { + store := newStore(t, t.TempDir()).(*filesystemStore) + first := testRecord(t, 1) + if err := store.Save(first); err != nil { + t.Fatal(err) + } + for _, tc := range []struct { + name string + hooks atomicWriteHooks + }{ + {name: "before temporary file", hooks: atomicWriteHooks{BeforeCreateTemp: func() error { return errors.New("interrupted before temporary file") }}}, + {name: "before rename", hooks: atomicWriteHooks{BeforeRename: func() error { return errors.New("interrupted before rename") }}}, + } { + t.Run(tc.name, func(t *testing.T) { + store.write = func(target string, data []byte) error { return writeAtomicWithHooks(target, data, tc.hooks) } + if err := store.Save(testRecord(t, 2)); err == nil { + t.Fatal("Save() error = nil") + } + store.write = writeAtomic + got, decision, err := store.Load(testSourceDigest) + if err != nil || decision.Status != pipeline.ChunkPlanHit || !reflect.DeepEqual(got, first) { + t.Fatalf("record after interruption=%#v decision=%#v error=%v", got, decision, err) + } + }) + } +} + func newStore(t *testing.T, root string) pipeline.ChunkPlanStore { t.Helper() store, err := NewFilesystemStore(root) diff --git a/internal/framework/pipeline/runner_chunk_plan_test.go b/internal/framework/pipeline/runner_chunk_plan_test.go index 4410f96..6acef92 100644 --- a/internal/framework/pipeline/runner_chunk_plan_test.go +++ b/internal/framework/pipeline/runner_chunk_plan_test.go @@ -2,6 +2,7 @@ package pipeline import ( "context" + "encoding/json" "errors" "reflect" "strings" @@ -242,6 +243,22 @@ func TestRunnerRegeneratesStructurallyInvalidHit(t *testing.T) { } } +func TestRunnerReusesCrossDomainAnnotationsAsOptionalData(t *testing.T) { + prepared, plan := preparedTerminalDebugPipeline(t) + plan.Annotations = source.ChunkAnnotations{"dnd/scenes": json.RawMessage(`{"boundary_caveats":["uncertain"]}`)} + plan.Ranges[0].Annotations = source.ChunkAnnotations{"dnd/scenes": json.RawMessage(`{"title":"Opening"}`)} + record := chunkPlanRecord(t, prepared, plan) + record.Producer.ChunkModule = "dnd/scenes" + store := &recordingChunkPlanStore{record: record, decision: ChunkPlanDecision{Status: ChunkPlanHit}} + output, err := New().Run(context.Background(), RunInput{Prepared: prepared, RawInput: []byte("input"), ChunkCacheMode: ChunkCacheAuto, ChunkPlans: store}) + if err != nil { + t.Fatal(err) + } + if output.Manifest.ChunkPlan.ProducerModule != "dnd/scenes" || output.Manifest.ChunkPlan.Action != "reused" { + t.Fatalf("cross-domain annotation plan was not reused: %#v", output.Manifest.ChunkPlan) + } +} + func TestRunnerRetriesBeforePublishingAcceptedPlan(t *testing.T) { prepared, plan := preparedTerminalDebugPipeline(t) prepared.resolved.Chunk.Retries = 1