package appendorder import ( "context" "encoding/json" "reflect" "strings" "testing" "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" "gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline" ) func TestModuleSpecAndRegister(t *testing.T) { want := pipeline.ModuleSpec{ Key: Key, Stage: pipeline.StageMerge, Provides: []string{"merged"}, } if got := ModuleSpec(); !reflect.DeepEqual(got, want) { t.Fatalf("ModuleSpec() = %#v, want %#v", got, want) } registry := pipeline.NewMergerRegistry() if err := Register(registry); err != nil { t.Fatalf("Register() error = %v, want nil", err) } spec, ok := registry.Spec(Key) if !ok { t.Fatalf("Spec(%q) ok = false, want true", Key) } if !reflect.DeepEqual(spec, want) { t.Fatalf("registered spec = %#v, want %#v", spec, want) } } func TestMergePassesThroughSingleExtractOutput(t *testing.T) { input := extractOutput("chunk-0", 0, `{"name":"original"}`) result, err := New().Merge(context.Background(), contracts.MergeRequest{ LaneID: "events", ExtractOutputs: []contracts.ExtractOutput{input}, }) if err != nil { t.Fatalf("Merge() error = %v, want nil", err) } if result.Output.LaneID != "events" || result.Output.MergerKey != Key { t.Fatalf("output provenance = %#v, want lane and merger", result.Output) } if string(result.Output.Payload.Content) != `{"name":"original"}` { t.Fatalf("content = %s, want original content", result.Output.Payload.Content) } if result.Output.Payload.Metadata["name"] != "chunk-0" { t.Fatalf("metadata = %#v, want original metadata", result.Output.Payload.Metadata) } } func TestMergeDefensivelyCopiesRawPayload(t *testing.T) { input := extractOutput("chunk-0", 0, `{"name":"original"}`) result, err := New().Merge(context.Background(), contracts.MergeRequest{ LaneID: "events", ExtractOutputs: []contracts.ExtractOutput{input}, }) if err != nil { t.Fatalf("Merge() error = %v, want nil", err) } input.Payload.Content[0] = '[' input.Payload.Metadata["name"] = "changed" if string(result.Output.Payload.Content) != `{"name":"original"}` { t.Fatalf("content changed after input mutation: %s", result.Output.Payload.Content) } if result.Output.Payload.Metadata["name"] != "chunk-0" { t.Fatalf("metadata changed after input mutation: %#v", result.Output.Payload.Metadata) } } func TestMergeConcatenatesCommonTopLevelArrayFieldInChunkOrder(t *testing.T) { result, err := New().Merge(context.Background(), contracts.MergeRequest{ LaneID: "events", ExtractOutputs: []contracts.ExtractOutput{ extractOutput("chunk-1", 1, `{"events":[{"name":"second"}]}`), extractOutput("chunk-0", 0, `{"events":[{"name":"first"}]}`), }, }) if err != nil { t.Fatalf("Merge() error = %v, want nil", err) } if result.Output.Payload.MediaType != "application/json" { t.Fatalf("MediaType = %q, want application/json", result.Output.Payload.MediaType) } var decoded struct { Events []struct { Name string `json:"name"` } `json:"events"` } if err := json.Unmarshal(result.Output.Payload.Content, &decoded); err != nil { t.Fatalf("Unmarshal() error = %v, want nil", err) } if len(decoded.Events) != 2 || decoded.Events[0].Name != "first" || decoded.Events[1].Name != "second" { t.Fatalf("events = %#v, want concatenated chunk order", decoded.Events) } if result.Output.Schema.ID != "schema-id" { t.Fatalf("schema = %#v, want common extract schema", result.Output.Schema) } } func TestMergeFallsBackToOrderedJSONValueArrayWhenShapesDiffer(t *testing.T) { result, err := New().Merge(context.Background(), contracts.MergeRequest{ LaneID: "events", ExtractOutputs: []contracts.ExtractOutput{ extractOutput("chunk-1", 1, `{"notes":["second"]}`), extractOutput("chunk-0", 0, `{"events":[{"name":"first"}]}`), }, }) if err != nil { t.Fatalf("Merge() error = %v, want nil", err) } var decoded []map[string]any if err := json.Unmarshal(result.Output.Payload.Content, &decoded); err != nil { t.Fatalf("Unmarshal() error = %v, want nil", err) } if len(decoded) != 2 { t.Fatalf("len(decoded) = %d, want 2", len(decoded)) } if _, ok := decoded[0]["events"]; !ok { t.Fatalf("decoded[0] = %#v, want first chunk value", decoded[0]) } if _, ok := decoded[1]["notes"]; !ok { t.Fatalf("decoded[1] = %#v, want second chunk value", decoded[1]) } } func TestMergeRejectsInvalidJSONAndNonJSONMediaTypes(t *testing.T) { tests := []struct { name string output contracts.ExtractOutput want string }{ { name: "invalid JSON", output: extractOutput("chunk-0", 0, `{"events":[`), want: "invalid JSON", }, { name: "non JSON media type", output: func() contracts.ExtractOutput { output := extractOutput("chunk-0", 0, `{"events":[]}`) output.Payload.MediaType = "text/plain" return output }(), want: "unsupported media type", }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { _, err := New().Merge(context.Background(), contracts.MergeRequest{ LaneID: "events", ExtractOutputs: []contracts.ExtractOutput{test.output}, }) if err == nil { t.Fatal("Merge() error = nil, want error") } if !strings.Contains(err.Error(), test.want) { t.Fatalf("Merge() error = %q, want %q", err.Error(), test.want) } }) } } func extractOutput(chunkID string, chunkIndex int, content string) contracts.ExtractOutput { return contracts.ExtractOutput{ LaneID: "events", ExtractorKey: "extract", SourceID: "source-1", ChunkID: chunkID, ChunkIndex: chunkIndex, Schema: contracts.ResponseSchema{ID: "schema-id", Name: "schema-name", Version: "v1"}, Payload: contracts.RawPayload{ Content: []byte(content), MediaType: "application/json", Metadata: map[string]any{"name": chunkID}, }, } }