diff --git a/assets/dnd/item-registry/normalize/prompts/candidates.md b/assets/dnd/item-registry/normalize/prompts/candidates.md deleted file mode 100644 index 33f4ca3..0000000 --- a/assets/dnd/item-registry/normalize/prompts/candidates.md +++ /dev/null @@ -1,2 +0,0 @@ -Item candidates: -{{ input "candidates" }} diff --git a/assets/dnd/item-registry/normalize/prompts/instructions.md b/assets/dnd/item-registry/normalize/prompts/instructions.md index 4823409..640e355 100644 --- a/assets/dnd/item-registry/normalize/prompts/instructions.md +++ b/assets/dnd/item-registry/normalize/prompts/instructions.md @@ -1,8 +1,9 @@ -Use candidate names and cited transcript windows only to determine whether -candidates identify the same item type or unique designation. Do not treat -nearby evidence, similar objects, or a shared owner as sufficient. Keep -currency denominations, materially different item types, and uncertain aliases -separate. Do not infer an item property or uniqueness. +Use the supplied positive integer candidate IDs and their cited transcript +windows only to determine whether candidates identify the same item type or +unique designation. Do not treat nearby evidence, similar objects, or a shared +owner as sufficient. Keep currency denominations, materially different item +types, and uncertain aliases separate. Do not infer an item property or +uniqueness. When selecting a canonical display name, choose one supplied candidate name that is the clearest established designation. diff --git a/assets/dnd/item-registry/normalize/prompts/prompt.yaml b/assets/dnd/item-registry/normalize/prompts/prompt.yaml index a8310f3..c1d8ece 100644 --- a/assets/dnd/item-registry/normalize/prompts/prompt.yaml +++ b/assets/dnd/item-registry/normalize/prompts/prompt.yaml @@ -12,19 +12,19 @@ messages: - role: system content_file: ./sharedassets/common-dnd-system.md - role: user - content_file: ./instructions.md + content_file: ./sharedassets/protocol.md - role: user - content_file: ./sharedassets/common-dnd-entity-reconciliation.md + content_file: ./instructions.md cache_control: type: ephemeral - role: user - content_file: ./candidates.md + content_file: ./sharedassets/candidates.md - role: user - content_file: ./sharedassets/common-dnd-transcript-windows.md + content_file: ./sharedassets/transcript-windows.md cache_control: type: ephemeral output: format: json validation_mode: json_schema - schema_path: dnd_entity_reconcile_llm.v1.json + schema_path: semantic_reconciliation_llm.v1.json repair_attempts: 0 diff --git a/internal/modules/dnd/normalize/itemregistry/normalizer.go b/internal/modules/dnd/normalize/itemregistry/normalizer.go index e634579..c8a1b8d 100644 --- a/internal/modules/dnd/normalize/itemregistry/normalizer.go +++ b/internal/modules/dnd/normalize/itemregistry/normalizer.go @@ -3,7 +3,6 @@ package itemregistry import ( "context" - "errors" "fmt" "reflect" "sort" @@ -13,20 +12,19 @@ import ( "gitea.maximumdirect.net/eric/notarius/internal/core/source" "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" "gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline" + "gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/items/identity" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/diagnostics" - "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/entityreconcile" ) const ( - Key = "dnd/item-registry" - PromptID = "dnd.item_registry.normalize" - normalizationPolicy = "dnd.item_registry.normalize.v2" - semanticContextPolicy = "dnd.entity_reconcile.context.v1" - semanticContextRadius = 2 - NormalizationPolicy = normalizationPolicy + Key = "dnd/item-registry" + PromptID = "dnd.item_registry.normalize" + PromptVersion = "v1" + normalizationPolicy = "dnd.item_registry.normalize.v3" + NormalizationPolicy = normalizationPolicy ReasonCodeItemFieldsNormalized = "item_fields_normalized" ReasonCodeItemIDRecomputed = "item_id_recomputed" @@ -47,9 +45,7 @@ var _ pipeline.CheckpointFingerprintProvider = (*Normalizer)(nil) type Options struct{} type Normalizer struct { - llm contracts.StructuredLLMClient - promptSHA string - responseSchemaSHA string + engine *semanticreconcile.Engine } func New(llmClient contracts.StructuredLLMClient, _ Options) (*Normalizer, error) { @@ -60,46 +56,45 @@ func New(llmClient contracts.StructuredLLMClient, _ Options) (*Normalizer, error if err != nil { return nil, normalizerErrorf("load prompt metadata: %w", err) } - responseSchema, err := entityreconcile.LoadResponseSchema() + engine, err := semanticreconcile.NewEngine(llmClient, semanticreconcile.PromptSpec{ + ID: PromptID, Version: PromptVersion, SHA256: promptSHA, + }, semanticreconcile.DefaultLimits()) if err != nil { - return nil, normalizerErrorf("load response schema: %w", err) + return nil, normalizerErrorf("construct semantic reconciliation engine: %w", err) } - return &Normalizer{llm: llmClient, promptSHA: promptSHA, responseSchemaSHA: responseSchema.SHA256}, nil + return &Normalizer{engine: engine}, nil } func (n *Normalizer) Key() string { return Key } func (n *Normalizer) ReferenceSlots() []contracts.ReferenceSlot { return nil } func (n *Normalizer) ManifestMetadata() map[string]any { - if n == nil { + if n == nil || n.engine == nil { return nil } - return map[string]any{ - "prompt_id": PromptID, "prompt_version": entityreconcile.SchemaVersion, "prompt_sha256": n.promptSHA, - "response_schema_key": string(entityreconcile.ResponseSchemaKey), "response_schema_id": entityreconcile.ResponseSchemaID, - "response_schema_name": entityreconcile.ResponseSchemaName, "response_schema_version": entityreconcile.SchemaVersion, - "response_schema_sha256": n.responseSchemaSHA, "identity_policy": identity.Policy, - "normalization_policy": normalizationPolicy, "semantic_context_policy": semanticContextPolicy, "semantic_context_radius": semanticContextRadius, - } + metadata := n.engine.ManifestMetadata() + metadata["identity_policy"] = identity.Policy + metadata["normalization_policy"] = normalizationPolicy + return metadata } func (n *Normalizer) CheckpointFingerprints() []pipeline.CheckpointFingerprint { - if n == nil { + if n == nil || n.engine == nil { return nil } - return []pipeline.CheckpointFingerprint{ - {Name: "prompt", Value: n.promptSHA}, {Name: "response_schema", Value: n.responseSchemaSHA}, - {Name: "identity_policy", Value: identity.Policy}, {Name: "normalization_policy", Value: normalizationPolicy}, - {Name: "semantic_context_policy", Value: fmt.Sprintf("%s:%d", semanticContextPolicy, semanticContextRadius)}, - } + fingerprints := n.engine.CheckpointFingerprints() + return append(fingerprints, + pipeline.CheckpointFingerprint{Name: "identity_policy", Value: identity.Policy}, + pipeline.CheckpointFingerprint{Name: "normalization_policy", Value: normalizationPolicy}, + ) } func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalizeRequest[dnd.ItemRegistry]) (contracts.TypedNormalizeResult[dnd.ItemRegistry], error) { if n == nil { return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("normalizer must not be nil") } - if n.llm == nil { - return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("LLM client must not be nil") + if n.engine == nil { + return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("semantic reconciliation engine must not be nil") } if ctx == nil { return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("context must not be nil") @@ -111,35 +106,44 @@ func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalize order := shared.NewSourceRefOrder(req.Source) records, warnings := preprocessRecords(req.MergeOutput.Value, order) deterministic := recordList(records) - materials, ready, err := entityreconcile.BuildContext(req.Source, reconciliationCandidates(records), semanticContextRadius) - if err != nil { - return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("build semantic context: %w", err) - } - if !ready { + if len(records) < 2 || req.Source == nil { return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil } - var response entityreconcile.ProposalResponse - if _, err := n.llm.CompleteStructured(ctx, contracts.StructuredCompletionRequest{ - StageName: Key, PromptID: PromptID, PromptVersion: entityreconcile.SchemaVersion, + candidates, envelopes, err := reconciliationInputs(records) + if err != nil { + return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("prepare semantic reconciliation inputs: %w", err) + } + reconciliation, err := n.engine.Reconcile(ctx, semanticreconcile.Request{ + StageName: Key, Source: req.Source, Candidates: candidates, ProfileID: req.LLMProfile, SessionID: req.SessionID, - Inputs: contracts.LLMInputSet{"candidates": materials.Candidates, "transcript": materials.Transcript}, - }, &response); err != nil { - if errors.Is(err, contracts.ErrInvalidStructuredOutput) { - return n.invalidStructuredResult(deterministic, warnings), nil - } - return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("complete structured output: %w", err) + }) + if err != nil { + return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("reconcile semantic duplicates: %w", err) + } + + switch reconciliation.Disposition() { + case semanticreconcile.SkippedInsufficientCandidates: + return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil + case semanticreconcile.SkippedLimitExceeded: + return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: deterministic, Warnings: limitWarningsWithSemanticFallback(warnings)}, nil + case semanticreconcile.RetryableInvalidStructuredOutput: + return n.invalidStructuredResult(deterministic, warnings), nil + case semanticreconcile.Complete, semanticreconcile.RetryableDiscardedProposalGroups: + default: + return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("unknown semantic reconciliation disposition %d", reconciliation.Disposition()) + } + + applied, semanticWarnings, rejectedGroups, err := applyReconciliationPlan(reconciliation.Plan(), records, envelopes, order) + if err != nil { + return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("apply semantic reconciliation plan: %w", err) } - assessment := materials.Assess(response) - applied, semanticWarnings, rejectedGroups := applySafeGroups(records, reconciliationGroups(assessment, materials.CandidateKeys()), order) warnings = append(warnings, semanticWarnings...) - if rejectedGroups > 0 { - return currencyRetryResult(recordList(applied), warnings, rejectedGroups), nil - } - if assessment.DiscardedGroups() == 0 { + discardedGroups := reconciliation.DiscardedGroupCount() + rejectedGroups + if discardedGroups == 0 { return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: recordList(applied), Warnings: limitWarnings(warnings)}, nil } - return retryResult(recordList(applied), warnings, assessment), nil + return retryResult(recordList(applied), warnings, reconciliation, rejectedGroups), nil } func (n *Normalizer) invalidStructuredResult(value dnd.ItemRegistry, warnings []contracts.Warning) contracts.TypedNormalizeResult[dnd.ItemRegistry] { @@ -149,17 +153,15 @@ func (n *Normalizer) invalidStructuredResult(value dnd.ItemRegistry, warnings [] }} } -func retryResult(value dnd.ItemRegistry, warnings []contracts.Warning, assessment entityreconcile.Assessment) contracts.TypedNormalizeResult[dnd.ItemRegistry] { +func retryResult(value dnd.ItemRegistry, warnings []contracts.Warning, reconciliation semanticreconcile.Result, rejectedGroups int) contracts.TypedNormalizeResult[dnd.ItemRegistry] { + details := reconciliationIssues(reconciliation.Issues()) + if rejectedGroups > 0 { + details = append(details, "currency may only be consolidated with aliases of one denomination") + } + discardedGroups := reconciliation.DiscardedGroupCount() + rejectedGroups return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: value, Warnings: limitWarningsForRetry(warnings), Retry: &contracts.NormalizeRetry{ - ReasonCode: ReasonCodeItemSemanticProposalInvalid, Message: diagnostics.Aggregate("semantic proposal requires retry", reconciliationIssues(assessment)), - FallbackWarnings: []contracts.Warning{semanticFallbackWarning(assessment.DiscardedGroups())}, - }} -} - -func currencyRetryResult(value dnd.ItemRegistry, warnings []contracts.Warning, rejectedGroups int) contracts.TypedNormalizeResult[dnd.ItemRegistry] { - return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: value, Warnings: limitWarningsForRetry(warnings), Retry: &contracts.NormalizeRetry{ - ReasonCode: ReasonCodeItemSemanticProposalInvalid, Message: "semantic proposal requires retry: currency may only be consolidated with aliases of one denomination", - FallbackWarnings: []contracts.Warning{semanticFallbackWarning(rejectedGroups)}, + ReasonCode: ReasonCodeItemSemanticProposalInvalid, Message: diagnostics.Aggregate("semantic proposal requires retry", details), + FallbackWarnings: []contracts.Warning{semanticFallbackWarning(discardedGroups)}, }} } @@ -187,6 +189,10 @@ func limitWarningsForRetry(warnings []contracts.Warning) []contracts.Warning { return append(bounded, contracts.Warning{Scope: "items", ReasonCode: ReasonCodeItemNormalizationWarningsOmitted, Message: fmt.Sprintf("%d additional warning(s) omitted", len(warnings)-displayed)}) } +func limitWarningsWithSemanticFallback(warnings []contracts.Warning) []contracts.Warning { + return append(limitWarningsForRetry(warnings), semanticFallbackWarning(-1)) +} + type normalizedRecord struct { item dnd.Item inputIndexes []int diff --git a/internal/modules/dnd/normalize/itemregistry/normalizer_test.go b/internal/modules/dnd/normalize/itemregistry/normalizer_test.go index 2def3a6..f6bf18f 100644 --- a/internal/modules/dnd/normalize/itemregistry/normalizer_test.go +++ b/internal/modules/dnd/normalize/itemregistry/normalizer_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "reflect" "strconv" "strings" @@ -14,10 +15,10 @@ import ( "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" "gitea.maximumdirect.net/eric/notarius/internal/framework/llm" "gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline" + "gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/items/identity" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/diagnostics" - "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/entityreconcile" identityvalidator "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/validate/itemregistry/identity" "gitea.maximumdirect.net/eric/promptkit" ) @@ -35,9 +36,15 @@ func TestModuleContractAndMetadata(t *testing.T) { t.Fatal("New() accepted nil client") } metadata := newNormalizer(t, &recordingNormalizerClient{}).ManifestMetadata() - if metadata["identity_policy"] != identity.Policy || metadata["response_schema_id"] != entityreconcile.ResponseSchemaID || metadata["normalization_policy"] != normalizationPolicy || metadata["semantic_context_radius"] != semanticContextRadius { + limits, ok := metadata["semantic_reconciliation_limits"].(map[string]any) + if !ok || metadata["identity_policy"] != identity.Policy || metadata["normalization_policy"] != normalizationPolicy || metadata["prompt_id"] != PromptID || metadata["prompt_version"] != PromptVersion || metadata["response_schema_key"] != string(semanticreconcile.ResponseSchemaKey) || metadata["response_schema_id"] != semanticreconcile.ResponseSchemaID || metadata["response_schema_name"] != semanticreconcile.ResponseSchemaName || metadata["semantic_reconciliation_policy"] != semanticreconcile.Policy || len(limits) != 3 { t.Fatalf("metadata = %#v", metadata) } + for _, name := range []string{"prompt", "response_schema", "semantic_reconciliation_policy", "semantic_reconciliation_limits", "identity_policy", "normalization_policy"} { + if !hasFingerprint(newNormalizer(t, &recordingNormalizerClient{}).CheckpointFingerprints(), name) { + t.Fatalf("fingerprints missing %q", name) + } + } } func TestNormalizeConsolidatesEqualNamesAcrossEvidenceWithoutMutation(t *testing.T) { @@ -73,10 +80,10 @@ func TestNormalizeConsolidatesEqualNamesAcrossEvidenceWithoutMutation(t *testing } var candidates struct { Candidates []struct { - Name string `json:"name"` + Label string `json:"label"` } `json:"candidates"` } - if err := json.Unmarshal(client.requests[0].Inputs["candidates"].Content, &candidates); err != nil || len(candidates.Candidates) != 2 || candidates.Candidates[0].Name != "Rope" || candidates.Candidates[1].Name != "Lantern" { + if err := json.Unmarshal(client.requests[0].Inputs["candidates"].Content, &candidates); err != nil || len(candidates.Candidates) != 2 || candidates.Candidates[0].Label != "Rope" || candidates.Candidates[1].Label != "Lantern" { t.Fatalf("semantic candidates = %#v, %v; want one candidate per comparison name", candidates, err) } repeated, repeatErr := newNormalizer(t, &recordingNormalizerClient{}).Normalize(context.Background(), normalizeRequestWithSource(result.Value, doc)) @@ -106,7 +113,7 @@ func TestNormalizeAppliesSafeAliasProposal(t *testing.T) { {Name: "Compass of the Stars", SourceRefs: []source.SourceRef{{SourceID: doc.ID, StartUnitID: 20, EndUnitID: 20}}}, {Name: "Gold Pieces", SourceRefs: []source.SourceRef{{SourceID: doc.ID, StartUnitID: 30, EndUnitID: 30}}}, }} - client := &recordingNormalizerClient{response: `{"duplicate_groups":[{"members":["candidate-000001","candidate-000002"],"canonical":"candidate-000002"}]}`} + client := &recordingNormalizerClient{response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":2}]}`} result, err := newNormalizer(t, client).Normalize(context.Background(), normalizeRequestWithSource(input, doc)) if err != nil || result.Retry != nil || len(result.Value.Items) != 2 { t.Fatalf("Normalize() = %#v, %v", result, err) @@ -142,7 +149,7 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) { {Name: "Gold Piece", SourceRefs: ref(20)}, {Name: "Gold Pieces", SourceRefs: ref(30)}, }, - response: `{"duplicate_groups":[{"members":["candidate-000001","candidate-000002","candidate-000003"],"canonical":"candidate-000002"}]}`, + response: `{"duplicate_groups":[{"candidate_ids":[1,2,3],"canonical_candidate_id":2}]}`, wantNames: []string{"Gold Piece"}, wantRefCounts: []int{3}, warning: ReasonCodeDuplicateItemCollapsed, @@ -153,7 +160,7 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) { {Name: "Gold Pieces", SourceRefs: ref(10)}, {Name: "Silver Pieces", SourceRefs: ref(20)}, }, - response: `{"duplicate_groups":[{"members":["candidate-000001","candidate-000002"],"canonical":"candidate-000001"}]}`, + response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":1}]}`, wantNames: []string{"Gold Pieces", "Silver Pieces"}, wantRefCounts: []int{1, 1}, wantRetry: true, @@ -165,7 +172,7 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) { {Name: "Gold Pieces", SourceRefs: ref(10)}, {Name: "Longsword", SourceRefs: ref(20)}, }, - response: `{"duplicate_groups":[{"members":["candidate-000001","candidate-000002"],"canonical":"candidate-000001"}]}`, + response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":1}]}`, wantNames: []string{"Gold Pieces", "Longsword"}, wantRefCounts: []int{1, 1}, wantRetry: true, @@ -177,7 +184,7 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) { {Name: "Star Compass", SourceRefs: ref(10)}, {Name: "Compass of the Stars", SourceRefs: ref(20)}, }, - response: `{"duplicate_groups":[{"members":["candidate-000001","candidate-000002"],"canonical":"candidate-000002"}]}`, + response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":2}]}`, wantNames: []string{"Compass of the Stars"}, wantRefCounts: []int{2}, warning: ReasonCodeDuplicateItemCollapsed, @@ -188,7 +195,7 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) { {Name: "Gold Pieces", SourceRefs: ref(10)}, {Name: "Longsword", SourceRefs: ref(20)}, }, - response: `{"duplicate_groups":[{"members":["candidate-000001","candidate-000002"],"canonical":"candidate-000002"}]}`, + response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":2}]}`, wantNames: []string{"Gold Pieces", "Longsword"}, wantRefCounts: []int{1, 1}, wantRetry: true, @@ -231,13 +238,67 @@ func TestNormalizePreservesCandidatesForUnsafeProposalGroups(t *testing.T) { {Name: "Compass", SourceRefs: []source.SourceRef{{SourceID: doc.ID, StartUnitID: 20, EndUnitID: 20}}}, {Name: "Rope", SourceRefs: []source.SourceRef{{SourceID: doc.ID, StartUnitID: 30, EndUnitID: 30}}}, }} - client := &recordingNormalizerClient{response: `{"duplicate_groups":[{"members":["candidate-000001","candidate-000002"],"canonical":"candidate-000001"},{"members":["candidate-000002","candidate-000003"],"canonical":"candidate-000003"},{"members":["candidate-000001","candidate-000003"],"canonical":"candidate-000099"}]}`} + client := &recordingNormalizerClient{response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":1},{"candidate_ids":[2,3],"canonical_candidate_id":3},{"candidate_ids":[1,3],"canonical_candidate_id":99}]}`} result, err := newNormalizer(t, client).Normalize(context.Background(), normalizeRequestWithSource(input, doc)) if err != nil || result.Retry == nil || len(result.Value.Items) != 3 || !strings.Contains(result.Retry.Message, "overlapping_member") || !strings.Contains(result.Retry.Message, "canonical_unknown") || strings.Contains(result.Retry.Message, "Star Compass") || len(result.Retry.Message) > 4096 { t.Fatalf("Normalize() = %#v, %v; want deterministic retry fallback", result, err) } } +func TestNormalizeAppliesIndependentGroupAndCountsAllOmissions(t *testing.T) { + doc := semanticDocument() + ref := func(unitID int) []source.SourceRef { + return []source.SourceRef{{SourceID: doc.ID, StartUnitID: unitID, EndUnitID: unitID}} + } + input := dnd.ItemRegistry{Items: []dnd.Item{ + {Name: "Star Compass", SourceRefs: ref(10)}, + {Name: "Compass of the Stars", SourceRefs: ref(20)}, + {Name: "Gold Pieces", SourceRefs: ref(30)}, + {Name: "Silver Pieces", SourceRefs: ref(30)}, + {Name: "Rope", SourceRefs: ref(10)}, + }} + client := &recordingNormalizerClient{response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":2},{"candidate_ids":[3,4],"canonical_candidate_id":3},{"candidate_ids":[5,99],"canonical_candidate_id":5}]}`} + result, err := newNormalizer(t, client).Normalize(context.Background(), normalizeRequestWithSource(input, doc)) + if err != nil || result.Retry == nil { + t.Fatalf("Normalize() = %#v, %v; want retry with independently accepted output", result, err) + } + wantNames := []string{"Compass of the Stars", "Gold Pieces", "Silver Pieces", "Rope"} + if len(result.Value.Items) != len(wantNames) { + t.Fatalf("items = %#v, want %v", result.Value.Items, wantNames) + } + for index, name := range wantNames { + if result.Value.Items[index].Name != name { + t.Fatalf("item %d = %#v, want %q", index, result.Value.Items[index], name) + } + } + if !hasWarning(result.Warnings, ReasonCodeDuplicateItemCollapsed) || !hasWarning(result.Warnings, ReasonCodeItemSemanticProposalInvalid) { + t.Fatalf("warnings = %#v, want accepted and guarded-group diagnostics", result.Warnings) + } + if len(result.Retry.FallbackWarnings) != 1 || !strings.Contains(result.Retry.FallbackWarnings[0].Message, "2 proposal group(s)") { + t.Fatalf("retry = %#v, want one guarded and one malformed group counted", result.Retry) + } +} + +func TestNormalizeLimitSkipDoesNotCallLLMAndAddsBoundedFallbackWarning(t *testing.T) { + client := &recordingNormalizerClient{} + doc := &source.SourceDocument{ID: "session", Units: []source.SourceUnit{{ID: 1, Kind: "narration", Text: "A crowded storeroom"}}} + limit := semanticreconcile.DefaultLimits().MaximumCandidates + input := dnd.ItemRegistry{Items: make([]dnd.Item, limit+1)} + for index := range input.Items { + input.Items[index] = dnd.Item{Name: fmt.Sprintf("Item %d", index), SourceRefs: []source.SourceRef{{SourceID: doc.ID, StartUnitID: 1, EndUnitID: 1}}} + } + result, err := newNormalizer(t, client).Normalize(context.Background(), normalizeRequestWithSource(input, doc)) + if err != nil || result.Retry != nil { + t.Fatalf("Normalize() = %#v, %v; want deterministic limit fallback", result, err) + } + if len(client.requests) != 0 || len(result.Value.Items) != limit+1 { + t.Fatalf("completion calls = %d, items = %d; want no call and all records", len(client.requests), len(result.Value.Items)) + } + if !hasWarning(result.Warnings, ReasonCodeItemSemanticReconciliationExhausted) || len(result.Warnings) > diagnostics.MaxWarnings { + t.Fatalf("warnings = %#v, want bounded reconciliation fallback", result.Warnings) + } +} + func TestNormalizeRetryFallbackErrorsWarningsAndIdempotence(t *testing.T) { doc := semanticDocument() input := dnd.ItemRegistry{Items: []dnd.Item{{Name: "Star Compass", SourceRefs: []source.SourceRef{{SourceID: doc.ID, StartUnitID: 10, EndUnitID: 10}}}, {Name: "Compass", SourceRefs: []source.SourceRef{{SourceID: doc.ID, StartUnitID: 20, EndUnitID: 20}}}}} @@ -261,7 +322,7 @@ func TestNormalizeRetryFallbackErrorsWarningsAndIdempotence(t *testing.T) { func TestRegisterPromptAssetsPreparesItemNormalizationPrompt(t *testing.T) { registry := llm.NewAssetRegistry() - if err := entityreconcile.RegisterSchemaAssets(registry); err != nil { + if err := semanticreconcile.RegisterAssets(registry); err != nil { t.Fatal(err) } if err := RegisterPromptAssets(registry); err != nil { @@ -276,13 +337,16 @@ func TestRegisterPromptAssetsPreparesItemNormalizationPrompt(t *testing.T) { if err != nil { t.Fatal(err) } - prepared, err := engine.Prepare(context.Background(), promptkit.RunRequest{PromptID: PromptID, PromptVersion: entityreconcile.SchemaVersion, ProfileID: "item-normalize-test", Inputs: map[string]promptkit.ArtifactRef{"candidates": promptkit.Inline(`{"candidates":[{"name":"Rope","source_refs":[{"start_unit_id":1,"end_unit_id":1}]}]}`), "transcript": promptkit.Inline(`{"windows":[{"units":[]}]}`)}}) + prepared, err := engine.Prepare(context.Background(), promptkit.RunRequest{PromptID: PromptID, PromptVersion: PromptVersion, ProfileID: "item-normalize-test", Inputs: map[string]promptkit.ArtifactRef{"candidates": promptkit.Inline(`{"candidates":[{"candidate_id":1,"label":"Rope","source_refs":[{"start_unit_id":1,"end_unit_id":1}]}]}`), "transcript": promptkit.Inline(`{"windows":[{"units":[]}]}`)}}) if err != nil { t.Fatal(err) } - if prepared.OutputContract.SchemaPath != "dnd_entity_reconcile_llm.v1.json" || !strings.Contains(prepared.Messages[1].Content, "currency denominations") || !strings.Contains(prepared.Messages[1].Content, "materially different item types") { + if prepared.OutputContract.SchemaPath != "semantic_reconciliation_llm.v1.json" || !strings.Contains(prepared.Messages[1].Content, "candidate_id") || !strings.Contains(prepared.Messages[1].Content, "integer") || !strings.Contains(prepared.Messages[2].Content, "currency denominations") || !strings.Contains(prepared.Messages[2].Content, "materially different item") || prepared.Messages[2].CacheControl == nil || prepared.Messages[2].CacheControl.Type != promptkit.CacheControlEphemeral { t.Fatalf("prepared prompt = %#v", prepared) } + if prepared.Messages[4].CacheControl == nil || prepared.Messages[4].CacheControl.Type != promptkit.CacheControlEphemeral || !strings.Contains(prepared.Messages[3].Content, `"Rope"`) || strings.Contains(prepared.Messages[3].Content, `"windows"`) || !strings.Contains(prepared.Messages[4].Content, `"windows"`) || strings.Contains(prepared.Messages[4].Content, `"Rope"`) { + t.Fatalf("prepared prompt = %#v, want isolated candidate and transcript presentation", prepared) + } } type recordingNormalizerClient struct { @@ -300,53 +364,13 @@ func (c *recordingNormalizerClient) CompleteStructured(_ context.Context, reques if response == "" { response = `{"duplicate_groups":[]}` } - content, err := contextualProposalResponse(response, request.Inputs["candidates"].Content) - if err != nil { - return contracts.StructuredCompletionResponse{}, err - } + content := []byte(response) if err := json.Unmarshal(content, output); err != nil { return contracts.StructuredCompletionResponse{}, err } return contracts.StructuredCompletionResponse{Content: content}, nil } -func contextualProposalResponse(response string, candidateContent []byte) ([]byte, error) { - if !strings.Contains(response, "candidate-") { - return []byte(response), nil - } - var selection struct { - DuplicateGroups []struct { - Members []string `json:"members"` - Canonical string `json:"canonical"` - } `json:"duplicate_groups"` - } - if err := json.Unmarshal([]byte(response), &selection); err != nil { - return nil, err - } - var candidates struct { - Candidates []entityreconcile.Selector `json:"candidates"` - } - if err := json.Unmarshal(candidateContent, &candidates); err != nil { - return nil, err - } - selector := func(key string) entityreconcile.Selector { - index, err := strconv.Atoi(strings.TrimPrefix(key, "candidate-")) - if err != nil || index < 1 || index > len(candidates.Candidates) { - return entityreconcile.Selector{Name: key, SourceRefs: []entityreconcile.SourceRange{}} - } - return candidates.Candidates[index-1].Clone() - } - proposal := entityreconcile.ProposalResponse{DuplicateGroups: make([]entityreconcile.DuplicateGroup, len(selection.DuplicateGroups))} - for index, group := range selection.DuplicateGroups { - members := make([]entityreconcile.Selector, len(group.Members)) - for memberIndex, key := range group.Members { - members[memberIndex] = selector(key) - } - proposal.DuplicateGroups[index] = entityreconcile.DuplicateGroup{Members: members, Canonical: selector(group.Canonical)} - } - return json.Marshal(proposal) -} - func newNormalizer(t *testing.T, client contracts.StructuredLLMClient) *Normalizer { t.Helper() normalizer, err := New(client, Options{}) @@ -374,3 +398,12 @@ func hasWarning(warnings []contracts.Warning, reason string) bool { } return false } + +func hasFingerprint(fingerprints []pipeline.CheckpointFingerprint, name string) bool { + for _, fingerprint := range fingerprints { + if fingerprint.Name == name && fingerprint.Value != "" { + return true + } + } + return false +} diff --git a/internal/modules/dnd/normalize/itemregistry/prompt_assets.go b/internal/modules/dnd/normalize/itemregistry/prompt_assets.go index ff03c89..4623172 100644 --- a/internal/modules/dnd/normalize/itemregistry/prompt_assets.go +++ b/internal/modules/dnd/normalize/itemregistry/prompt_assets.go @@ -8,19 +8,26 @@ import ( rootassets "gitea.maximumdirect.net/eric/notarius/assets" "gitea.maximumdirect.net/eric/notarius/internal/framework/llm" "gitea.maximumdirect.net/eric/notarius/internal/framework/promptfs" + "gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared" ) const promptAssetRoot = "assets/prompts" -var promptAssetManifest = shared.PromptAssetManifest{ - ModuleDir: PromptID, - ModuleFiles: []promptfs.ModulePromptFile{ - {Name: "prompt.yaml", Path: "prompts/prompt.yaml"}, - {Name: "instructions.md", Path: "prompts/instructions.md"}, - {Name: "candidates.md", Path: "prompts/candidates.md"}, - }, - SharedFiles: []string{"common-dnd-system.md", "common-dnd-entity-reconciliation.md", "common-dnd-transcript-windows.md"}, +func promptAssetManifest() (shared.PromptAssetManifest, error) { + sharedFiles, err := semanticreconcile.SharedPromptFiles() + if err != nil { + return shared.PromptAssetManifest{}, fmt.Errorf("load shared semantic reconciliation prompt assets: %w", err) + } + return shared.PromptAssetManifest{ + ModuleDir: PromptID, + ModuleFiles: []promptfs.ModulePromptFile{ + {Name: "prompt.yaml", Path: "prompts/prompt.yaml"}, + {Name: "instructions.md", Path: "prompts/instructions.md"}, + }, + SharedFiles: []string{"common-dnd-system.md"}, + ExternalSharedFiles: sharedFiles, + }, nil } func moduleAssetFS() (fs.FS, error) { @@ -36,7 +43,11 @@ func RegisterPromptAssets(registry *llm.AssetRegistry) error { if err != nil { return err } - promptFS, err := promptAssetManifest.PromptFS(assets) + manifest, err := promptAssetManifest() + if err != nil { + return err + } + promptFS, err := manifest.PromptFS(assets) if err != nil { return fmt.Errorf("prepare item normalization prompt assets: %w", err) } @@ -50,7 +61,12 @@ func promptAssetMetadata() (string, error) { promptAssetHashErr = err return } - promptAssetHash, promptAssetHashErr = promptAssetManifest.Hash(assets) + manifest, err := promptAssetManifest() + if err != nil { + promptAssetHashErr = err + return + } + promptAssetHash, promptAssetHashErr = manifest.Hash(assets) }) return promptAssetHash, promptAssetHashErr } diff --git a/internal/modules/dnd/normalize/itemregistry/reconciliation.go b/internal/modules/dnd/normalize/itemregistry/reconciliation.go index df4cc50..5cbe6c3 100644 --- a/internal/modules/dnd/normalize/itemregistry/reconciliation.go +++ b/internal/modules/dnd/normalize/itemregistry/reconciliation.go @@ -2,104 +2,106 @@ package itemregistry import ( "fmt" + "sort" "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" + "gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile" + "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/items/identity" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/diagnostics" - "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/entityreconcile" ) -type safeReconciliationGroup struct { - members []int - canonical int -} +const incompatibleCurrencyGroup semanticreconcile.RejectionCategory = "incompatible_currency" -func reconciliationCandidates(records []normalizedRecord) []entityreconcile.Candidate { - candidates := make([]entityreconcile.Candidate, len(records)) +func reconciliationInputs(records []normalizedRecord) ([]semanticreconcile.Candidate, []semanticreconcile.Record[dnd.Item], error) { + candidates := make([]semanticreconcile.Candidate, len(records)) + envelopes := make([]semanticreconcile.Record[dnd.Item], len(records)) for index, record := range records { - candidates[index] = entityreconcile.Candidate{Name: record.item.Name, SourceRefs: cloneSourceRefs(record.item.SourceRefs)} + candidates[index] = semanticreconcile.Candidate{ + Label: record.item.Name, + SourceRefs: cloneSourceRefs(record.item.SourceRefs), + } + envelope, err := semanticreconcile.NewRecord(record.item, record.inputIndexes, record.earliest, cloneItem) + if err != nil { + return nil, nil, fmt.Errorf("record %d: %w", index, err) + } + envelopes[index] = envelope } - return candidates + return candidates, envelopes, nil } -func reconciliationGroups(assessment entityreconcile.Assessment, candidateKeys []string) []safeReconciliationGroup { - positions := make(map[string]int, len(candidateKeys)) - for index, key := range candidateKeys { - positions[key] = index - } - safeGroups := assessment.SafeGroups() - groups := make([]safeReconciliationGroup, 0, len(safeGroups)) - for _, group := range safeGroups { - members := group.Members() - memberPositions := make([]int, len(members)) - valid := true - for index, key := range members { - position, ok := positions[key] - if !ok { - valid = false - break +func applyReconciliationPlan(plan semanticreconcile.Plan, records []normalizedRecord, envelopes []semanticreconcile.Record[dnd.Item], order shared.SourceRefOrder) ([]normalizedRecord, []contracts.Warning, int, error) { + application, err := semanticreconcile.ApplyPlan(plan, envelopes, semanticreconcile.ApplicationPolicy[dnd.Item]{ + CloneValue: cloneItem, + RejectGroup: func(members []dnd.Item, _ dnd.Item) semanticreconcile.RejectionCategory { + if !canConsolidate(members) { + return incompatibleCurrencyGroup } - memberPositions[index] = position - } - canonical, ok := positions[group.Canonical()] - if valid && ok { - groups = append(groups, safeReconciliationGroup{members: memberPositions, canonical: canonical}) - } - } - return groups -} - -func reconciliationIssues(assessment entityreconcile.Assessment) []string { - issues := assessment.Issues() - details := make([]string, len(issues)) - for index, issue := range issues { - details[index] = fmt.Sprintf("group %d: %s", issue.GroupIndex, issue.Category) - } - return details -} - -func applySafeGroups(records []normalizedRecord, groups []safeReconciliationGroup, order shared.SourceRefOrder) ([]normalizedRecord, []contracts.Warning, int) { - byMember := make(map[int]safeReconciliationGroup, len(groups)*2) - for _, group := range groups { - for _, member := range group.members { - byMember[member] = group - } - } - output := make([]normalizedRecord, 0, len(records)-len(groups)) - warnings := make([]contracts.Warning, 0, len(groups)) - rejectedGroups := 0 - for index, record := range records { - group, grouped := byMember[index] - if !grouped { - output = append(output, cloneRecord(record)) - continue - } - if group.members[0] != index { - continue - } - if !canConsolidate(records, group) { - for _, member := range group.members { - output = append(output, cloneRecord(records[member])) + return "" + }, + ConsolidateGroup: func(members []dnd.Item, canonical dnd.Item) (dnd.Item, error) { + output := cloneItem(canonical) + output.SourceRefs = nil + for _, member := range members { + output.SourceRefs = append(output.SourceRefs, member.SourceRefs...) } - warnings = append(warnings, contracts.Warning{Scope: itemScope(records[group.members[0]].earliest), ReasonCode: ReasonCodeItemSemanticProposalInvalid, Message: "proposal group preserved because currency may only be consolidated with aliases of one denomination"}) - rejectedGroups++ - continue - } - consolidated := consolidateSemanticGroup(records, group, order) - output = append(output, consolidated) - warnings = append(warnings, semanticDuplicateWarning(consolidated, records[group.canonical])) + output.SourceRefs = order.Canonicalize(output.SourceRefs) + output.ID = identity.DeriveID(output.Name) + return output, nil + }, + }) + if err != nil { + return nil, nil, 0, err } - return output, warnings, rejectedGroups + + applied := application.Records() + output := make([]normalizedRecord, len(applied)) + for index, record := range applied { + output[index] = normalizedRecord{ + item: record.Value(), + inputIndexes: record.OriginalInputIndexes(), + earliest: record.EarliestInputPosition(), + } + } + type orderedWarning struct { + position int + warning contracts.Warning + } + orderedWarnings := make([]orderedWarning, 0, len(application.AppliedGroups())+len(application.RejectedGroups())) + for _, event := range application.AppliedGroups() { + provenance := event.Provenance() + orderedWarnings = append(orderedWarnings, orderedWarning{ + position: provenance.EarliestInputPosition(), + warning: semanticDuplicateWarning(provenance, records[provenance.CanonicalPosition()]), + }) + } + for _, event := range application.RejectedGroups() { + provenance := event.Provenance() + orderedWarnings = append(orderedWarnings, orderedWarning{ + position: provenance.EarliestInputPosition(), + warning: contracts.Warning{ + Scope: itemScope(provenance.EarliestInputPosition()), + ReasonCode: ReasonCodeItemSemanticProposalInvalid, + Message: "proposal group preserved because currency may only be consolidated with aliases of one denomination", + }, + }) + } + sort.SliceStable(orderedWarnings, func(left, right int) bool { return orderedWarnings[left].position < orderedWarnings[right].position }) + warnings := make([]contracts.Warning, len(orderedWarnings)) + for index, entry := range orderedWarnings { + warnings[index] = entry.warning + } + return output, warnings, len(application.RejectedGroups()), nil } -func canConsolidate(records []normalizedRecord, group safeReconciliationGroup) bool { +func canConsolidate(items []dnd.Item) bool { denomination := "" hasCurrency := false hasNonCurrency := false hasConflictingDenominations := false - for _, member := range group.members { - current := currencyDenomination(records[member].item.Name) + for _, item := range items { + current := currencyDenomination(item.Name) if current == "" { hasNonCurrency = true continue @@ -131,29 +133,26 @@ func currencyDenomination(name string) string { } } -func consolidateSemanticGroup(records []normalizedRecord, group safeReconciliationGroup, order shared.SourceRefOrder) normalizedRecord { - output := cloneRecord(records[group.members[0]]) - output.item.Name = records[group.canonical].item.Name - for _, member := range group.members[1:] { - output.item.SourceRefs = append(output.item.SourceRefs, records[member].item.SourceRefs...) - output.inputIndexes = append(output.inputIndexes, records[member].inputIndexes...) - if records[member].earliest < output.earliest { - output.earliest = records[member].earliest - } +func reconciliationIssues(issues []semanticreconcile.Issue) []string { + details := make([]string, len(issues)) + for index, issue := range issues { + details[index] = fmt.Sprintf("group %d: %s", issue.GroupIndex, issue.Category) } - output.inputIndexes = sortedUniqueIndexes(output.inputIndexes) - output.item.SourceRefs = order.Canonicalize(output.item.SourceRefs) - output.item.ID = identity.DeriveID(output.item.Name) - return output + return details } -func semanticDuplicateWarning(record normalizedRecord, canonical normalizedRecord) contracts.Warning { - details := make([]string, 0, len(record.inputIndexes)+1) - for _, inputIndex := range record.inputIndexes { +func semanticDuplicateWarning(provenance semanticreconcile.GroupProvenance, canonical normalizedRecord) contracts.Warning { + inputIndexes := provenance.OriginalInputIndexes() + details := make([]string, 0, len(inputIndexes)+1) + for _, inputIndex := range inputIndexes { details = append(details, fmt.Sprintf("input index %d", inputIndex)) } - if canonical.earliest != record.earliest { + if canonical.earliest != provenance.EarliestInputPosition() { details = append(details, fmt.Sprintf("canonical display name from input index %d", canonical.earliest)) } - return contracts.Warning{Scope: itemScope(record.earliest), ReasonCode: ReasonCodeDuplicateItemCollapsed, Message: diagnostics.Aggregate("semantic duplicate consolidation", details)} + return contracts.Warning{ + Scope: itemScope(provenance.EarliestInputPosition()), + ReasonCode: ReasonCodeDuplicateItemCollapsed, + Message: diagnostics.Aggregate("semantic duplicate consolidation", details), + } }