package integration_test import ( "context" "encoding/json" "fmt" "os" "strings" "testing" "gitea.maximumdirect.net/eric/notarius/internal/core/config" "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" "gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd" itemeventcodec "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/codec/itemevents" itemeventextract "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/extract/itemevents" itemeventnormalize "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/normalize/itemevents" itemeventshape "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/validate/itemevents/shape" itemeventrefs "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/validate/itemevents/source_refs" "gitea.maximumdirect.net/eric/notarius/internal/modules/seriatim/input/transcript" ) func TestProductionItemEventPipelineProducesNormalizedDurableArtifact(t *testing.T) { raw := readItemEventFixture(t) registries := productionNPCRegistries(t) resolved := resolveItemEventPipeline(t, registries) client := &itemEventProductionLLMClient{responses: [][]byte{[]byte(`{"events":[ {"name":"Healing Potion","kind":"transferred","quantity":1,"from":"Aria","to":"Borin","source_refs":[{"start_segment":3,"end_segment":3}]}, {"name":" Gold Pieces ","kind":"acquired","quantity":25,"to":" party ","source_refs":[{"start_segment":2,"end_segment":2}]}, {"name":"Missing Relic","kind":"discovered","source_refs":[{"start_segment":4,"end_segment":4}]}, {"name":"Ancient Coin","kind":"discovered","source_refs":[{"start_segment":1,"end_segment":1}]}, {"name":"Gold Pieces","kind":"acquired","quantity":25,"to":"party","source_refs":[{"start_segment":2,"end_segment":2}]} ]}`)}} output, err := runPreparedPipeline(t, registries, resolved, client, pipeline.RunInput{RawInput: raw}) if err != nil { t.Fatalf("Run() error = %v", err) } if len(output.Rejected) != 0 || len(output.NormalizeOutputs) != 1 { t.Fatalf("run output = %#v, rejected = %#v", output.NormalizeOutputs, output.Rejected) } serialized := normalizedLane(t, output, "item-events") if serialized.NormalizerKey != itemeventnormalize.Key || serialized.Artifact.Schema.ID != itemeventcodec.SchemaID || serialized.Artifact.Schema.Version != itemeventcodec.SchemaVersion { t.Fatalf("serialized item events = %#v", serialized) } events, err := itemeventcodec.New().Decode(serialized.Artifact.Content) if err != nil { t.Fatalf("Decode(item event output) error = %v", err) } if len(events.Events) != 4 { t.Fatalf("item events = %#v, want one exact duplicate removed", events) } first, second, third, fourth := events.Events[0], events.Events[1], events.Events[2], events.Events[3] if first.Name != "Ancient Coin" || second.Name != "Gold Pieces" || second.Kind != dnd.ItemEventKindAcquired || second.Quantity == nil || *second.Quantity != 25 || second.To != "party" || third.Name != "Healing Potion" || third.Kind != dnd.ItemEventKindTransferred || third.From != "Aria" || third.To != "Borin" || fourth.Name != "Missing Relic" { t.Fatalf("normalized item event sequence = %#v", events.Events) } for _, event := range events.Events { for _, ref := range event.SourceRefs { if ref.SourceID != "item-events-session" { t.Fatalf("source reference = %#v, want current transcript source", ref) } } } if !hasItemEventWarning(output.Warnings, "item_event_source_unrelated") || !hasItemEventWarning(output.Warnings, "duplicate_item_event_collapsed") { t.Fatalf("warnings = %#v, want advisory relatedness and duplicate-collapse warnings", output.Warnings) } if output.Manifest.ValidationStatus != "approved" || len(output.Manifest.ArtifactLanes) != 1 { t.Fatalf("manifest = %#v", output.Manifest) } lane := output.Manifest.ArtifactLanes[0] if lane.ID != "item-events" || lane.Extractor != itemeventextract.Key || lane.Merger != pipeline.DefaultMergeModule || lane.Normalizer != itemeventnormalize.Key { t.Fatalf("manifest lane = %#v", lane) } if len(client.requests) != 1 || client.requests[0].PromptID != itemeventextract.PromptID { t.Fatalf("LLM requests = %#v", client.requests) } } func TestProductionItemEventPipelineRetriesBlockingCandidates(t *testing.T) { registries := productionNPCRegistries(t) resolved := resolveItemEventPipeline(t, registries) for _, test := range []struct { name string response string reasonCode string validatorName string }{ {name: "shape", response: `{"events":[{"name":"","kind":"discovered","source_refs":[{"start_segment":1,"end_segment":1}]}]}`, reasonCode: itemeventshape.ReasonCode, validatorName: itemeventshape.Key}, {name: "source references", response: `{"events":[{"name":"Ancient Coin","kind":"discovered","source_refs":[{"start_segment":99,"end_segment":99}]}]}`, reasonCode: itemeventrefs.ReasonCode, validatorName: itemeventrefs.Key}, } { t.Run(test.name, func(t *testing.T) { client := &itemEventProductionLLMClient{responses: [][]byte{[]byte(test.response), []byte(test.response), []byte(test.response)}} output, err := runPreparedPipeline(t, registries, resolved, client, pipeline.RunInput{RawInput: readItemEventFixture(t)}) if err != nil { t.Fatalf("Run() error = %v; want rejected candidate outcome", err) } if len(output.NormalizeOutputs) != 0 || len(output.Rejected) != 1 || len(client.requests) != 3 { t.Fatalf("output=%#v rejected=%#v requests=%d", output.NormalizeOutputs, output.Rejected, len(client.requests)) } rejection := output.Rejected[0] if rejection.ReasonCode != test.reasonCode || rejection.ValidatorName != test.validatorName || rejection.AttemptCount != 3 { t.Fatalf("rejection = %#v", rejection) } }) } } func resolveItemEventPipeline(t *testing.T, registries pipeline.Registries) pipeline.ResolvedPipeline { t.Helper() configValue := config.Default() configValue.Pipelines["dnd-item-events-fixture"] = pipeline.PipelineProfile{ Input: pipeline.Binding(transcript.Key), Chunk: pipeline.ModuleBinding{Module: pipeline.DefaultChunkModule, Options: map[string]any{"max_units": 100}}, Artifacts: map[string]pipeline.ArtifactLaneProfile{ "item-events": { Extract: pipeline.ModuleBinding{Module: itemeventextract.Key, Retries: 2}, Normalize: pipeline.Binding(itemeventnormalize.Key), }, }, } effective, err := configValue.Resolve(config.ResolveInput{PipelineID: "dnd-item-events-fixture", Catalog: moduleCatalog(registries)}) if err != nil { t.Fatalf("Resolve() error = %v", err) } return effective.ResolvedPipeline } func readItemEventFixture(t *testing.T) []byte { t.Helper() raw, err := os.ReadFile("testdata/seriatim_item_events_session.json") if err != nil { t.Fatal(err) } return raw } func hasItemEventWarning(warnings []contracts.Warning, reasonCode string) bool { for _, warning := range warnings { if warning.ReasonCode == reasonCode && strings.HasPrefix(warning.Scope, "events[") { return true } } return false } type itemEventProductionLLMClient struct { responses [][]byte requests []contracts.StructuredCompletionRequest } func (client *itemEventProductionLLMClient) CompleteStructured(ctx context.Context, request contracts.StructuredCompletionRequest, out any) (contracts.StructuredCompletionResponse, error) { if err := ctx.Err(); err != nil { return contracts.StructuredCompletionResponse{}, err } if request.PromptID != itemeventextract.PromptID { return contracts.StructuredCompletionResponse{}, fmt.Errorf("unexpected prompt %q", request.PromptID) } index := len(client.requests) if index >= len(client.responses) { return contracts.StructuredCompletionResponse{}, fmt.Errorf("missing item event response %d", index) } content := append([]byte(nil), client.responses[index]...) if err := json.Unmarshal(content, out); err != nil { return contracts.StructuredCompletionResponse{}, fmt.Errorf("populate item event structured target: %w", err) } client.requests = append(client.requests, request) return contracts.StructuredCompletionResponse{Content: content, Provider: "test", Model: "item-events-fake"}, nil } var _ contracts.StructuredLLMClient = (*itemEventProductionLLMClient)(nil)