Classify registry normalization diagnostics
This commit is contained in:
@@ -404,7 +404,7 @@ and cannot increment process warning groups.
|
|||||||
This repetitive but cohesive migration is appropriately sized for one
|
This repetitive but cohesive migration is appropriately sized for one
|
||||||
`gpt-5.6-terra` prompt. Do not combine it with normalizer migration.
|
`gpt-5.6-terra` prompt. Do not combine it with normalizer migration.
|
||||||
|
|
||||||
## Stage 6 — Migrate D&D Registry Normalizers And Semantic Reconciliation
|
## Stage 6 — Migrate D&D Registry Normalizers And Semantic Reconciliation ✅
|
||||||
|
|
||||||
### Goal
|
### Goal
|
||||||
|
|
||||||
|
|||||||
@@ -31,8 +31,8 @@ const (
|
|||||||
ReasonCodeSourceReferencesNormalized = "source_references_normalized"
|
ReasonCodeSourceReferencesNormalized = "source_references_normalized"
|
||||||
ReasonCodeDuplicateItemCollapsed = "duplicate_item_collapsed"
|
ReasonCodeDuplicateItemCollapsed = "duplicate_item_collapsed"
|
||||||
ReasonCodeItemSemanticProposalInvalid = "item_semantic_proposal_invalid"
|
ReasonCodeItemSemanticProposalInvalid = "item_semantic_proposal_invalid"
|
||||||
|
ReasonCodeItemSemanticRetryProposalInvalid = "item_semantic_retry_proposal_invalid"
|
||||||
ReasonCodeItemSemanticReconciliationExhausted = "item_semantic_reconciliation_exhausted"
|
ReasonCodeItemSemanticReconciliationExhausted = "item_semantic_reconciliation_exhausted"
|
||||||
ReasonCodeItemNormalizationWarningsOmitted = "item_normalization_warnings_omitted"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
var requiredCapabilities = []string{"merged"}
|
var requiredCapabilities = []string{"merged"}
|
||||||
@@ -107,10 +107,10 @@ func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalize
|
|||||||
}
|
}
|
||||||
|
|
||||||
order := shared.NewSourceRefOrder(req.Source)
|
order := shared.NewSourceRefOrder(req.Source)
|
||||||
records, warnings := preprocessRecords(req.MergeOutput.Value, order)
|
records, findings := preprocessRecords(req.MergeOutput.Value, order)
|
||||||
deterministic := recordList(records)
|
deterministic := recordList(records)
|
||||||
if len(records) < 2 {
|
if len(records) < 2 {
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil
|
return normalizationResult(deterministic, findings, nil, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
candidates, envelopes, err := reconciliationInputs(records)
|
candidates, envelopes, err := reconciliationInputs(records)
|
||||||
@@ -127,48 +127,42 @@ func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalize
|
|||||||
|
|
||||||
switch reconciliation.Disposition() {
|
switch reconciliation.Disposition() {
|
||||||
case semanticreconcile.SkippedInsufficientCandidates:
|
case semanticreconcile.SkippedInsufficientCandidates:
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil
|
return normalizationResult(deterministic, findings, nil, nil)
|
||||||
case semanticreconcile.SkippedLimitExceeded:
|
case semanticreconcile.SkippedLimitExceeded:
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: deterministic, Warnings: limitWarningsWithSemanticFallback(warnings)}, nil
|
return fallbackResult(deterministic, findings, nil, semanticFallbackFinding(-1))
|
||||||
case semanticreconcile.RetryableInvalidStructuredOutput:
|
case semanticreconcile.RetryableInvalidStructuredOutput:
|
||||||
return n.invalidStructuredResult(deterministic, warnings), nil
|
return n.invalidStructuredResult(deterministic, findings)
|
||||||
case semanticreconcile.Complete, semanticreconcile.RetryableDiscardedProposalGroups:
|
case semanticreconcile.Complete, semanticreconcile.RetryableDiscardedProposalGroups:
|
||||||
default:
|
default:
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("unknown semantic reconciliation disposition %d", reconciliation.Disposition())
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("unknown semantic reconciliation disposition %d", reconciliation.Disposition())
|
||||||
}
|
}
|
||||||
|
|
||||||
applied, semanticWarnings, rejectedGroups, err := applyReconciliationPlan(reconciliation.Plan(), records, envelopes, order)
|
applied, semanticFindings, advisoryFindings, rejectedGroups, err := applyReconciliationPlan(reconciliation.Plan(), records, envelopes, order)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("apply semantic reconciliation plan: %w", err)
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("apply semantic reconciliation plan: %w", err)
|
||||||
}
|
}
|
||||||
warnings = append(warnings, semanticWarnings...)
|
findings = append(findings, semanticFindings...)
|
||||||
discardedGroups := reconciliation.DiscardedGroupCount() + rejectedGroups
|
discardedGroups := reconciliation.DiscardedGroupCount() + rejectedGroups
|
||||||
if discardedGroups == 0 {
|
if discardedGroups == 0 {
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: recordList(applied), Warnings: limitWarnings(warnings), ModelCandidate: reconciliation.ModelCandidate()}, nil
|
return normalizationResult(recordList(applied), findings, advisoryFindings, reconciliation.ModelCandidate())
|
||||||
}
|
}
|
||||||
return retryResult(recordList(applied), warnings, reconciliation, rejectedGroups), nil
|
return retryResult(recordList(applied), findings, advisoryFindings, reconciliation, rejectedGroups)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (n *Normalizer) invalidStructuredResult(value dnd.ItemRegistry, warnings []contracts.Warning) contracts.TypedNormalizeResult[dnd.ItemRegistry] {
|
func (n *Normalizer) invalidStructuredResult(value dnd.ItemRegistry, findings []contracts.Warning) (contracts.TypedNormalizeResult[dnd.ItemRegistry], error) {
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: value, Warnings: limitWarningsForRetry(warnings), Retry: &contracts.NormalizeRetry{
|
return retryResultWithFallback(value, findings, nil, nil, ReasonCodeItemSemanticRetryProposalInvalid, "semantic proposal requires retry: invalid structured output", semanticFallbackFinding(-1))
|
||||||
ReasonCode: ReasonCodeItemSemanticProposalInvalid, Message: "semantic proposal requires retry: invalid structured output",
|
|
||||||
FallbackWarnings: []contracts.Warning{semanticFallbackWarning(-1)},
|
|
||||||
}}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func retryResult(value dnd.ItemRegistry, warnings []contracts.Warning, reconciliation semanticreconcile.Result, rejectedGroups int) contracts.TypedNormalizeResult[dnd.ItemRegistry] {
|
func retryResult(value dnd.ItemRegistry, findings, advisoryFindings []contracts.Warning, reconciliation semanticreconcile.Result, rejectedGroups int) (contracts.TypedNormalizeResult[dnd.ItemRegistry], error) {
|
||||||
details := semanticreconcile.IssueDetails(reconciliation.Issues())
|
details := semanticreconcile.IssueDetails(reconciliation.Issues())
|
||||||
if rejectedGroups > 0 {
|
if rejectedGroups > 0 {
|
||||||
details = append(details, "currency may only be consolidated with aliases of one denomination")
|
details = append(details, "currency may only be consolidated with aliases of one denomination")
|
||||||
}
|
}
|
||||||
discardedGroups := reconciliation.DiscardedGroupCount() + rejectedGroups
|
discardedGroups := reconciliation.DiscardedGroupCount() + rejectedGroups
|
||||||
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: value, Warnings: limitWarningsForRetry(warnings), ModelCandidate: reconciliation.ModelCandidate(), Retry: &contracts.NormalizeRetry{
|
return retryResultWithFallback(value, findings, advisoryFindings, reconciliation.ModelCandidate(), ReasonCodeItemSemanticRetryProposalInvalid, diagnostics.Aggregate("semantic proposal requires retry", details), semanticFallbackFinding(discardedGroups))
|
||||||
ReasonCode: ReasonCodeItemSemanticProposalInvalid, Message: diagnostics.Aggregate("semantic proposal requires retry", details),
|
|
||||||
FallbackWarnings: []contracts.Warning{semanticFallbackWarning(discardedGroups)},
|
|
||||||
}}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func semanticFallbackWarning(discarded int) contracts.Warning {
|
func semanticFallbackFinding(discarded int) contracts.Warning {
|
||||||
message := "semantic proposal could not be applied"
|
message := "semantic proposal could not be applied"
|
||||||
if discarded >= 0 {
|
if discarded >= 0 {
|
||||||
message = fmt.Sprintf("%d proposal group(s) omitted after semantic proposal retry exhaustion", discarded)
|
message = fmt.Sprintf("%d proposal group(s) omitted after semantic proposal retry exhaustion", discarded)
|
||||||
@@ -176,24 +170,43 @@ func semanticFallbackWarning(discarded int) contracts.Warning {
|
|||||||
return contracts.Warning{Scope: "items", ReasonCode: ReasonCodeItemSemanticReconciliationExhausted, Message: message}
|
return contracts.Warning{Scope: "items", ReasonCode: ReasonCodeItemSemanticReconciliationExhausted, Message: message}
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarnings(warnings []contracts.Warning) []contracts.Warning {
|
func normalizationResult(value dnd.ItemRegistry, findings, advisoryFindings []contracts.Warning, candidate *contracts.ModelCandidate) (contracts.TypedNormalizeResult[dnd.ItemRegistry], error) {
|
||||||
return diagnostics.LimitWarnings(warnings, "items", ReasonCodeItemNormalizationWarningsOmitted)
|
diagnosticGroups, err := diagnostics.Collect(findings, contracts.DiagnosticDispositionObservation, contracts.DiagnosticCategoryNormalization)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("collect normalization diagnostics: %w", err)
|
||||||
|
}
|
||||||
|
advisoryGroups, err := diagnostics.Collect(advisoryFindings, contracts.DiagnosticDispositionAdvisory, contracts.DiagnosticCategoryDataQuality)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("collect data-quality diagnostics: %w", err)
|
||||||
|
}
|
||||||
|
diagnosticGroups = append(diagnosticGroups, advisoryGroups...)
|
||||||
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{Value: value, Diagnostics: diagnosticGroups, ModelCandidate: candidate}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarningsForRetry(warnings []contracts.Warning) []contracts.Warning {
|
func fallbackResult(value dnd.ItemRegistry, findings, advisoryFindings []contracts.Warning, fallback contracts.Warning) (contracts.TypedNormalizeResult[dnd.ItemRegistry], error) {
|
||||||
if warnings == nil {
|
result, err := normalizationResult(value, findings, advisoryFindings, nil)
|
||||||
return nil
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, err
|
||||||
}
|
}
|
||||||
if len(warnings) < diagnostics.MaxWarnings {
|
fallbackGroups, err := diagnostics.Collect([]contracts.Warning{fallback}, contracts.DiagnosticDispositionWarning, contracts.DiagnosticCategoryFallback)
|
||||||
return append([]contracts.Warning(nil), warnings...)
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("collect fallback diagnostic: %w", err)
|
||||||
}
|
}
|
||||||
displayed := diagnostics.MaxWarnings - 2
|
result.Diagnostics = append(result.Diagnostics, fallbackGroups...)
|
||||||
bounded := append([]contracts.Warning(nil), warnings[:displayed]...)
|
return result, nil
|
||||||
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 {
|
func retryResultWithFallback(value dnd.ItemRegistry, findings, advisoryFindings []contracts.Warning, candidate *contracts.ModelCandidate, reasonCode, message string, fallback contracts.Warning) (contracts.TypedNormalizeResult[dnd.ItemRegistry], error) {
|
||||||
return append(limitWarningsForRetry(warnings), semanticFallbackWarning(-1))
|
result, err := normalizationResult(value, findings, advisoryFindings, candidate)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, err
|
||||||
|
}
|
||||||
|
fallbackGroups, err := diagnostics.Collect([]contracts.Warning{fallback}, contracts.DiagnosticDispositionWarning, contracts.DiagnosticCategoryFallback)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.ItemRegistry]{}, normalizerErrorf("collect fallback diagnostic: %w", err)
|
||||||
|
}
|
||||||
|
result.Retry = &contracts.NormalizeRetry{ReasonCode: reasonCode, Message: message, FallbackDiagnostics: fallbackGroups}
|
||||||
|
return result, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
type normalizedRecord struct {
|
type normalizedRecord struct {
|
||||||
@@ -207,18 +220,18 @@ func preprocessRecords(input dnd.ItemRegistry, order shared.SourceRefOrder) ([]n
|
|||||||
return nil, nil
|
return nil, nil
|
||||||
}
|
}
|
||||||
records := make([]normalizedRecord, len(input.Items))
|
records := make([]normalizedRecord, len(input.Items))
|
||||||
warnings := make([]contracts.Warning, 0)
|
findings := make([]contracts.Warning, 0)
|
||||||
for index, inputItem := range input.Items {
|
for index, inputItem := range input.Items {
|
||||||
item, fieldsChanged, refsChanged := normalizeRecord(inputItem, order)
|
item, fieldsChanged, refsChanged := normalizeRecord(inputItem, order)
|
||||||
records[index] = normalizedRecord{item: item, inputIndexes: []int{index}, earliest: index}
|
records[index] = normalizedRecord{item: item, inputIndexes: []int{index}, earliest: index}
|
||||||
if fieldsChanged {
|
if fieldsChanged {
|
||||||
warnings = append(warnings, contracts.Warning{Scope: itemScope(index), ReasonCode: ReasonCodeItemFieldsNormalized, Message: fmt.Sprintf("input index %d: item name normalized for %s", index, diagnostics.Quote(inputItem.Name))})
|
findings = append(findings, contracts.Warning{Scope: itemScope(index), ReasonCode: ReasonCodeItemFieldsNormalized, Message: fmt.Sprintf("input index %d: item name normalized for %s", index, diagnostics.Quote(inputItem.Name))})
|
||||||
}
|
}
|
||||||
if refsChanged {
|
if refsChanged {
|
||||||
warnings = append(warnings, contracts.Warning{Scope: itemScope(index), ReasonCode: ReasonCodeSourceReferencesNormalized, Message: fmt.Sprintf("input index %d: source references normalized (original count %d, final count %d)", index, len(inputItem.SourceRefs), len(item.SourceRefs))})
|
findings = append(findings, contracts.Warning{Scope: itemScope(index), ReasonCode: ReasonCodeSourceReferencesNormalized, Message: fmt.Sprintf("input index %d: source references normalized (original count %d, final count %d)", index, len(inputItem.SourceRefs), len(item.SourceRefs))})
|
||||||
}
|
}
|
||||||
if inputItem.ID != item.ID {
|
if inputItem.ID != item.ID {
|
||||||
warnings = append(warnings, contracts.Warning{Scope: itemScope(index), ReasonCode: ReasonCodeItemIDRecomputed, Message: fmt.Sprintf("input index %d: item ID recomputed from %s", index, diagnostics.Quote(item.Name))})
|
findings = append(findings, contracts.Warning{Scope: itemScope(index), ReasonCode: ReasonCodeItemIDRecomputed, Message: fmt.Sprintf("input index %d: item ID recomputed from %s", index, diagnostics.Quote(item.Name))})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
groups := comparisonNameGroups(records)
|
groups := comparisonNameGroups(records)
|
||||||
@@ -234,10 +247,10 @@ func preprocessRecords(input dnd.ItemRegistry, order shared.SourceRefOrder) ([]n
|
|||||||
retained.inputIndexes = sortedUniqueIndexes(retained.inputIndexes)
|
retained.inputIndexes = sortedUniqueIndexes(retained.inputIndexes)
|
||||||
output = append(output, retained)
|
output = append(output, retained)
|
||||||
if len(members) > 1 {
|
if len(members) > 1 {
|
||||||
warnings = append(warnings, duplicateWarning(retained.earliest, memberInputIndexes(records, members[1:])))
|
findings = append(findings, duplicateWarning(retained.earliest, memberInputIndexes(records, members[1:])))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return output, warnings
|
return output, findings
|
||||||
}
|
}
|
||||||
|
|
||||||
func normalizeRecord(input dnd.Item, order shared.SourceRefOrder) (dnd.Item, bool, bool) {
|
func normalizeRecord(input dnd.Item, order shared.SourceRefOrder) (dnd.Item, bool, bool) {
|
||||||
|
|||||||
@@ -18,7 +18,6 @@ import (
|
|||||||
"gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile"
|
"gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd"
|
"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/items/identity"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/diagnostics"
|
|
||||||
identityvalidator "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/validate/itemregistry/identity"
|
identityvalidator "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/validate/itemregistry/identity"
|
||||||
"gitea.maximumdirect.net/eric/promptkit"
|
"gitea.maximumdirect.net/eric/promptkit"
|
||||||
)
|
)
|
||||||
@@ -81,8 +80,8 @@ func TestNormalizeConsolidatesEqualNamesAcrossEvidenceWithoutMutation(t *testing
|
|||||||
}
|
}
|
||||||
rope := result.Value.Items[0]
|
rope := result.Value.Items[0]
|
||||||
wantRefs := []source.SourceRef{{SourceID: "session", StartUnitID: 1, EndUnitID: 1}, {SourceID: "session", StartUnitID: 2, EndUnitID: 2}, {SourceID: "session", StartUnitID: 3, EndUnitID: 3}}
|
wantRefs := []source.SourceRef{{SourceID: "session", StartUnitID: 1, EndUnitID: 1}, {SourceID: "session", StartUnitID: 2, EndUnitID: 2}, {SourceID: "session", StartUnitID: 3, EndUnitID: 3}}
|
||||||
if rope.Name != "Rope" || rope.ID != identity.DeriveID("Rope") || !reflect.DeepEqual(rope.SourceRefs, wantRefs) || result.Value.Items[1].Name != "Lantern" || !hasWarning(result.Warnings, ReasonCodeDuplicateItemCollapsed) {
|
if rope.Name != "Rope" || rope.ID != identity.DeriveID("Rope") || !reflect.DeepEqual(rope.SourceRefs, wantRefs) || result.Value.Items[1].Name != "Lantern" || !hasDiagnostic(result.Diagnostics, ReasonCodeDuplicateItemCollapsed, contracts.DiagnosticDispositionObservation) {
|
||||||
t.Fatalf("items = %#v, warnings = %#v; want earliest display name, canonical evidence union, and stable placement", result.Value.Items, result.Warnings)
|
t.Fatalf("items = %#v, diagnostics = %#v; want earliest display name, canonical evidence union, and stable placement", result.Value.Items, result.Diagnostics)
|
||||||
}
|
}
|
||||||
validation, validationErr := identityvalidator.New(identityvalidator.Options{}).Validate(context.Background(), contracts.TypedValidationRequest[dnd.ItemRegistry]{Value: result.Value})
|
validation, validationErr := identityvalidator.New(identityvalidator.Options{}).Validate(context.Background(), contracts.TypedValidationRequest[dnd.ItemRegistry]{Value: result.Value})
|
||||||
if validationErr != nil || !validation.Approved {
|
if validationErr != nil || !validation.Approved {
|
||||||
@@ -129,8 +128,8 @@ func TestNormalizeAppliesSafeAliasProposal(t *testing.T) {
|
|||||||
t.Fatalf("Normalize() = %#v, %v", result, err)
|
t.Fatalf("Normalize() = %#v, %v", result, err)
|
||||||
}
|
}
|
||||||
merged := result.Value.Items[0]
|
merged := result.Value.Items[0]
|
||||||
if merged.Name != "Compass of the Stars" || merged.ID != identity.DeriveID(merged.Name) || len(merged.SourceRefs) != 2 || !hasWarning(result.Warnings, ReasonCodeDuplicateItemCollapsed) {
|
if merged.Name != "Compass of the Stars" || merged.ID != identity.DeriveID(merged.Name) || len(merged.SourceRefs) != 2 || !hasDiagnostic(result.Diagnostics, ReasonCodeDuplicateItemCollapsed, contracts.DiagnosticDispositionObservation) {
|
||||||
t.Fatalf("merged item = %#v, warnings = %#v", merged, result.Warnings)
|
t.Fatalf("merged item = %#v, diagnostics = %#v", merged, result.Diagnostics)
|
||||||
}
|
}
|
||||||
encoded := string(client.requests[0].Inputs["candidates"].Content) + string(client.requests[0].Inputs["transcript"].Content)
|
encoded := string(client.requests[0].Inputs["candidates"].Content) + string(client.requests[0].Inputs["transcript"].Content)
|
||||||
if strings.Contains(encoded, doc.ID) || strings.Contains(encoded, "candidate-") || strings.Contains(encoded, merged.ID) || !strings.Contains(encoded, `"source_refs"`) {
|
if strings.Contains(encoded, doc.ID) || strings.Contains(encoded, "candidate-") || strings.Contains(encoded, merged.ID) || !strings.Contains(encoded, `"source_refs"`) {
|
||||||
@@ -150,7 +149,8 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
wantNames []string
|
wantNames []string
|
||||||
wantRefCounts []int
|
wantRefCounts []int
|
||||||
wantRetry bool
|
wantRetry bool
|
||||||
warning string
|
reasonCode string
|
||||||
|
disposition contracts.DiagnosticDisposition
|
||||||
}{
|
}{
|
||||||
{
|
{
|
||||||
name: "same denomination aliases",
|
name: "same denomination aliases",
|
||||||
@@ -162,7 +162,8 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
response: `{"duplicate_groups":[{"candidate_ids":[1,2,3],"canonical_candidate_id":2}]}`,
|
response: `{"duplicate_groups":[{"candidate_ids":[1,2,3],"canonical_candidate_id":2}]}`,
|
||||||
wantNames: []string{"Gold Piece"},
|
wantNames: []string{"Gold Piece"},
|
||||||
wantRefCounts: []int{3},
|
wantRefCounts: []int{3},
|
||||||
warning: ReasonCodeDuplicateItemCollapsed,
|
reasonCode: ReasonCodeDuplicateItemCollapsed,
|
||||||
|
disposition: contracts.DiagnosticDispositionObservation,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "different denominations",
|
name: "different denominations",
|
||||||
@@ -174,7 +175,8 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
wantNames: []string{"Gold Pieces", "Silver Pieces"},
|
wantNames: []string{"Gold Pieces", "Silver Pieces"},
|
||||||
wantRefCounts: []int{1, 1},
|
wantRefCounts: []int{1, 1},
|
||||||
wantRetry: true,
|
wantRetry: true,
|
||||||
warning: ReasonCodeItemSemanticProposalInvalid,
|
reasonCode: ReasonCodeItemSemanticProposalInvalid,
|
||||||
|
disposition: contracts.DiagnosticDispositionAdvisory,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "currency plus ordinary item",
|
name: "currency plus ordinary item",
|
||||||
@@ -186,7 +188,8 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
wantNames: []string{"Gold Pieces", "Longsword"},
|
wantNames: []string{"Gold Pieces", "Longsword"},
|
||||||
wantRefCounts: []int{1, 1},
|
wantRefCounts: []int{1, 1},
|
||||||
wantRetry: true,
|
wantRetry: true,
|
||||||
warning: ReasonCodeItemSemanticProposalInvalid,
|
reasonCode: ReasonCodeItemSemanticProposalInvalid,
|
||||||
|
disposition: contracts.DiagnosticDispositionAdvisory,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "ordinary items",
|
name: "ordinary items",
|
||||||
@@ -197,7 +200,8 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":2}]}`,
|
response: `{"duplicate_groups":[{"candidate_ids":[1,2],"canonical_candidate_id":2}]}`,
|
||||||
wantNames: []string{"Compass of the Stars"},
|
wantNames: []string{"Compass of the Stars"},
|
||||||
wantRefCounts: []int{2},
|
wantRefCounts: []int{2},
|
||||||
warning: ReasonCodeDuplicateItemCollapsed,
|
reasonCode: ReasonCodeDuplicateItemCollapsed,
|
||||||
|
disposition: contracts.DiagnosticDispositionObservation,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "currency plus ordinary canonical item",
|
name: "currency plus ordinary canonical item",
|
||||||
@@ -209,7 +213,8 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
wantNames: []string{"Gold Pieces", "Longsword"},
|
wantNames: []string{"Gold Pieces", "Longsword"},
|
||||||
wantRefCounts: []int{1, 1},
|
wantRefCounts: []int{1, 1},
|
||||||
wantRetry: true,
|
wantRetry: true,
|
||||||
warning: ReasonCodeItemSemanticProposalInvalid,
|
reasonCode: ReasonCodeItemSemanticProposalInvalid,
|
||||||
|
disposition: contracts.DiagnosticDispositionAdvisory,
|
||||||
},
|
},
|
||||||
} {
|
} {
|
||||||
t.Run(test.name, func(t *testing.T) {
|
t.Run(test.name, func(t *testing.T) {
|
||||||
@@ -223,7 +228,7 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
result, err := newNormalizer(t, &recordingNormalizerClient{response: test.response}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
result, err := newNormalizer(t, &recordingNormalizerClient{response: test.response}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
||||||
if err != nil || (result.Retry != nil) != test.wantRetry || !reflect.DeepEqual(input, before) || !hasWarning(result.Warnings, test.warning) {
|
if err != nil || (result.Retry != nil) != test.wantRetry || !reflect.DeepEqual(input, before) || !hasDiagnostic(result.Diagnostics, test.reasonCode, test.disposition) {
|
||||||
t.Fatalf("Normalize() = %#v, %v", result, err)
|
t.Fatalf("Normalize() = %#v, %v", result, err)
|
||||||
}
|
}
|
||||||
if len(result.Value.Items) != len(test.wantNames) {
|
if len(result.Value.Items) != len(test.wantNames) {
|
||||||
@@ -234,7 +239,7 @@ func TestNormalizeAppliesCurrencyReconciliationSafely(t *testing.T) {
|
|||||||
t.Fatalf("item %d = %#v, want name %q with %d source refs", index, item, test.wantNames[index], test.wantRefCounts[index])
|
t.Fatalf("item %d = %#v, want name %q with %d source refs", index, item, test.wantNames[index], test.wantRefCounts[index])
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if test.wantRetry && (len(result.Retry.FallbackWarnings) != 1 || result.Retry.FallbackWarnings[0].ReasonCode != ReasonCodeItemSemanticReconciliationExhausted) {
|
if test.wantRetry && (len(result.Retry.FallbackDiagnostics) != 1 || result.Retry.FallbackDiagnostics[0].ReasonCode != ReasonCodeItemSemanticReconciliationExhausted || result.Retry.ReasonCode != ReasonCodeItemSemanticRetryProposalInvalid) {
|
||||||
t.Fatalf("retry = %#v, want preserved-group fallback", result.Retry)
|
t.Fatalf("retry = %#v, want preserved-group fallback", result.Retry)
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
@@ -281,10 +286,10 @@ func TestNormalizeAppliesIndependentGroupAndCountsAllOmissions(t *testing.T) {
|
|||||||
t.Fatalf("item %d = %#v, want %q", index, result.Value.Items[index], name)
|
t.Fatalf("item %d = %#v, want %q", index, result.Value.Items[index], name)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if !hasWarning(result.Warnings, ReasonCodeDuplicateItemCollapsed) || !hasWarning(result.Warnings, ReasonCodeItemSemanticProposalInvalid) {
|
if !hasDiagnostic(result.Diagnostics, ReasonCodeDuplicateItemCollapsed, contracts.DiagnosticDispositionObservation) || !hasDiagnostic(result.Diagnostics, ReasonCodeItemSemanticProposalInvalid, contracts.DiagnosticDispositionAdvisory) {
|
||||||
t.Fatalf("warnings = %#v, want accepted and guarded-group diagnostics", result.Warnings)
|
t.Fatalf("diagnostics = %#v, want accepted and guarded-group diagnostics", result.Diagnostics)
|
||||||
}
|
}
|
||||||
if len(result.Retry.FallbackWarnings) != 1 || !strings.Contains(result.Retry.FallbackWarnings[0].Message, "2 proposal group(s)") {
|
if len(result.Retry.FallbackDiagnostics) != 1 || !strings.Contains(result.Retry.FallbackDiagnostics[0].Samples[0].Message, "2 proposal group(s)") {
|
||||||
t.Fatalf("retry = %#v, want one guarded and one malformed group counted", result.Retry)
|
t.Fatalf("retry = %#v, want one guarded and one malformed group counted", result.Retry)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -304,8 +309,8 @@ func TestNormalizeLimitSkipDoesNotCallLLMAndAddsBoundedFallbackWarning(t *testin
|
|||||||
if len(client.requests) != 0 || len(result.Value.Items) != limit+1 {
|
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))
|
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 {
|
if !hasDiagnostic(result.Diagnostics, ReasonCodeItemSemanticReconciliationExhausted, contracts.DiagnosticDispositionWarning) {
|
||||||
t.Fatalf("warnings = %#v, want bounded reconciliation fallback", result.Warnings)
|
t.Fatalf("diagnostics = %#v, want reconciliation fallback warning", result.Diagnostics)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -313,16 +318,13 @@ func TestNormalizeRetryFallbackErrorsWarningsAndIdempotence(t *testing.T) {
|
|||||||
doc := semanticDocument()
|
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}}}}}
|
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}}}}}
|
||||||
invalid, err := newNormalizer(t, &recordingNormalizerClient{err: contracts.ErrInvalidStructuredOutput}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
invalid, err := newNormalizer(t, &recordingNormalizerClient{err: contracts.ErrInvalidStructuredOutput}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
||||||
if err != nil || invalid.Retry == nil || invalid.Retry.ReasonCode != ReasonCodeItemSemanticProposalInvalid {
|
if err != nil || invalid.Retry == nil || invalid.Retry.ReasonCode != ReasonCodeItemSemanticRetryProposalInvalid {
|
||||||
t.Fatalf("invalid result = %#v, %v", invalid, err)
|
t.Fatalf("invalid result = %#v, %v", invalid, err)
|
||||||
}
|
}
|
||||||
_, err = newNormalizer(t, &recordingNormalizerClient{err: errors.New("provider unavailable")}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
_, err = newNormalizer(t, &recordingNormalizerClient{err: errors.New("provider unavailable")}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
||||||
if err == nil || !strings.Contains(err.Error(), "provider unavailable") {
|
if err == nil || !strings.Contains(err.Error(), "provider unavailable") {
|
||||||
t.Fatalf("provider error = %v", err)
|
t.Fatalf("provider error = %v", err)
|
||||||
}
|
}
|
||||||
if bounded := limitWarningsForRetry(make([]contracts.Warning, diagnostics.MaxWarnings+5)); len(bounded) != diagnostics.MaxWarnings-1 || bounded[len(bounded)-1].ReasonCode != ReasonCodeItemNormalizationWarningsOmitted {
|
|
||||||
t.Fatalf("retry warning limit = %#v", bounded)
|
|
||||||
}
|
|
||||||
first, err := newNormalizer(t, &recordingNormalizerClient{}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
first, err := newNormalizer(t, &recordingNormalizerClient{}).Normalize(context.Background(), normalizeRequestWithSource(input, doc))
|
||||||
second, secondErr := newNormalizer(t, &recordingNormalizerClient{}).Normalize(context.Background(), normalizeRequestWithSource(first.Value, doc))
|
second, secondErr := newNormalizer(t, &recordingNormalizerClient{}).Normalize(context.Background(), normalizeRequestWithSource(first.Value, doc))
|
||||||
if err != nil || secondErr != nil || !reflect.DeepEqual(first.Value, second.Value) {
|
if err != nil || secondErr != nil || !reflect.DeepEqual(first.Value, second.Value) {
|
||||||
@@ -407,9 +409,9 @@ func normalizeRequestWithSource(value dnd.ItemRegistry, doc *source.SourceDocume
|
|||||||
func semanticDocument() *source.SourceDocument {
|
func semanticDocument() *source.SourceDocument {
|
||||||
return &source.SourceDocument{ID: "item-session", Units: []source.SourceUnit{{ID: 10, Text: "The Star Compass points north."}, {ID: 20, Text: "The compass of the stars glows."}, {ID: 30, Text: "The chest holds gold pieces."}}}
|
return &source.SourceDocument{ID: "item-session", Units: []source.SourceUnit{{ID: 10, Text: "The Star Compass points north."}, {ID: 20, Text: "The compass of the stars glows."}, {ID: 30, Text: "The chest holds gold pieces."}}}
|
||||||
}
|
}
|
||||||
func hasWarning(warnings []contracts.Warning, reason string) bool {
|
func hasDiagnostic(diagnostics []contracts.ProducerDiagnostic, reason string, disposition contracts.DiagnosticDisposition) bool {
|
||||||
for _, warning := range warnings {
|
for _, diagnostic := range diagnostics {
|
||||||
if warning.ReasonCode == reason {
|
if diagnostic.ReasonCode == reason && diagnostic.Disposition == disposition {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -31,7 +31,7 @@ func reconciliationInputs(records []normalizedRecord) ([]semanticreconcile.Candi
|
|||||||
return candidates, envelopes, nil
|
return candidates, envelopes, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func applyReconciliationPlan(plan semanticreconcile.Plan, records []normalizedRecord, envelopes []semanticreconcile.Record[dnd.Item], order shared.SourceRefOrder) ([]normalizedRecord, []contracts.Warning, int, error) {
|
func applyReconciliationPlan(plan semanticreconcile.Plan, records []normalizedRecord, envelopes []semanticreconcile.Record[dnd.Item], order shared.SourceRefOrder) ([]normalizedRecord, []contracts.Warning, []contracts.Warning, int, error) {
|
||||||
application, err := semanticreconcile.ApplyPlan(plan, envelopes, semanticreconcile.ApplicationPolicy[dnd.Item]{
|
application, err := semanticreconcile.ApplyPlan(plan, envelopes, semanticreconcile.ApplicationPolicy[dnd.Item]{
|
||||||
CloneValue: cloneItem,
|
CloneValue: cloneItem,
|
||||||
RejectGroup: func(members []dnd.Item, _ dnd.Item) semanticreconcile.RejectionCategory {
|
RejectGroup: func(members []dnd.Item, _ dnd.Item) semanticreconcile.RejectionCategory {
|
||||||
@@ -52,7 +52,7 @@ func applyReconciliationPlan(plan semanticreconcile.Plan, records []normalizedRe
|
|||||||
},
|
},
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, nil, 0, err
|
return nil, nil, nil, 0, err
|
||||||
}
|
}
|
||||||
|
|
||||||
applied := application.Records()
|
applied := application.Records()
|
||||||
@@ -64,35 +64,41 @@ func applyReconciliationPlan(plan semanticreconcile.Plan, records []normalizedRe
|
|||||||
earliest: record.EarliestInputPosition(),
|
earliest: record.EarliestInputPosition(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
type orderedWarning struct {
|
type orderedFinding struct {
|
||||||
position int
|
position int
|
||||||
warning contracts.Warning
|
finding contracts.Warning
|
||||||
}
|
}
|
||||||
orderedWarnings := make([]orderedWarning, 0, len(application.AppliedGroups())+len(application.RejectedGroups()))
|
observations := make([]orderedFinding, 0, len(application.AppliedGroups()))
|
||||||
|
advisories := make([]orderedFinding, 0, len(application.RejectedGroups()))
|
||||||
for _, event := range application.AppliedGroups() {
|
for _, event := range application.AppliedGroups() {
|
||||||
provenance := event.Provenance()
|
provenance := event.Provenance()
|
||||||
orderedWarnings = append(orderedWarnings, orderedWarning{
|
observations = append(observations, orderedFinding{
|
||||||
position: provenance.EarliestInputPosition(),
|
position: provenance.EarliestInputPosition(),
|
||||||
warning: semanticDuplicateWarning(provenance, records[provenance.CanonicalPosition()]),
|
finding: semanticDuplicateFinding(provenance, records[provenance.CanonicalPosition()]),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
for _, event := range application.RejectedGroups() {
|
for _, event := range application.RejectedGroups() {
|
||||||
provenance := event.Provenance()
|
provenance := event.Provenance()
|
||||||
orderedWarnings = append(orderedWarnings, orderedWarning{
|
advisories = append(advisories, orderedFinding{
|
||||||
position: provenance.EarliestInputPosition(),
|
position: provenance.EarliestInputPosition(),
|
||||||
warning: contracts.Warning{
|
finding: contracts.Warning{
|
||||||
Scope: itemScope(provenance.EarliestInputPosition()),
|
Scope: itemScope(provenance.EarliestInputPosition()),
|
||||||
ReasonCode: ReasonCodeItemSemanticProposalInvalid,
|
ReasonCode: ReasonCodeItemSemanticProposalInvalid,
|
||||||
Message: "proposal group preserved because currency may only be consolidated with aliases of one denomination",
|
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 })
|
sort.SliceStable(observations, func(left, right int) bool { return observations[left].position < observations[right].position })
|
||||||
warnings := make([]contracts.Warning, len(orderedWarnings))
|
sort.SliceStable(advisories, func(left, right int) bool { return advisories[left].position < advisories[right].position })
|
||||||
for index, entry := range orderedWarnings {
|
observationFindings := make([]contracts.Warning, len(observations))
|
||||||
warnings[index] = entry.warning
|
for index, entry := range observations {
|
||||||
|
observationFindings[index] = entry.finding
|
||||||
}
|
}
|
||||||
return output, warnings, len(application.RejectedGroups()), nil
|
advisoryFindings := make([]contracts.Warning, len(advisories))
|
||||||
|
for index, entry := range advisories {
|
||||||
|
advisoryFindings[index] = entry.finding
|
||||||
|
}
|
||||||
|
return output, observationFindings, advisoryFindings, len(application.RejectedGroups()), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func canConsolidate(items []dnd.Item) bool {
|
func canConsolidate(items []dnd.Item) bool {
|
||||||
@@ -133,7 +139,7 @@ func currencyDenomination(name string) string {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func semanticDuplicateWarning(provenance semanticreconcile.GroupProvenance, canonical normalizedRecord) contracts.Warning {
|
func semanticDuplicateFinding(provenance semanticreconcile.GroupProvenance, canonical normalizedRecord) contracts.Warning {
|
||||||
inputIndexes := provenance.OriginalInputIndexes()
|
inputIndexes := provenance.OriginalInputIndexes()
|
||||||
details := make([]string, 0, len(inputIndexes)+1)
|
details := make([]string, 0, len(inputIndexes)+1)
|
||||||
for _, inputIndex := range inputIndexes {
|
for _, inputIndex := range inputIndexes {
|
||||||
|
|||||||
@@ -33,7 +33,6 @@ const (
|
|||||||
ReasonCodeDuplicateLocationCollapsed = "duplicate_location_collapsed"
|
ReasonCodeDuplicateLocationCollapsed = "duplicate_location_collapsed"
|
||||||
ReasonCodeLocationSemanticProposalInvalid = "location_semantic_proposal_invalid"
|
ReasonCodeLocationSemanticProposalInvalid = "location_semantic_proposal_invalid"
|
||||||
ReasonCodeLocationSemanticReconciliationExhausted = "location_semantic_reconciliation_exhausted"
|
ReasonCodeLocationSemanticReconciliationExhausted = "location_semantic_reconciliation_exhausted"
|
||||||
ReasonCodeLocationNormalizationWarningsOmitted = "location_normalization_warnings_omitted"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
var requiredCapabilities = []string{"merged"}
|
var requiredCapabilities = []string{"merged"}
|
||||||
@@ -108,10 +107,10 @@ func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalize
|
|||||||
}
|
}
|
||||||
|
|
||||||
order := shared.NewSourceRefOrder(req.Source)
|
order := shared.NewSourceRefOrder(req.Source)
|
||||||
records, warnings := preprocessRecords(req.MergeOutput.Value, order)
|
records, findings := preprocessRecords(req.MergeOutput.Value, order)
|
||||||
deterministic := recordList(records)
|
deterministic := recordList(records)
|
||||||
if len(records) < 2 {
|
if len(records) < 2 {
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil
|
return normalizationResult(deterministic, findings, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
candidates, envelopes, err := reconciliationInputs(records)
|
candidates, envelopes, err := reconciliationInputs(records)
|
||||||
@@ -128,42 +127,36 @@ func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalize
|
|||||||
|
|
||||||
switch reconciliation.Disposition() {
|
switch reconciliation.Disposition() {
|
||||||
case semanticreconcile.SkippedInsufficientCandidates:
|
case semanticreconcile.SkippedInsufficientCandidates:
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil
|
return normalizationResult(deterministic, findings, nil)
|
||||||
case semanticreconcile.SkippedLimitExceeded:
|
case semanticreconcile.SkippedLimitExceeded:
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{Value: deterministic, Warnings: limitWarningsWithSemanticFallback(warnings)}, nil
|
return fallbackResult(deterministic, findings, semanticFallbackFinding(-1))
|
||||||
case semanticreconcile.RetryableInvalidStructuredOutput:
|
case semanticreconcile.RetryableInvalidStructuredOutput:
|
||||||
return n.invalidStructuredResult(deterministic, warnings), nil
|
return n.invalidStructuredResult(deterministic, findings)
|
||||||
case semanticreconcile.Complete, semanticreconcile.RetryableDiscardedProposalGroups:
|
case semanticreconcile.Complete, semanticreconcile.RetryableDiscardedProposalGroups:
|
||||||
default:
|
default:
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, normalizerErrorf("unknown semantic reconciliation disposition %d", reconciliation.Disposition())
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, normalizerErrorf("unknown semantic reconciliation disposition %d", reconciliation.Disposition())
|
||||||
}
|
}
|
||||||
|
|
||||||
applied, semanticWarnings, err := applyReconciliationPlan(reconciliation.Plan(), records, envelopes, order)
|
applied, semanticFindings, err := applyReconciliationPlan(reconciliation.Plan(), records, envelopes, order)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, normalizerErrorf("apply semantic reconciliation plan: %w", err)
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, normalizerErrorf("apply semantic reconciliation plan: %w", err)
|
||||||
}
|
}
|
||||||
warnings = append(warnings, semanticWarnings...)
|
findings = append(findings, semanticFindings...)
|
||||||
if reconciliation.Disposition() == semanticreconcile.Complete {
|
if reconciliation.Disposition() == semanticreconcile.Complete {
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{Value: recordList(applied), Warnings: limitWarnings(warnings), ModelCandidate: reconciliation.ModelCandidate()}, nil
|
return normalizationResult(recordList(applied), findings, reconciliation.ModelCandidate())
|
||||||
}
|
}
|
||||||
return retryResult(recordList(applied), warnings, reconciliation), nil
|
return retryResult(recordList(applied), findings, reconciliation)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (n *Normalizer) invalidStructuredResult(value dnd.LocationRegistry, warnings []contracts.Warning) contracts.TypedNormalizeResult[dnd.LocationRegistry] {
|
func (n *Normalizer) invalidStructuredResult(value dnd.LocationRegistry, findings []contracts.Warning) (contracts.TypedNormalizeResult[dnd.LocationRegistry], error) {
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{Value: value, Warnings: limitWarningsForRetry(warnings), Retry: &contracts.NormalizeRetry{
|
return retryResultWithFallback(value, findings, nil, "semantic proposal requires retry: invalid structured output", semanticFallbackFinding(-1))
|
||||||
ReasonCode: ReasonCodeLocationSemanticProposalInvalid, Message: "semantic proposal requires retry: invalid structured output",
|
|
||||||
FallbackWarnings: []contracts.Warning{semanticFallbackWarning(-1)},
|
|
||||||
}}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func retryResult(value dnd.LocationRegistry, warnings []contracts.Warning, reconciliation semanticreconcile.Result) contracts.TypedNormalizeResult[dnd.LocationRegistry] {
|
func retryResult(value dnd.LocationRegistry, findings []contracts.Warning, reconciliation semanticreconcile.Result) (contracts.TypedNormalizeResult[dnd.LocationRegistry], error) {
|
||||||
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{Value: value, Warnings: limitWarningsForRetry(warnings), ModelCandidate: reconciliation.ModelCandidate(), Retry: &contracts.NormalizeRetry{
|
return retryResultWithFallback(value, findings, reconciliation.ModelCandidate(), diagnostics.Aggregate("semantic proposal requires retry", semanticreconcile.IssueDetails(reconciliation.Issues())), semanticFallbackFinding(reconciliation.DiscardedGroupCount()))
|
||||||
ReasonCode: ReasonCodeLocationSemanticProposalInvalid, Message: diagnostics.Aggregate("semantic proposal requires retry", semanticreconcile.IssueDetails(reconciliation.Issues())),
|
|
||||||
FallbackWarnings: []contracts.Warning{semanticFallbackWarning(reconciliation.DiscardedGroupCount())},
|
|
||||||
}}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func semanticFallbackWarning(discarded int) contracts.Warning {
|
func semanticFallbackFinding(discarded int) contracts.Warning {
|
||||||
message := "semantic proposal could not be applied"
|
message := "semantic proposal could not be applied"
|
||||||
if discarded >= 0 {
|
if discarded >= 0 {
|
||||||
message = fmt.Sprintf("%d proposal group(s) omitted after semantic proposal retry exhaustion", discarded)
|
message = fmt.Sprintf("%d proposal group(s) omitted after semantic proposal retry exhaustion", discarded)
|
||||||
@@ -171,24 +164,38 @@ func semanticFallbackWarning(discarded int) contracts.Warning {
|
|||||||
return contracts.Warning{Scope: "locations", ReasonCode: ReasonCodeLocationSemanticReconciliationExhausted, Message: message}
|
return contracts.Warning{Scope: "locations", ReasonCode: ReasonCodeLocationSemanticReconciliationExhausted, Message: message}
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarnings(warnings []contracts.Warning) []contracts.Warning {
|
func normalizationResult(value dnd.LocationRegistry, findings []contracts.Warning, candidate *contracts.ModelCandidate) (contracts.TypedNormalizeResult[dnd.LocationRegistry], error) {
|
||||||
return diagnostics.LimitWarnings(warnings, "locations", ReasonCodeLocationNormalizationWarningsOmitted)
|
diagnosticGroups, err := diagnostics.Collect(findings, contracts.DiagnosticDispositionObservation, contracts.DiagnosticCategoryNormalization)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, normalizerErrorf("collect normalization diagnostics: %w", err)
|
||||||
|
}
|
||||||
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{Value: value, Diagnostics: diagnosticGroups, ModelCandidate: candidate}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarningsForRetry(warnings []contracts.Warning) []contracts.Warning {
|
func fallbackResult(value dnd.LocationRegistry, findings []contracts.Warning, fallback contracts.Warning) (contracts.TypedNormalizeResult[dnd.LocationRegistry], error) {
|
||||||
if warnings == nil {
|
result, err := normalizationResult(value, findings, nil)
|
||||||
return nil
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, err
|
||||||
}
|
}
|
||||||
if len(warnings) < diagnostics.MaxWarnings {
|
fallbackGroups, err := diagnostics.Collect([]contracts.Warning{fallback}, contracts.DiagnosticDispositionWarning, contracts.DiagnosticCategoryFallback)
|
||||||
return append([]contracts.Warning(nil), warnings...)
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, normalizerErrorf("collect fallback diagnostic: %w", err)
|
||||||
}
|
}
|
||||||
displayed := diagnostics.MaxWarnings - 2
|
result.Diagnostics = append(result.Diagnostics, fallbackGroups...)
|
||||||
bounded := append([]contracts.Warning(nil), warnings[:displayed]...)
|
return result, nil
|
||||||
return append(bounded, contracts.Warning{Scope: "locations", ReasonCode: ReasonCodeLocationNormalizationWarningsOmitted, Message: fmt.Sprintf("%d additional warning(s) omitted", len(warnings)-displayed)})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarningsWithSemanticFallback(warnings []contracts.Warning) []contracts.Warning {
|
func retryResultWithFallback(value dnd.LocationRegistry, findings []contracts.Warning, candidate *contracts.ModelCandidate, message string, fallback contracts.Warning) (contracts.TypedNormalizeResult[dnd.LocationRegistry], error) {
|
||||||
return append(limitWarningsForRetry(warnings), semanticFallbackWarning(-1))
|
result, err := normalizationResult(value, findings, candidate)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, err
|
||||||
|
}
|
||||||
|
fallbackGroups, err := diagnostics.Collect([]contracts.Warning{fallback}, contracts.DiagnosticDispositionWarning, contracts.DiagnosticCategoryFallback)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.LocationRegistry]{}, normalizerErrorf("collect fallback diagnostic: %w", err)
|
||||||
|
}
|
||||||
|
result.Retry = &contracts.NormalizeRetry{ReasonCode: ReasonCodeLocationSemanticProposalInvalid, Message: message, FallbackDiagnostics: fallbackGroups}
|
||||||
|
return result, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
type normalizedRecord struct {
|
type normalizedRecord struct {
|
||||||
|
|||||||
@@ -16,7 +16,6 @@ import (
|
|||||||
"gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile"
|
"gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd"
|
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/locations/identity"
|
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/locations/identity"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/diagnostics"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestModuleContractAndMetadata(t *testing.T) {
|
func TestModuleContractAndMetadata(t *testing.T) {
|
||||||
@@ -73,7 +72,7 @@ func TestNormalizePreparesOnlyExactDuplicatesAndRetainsSameNameAndNestedPlaces(t
|
|||||||
if got := []string{result.Value.Locations[0].Name, result.Value.Locations[1].Name, result.Value.Locations[2].Name}; !reflect.DeepEqual(got, []string{"The Tavern", "The Tavern", "The Tavern Cellar"}) {
|
if got := []string{result.Value.Locations[0].Name, result.Value.Locations[1].Name, result.Value.Locations[2].Name}; !reflect.DeepEqual(got, []string{"The Tavern", "The Tavern", "The Tavern Cellar"}) {
|
||||||
t.Fatalf("locations = %#v, want same names and nested place retained", got)
|
t.Fatalf("locations = %#v, want same names and nested place retained", got)
|
||||||
}
|
}
|
||||||
if result.Value.Locations[0].ID == result.Value.Locations[1].ID || !hasWarning(result.Warnings, ReasonCodeDuplicateLocationCollapsed) || len(client.requests) != 0 {
|
if result.Value.Locations[0].ID == result.Value.Locations[1].ID || !hasDiagnostic(result.Diagnostics, ReasonCodeDuplicateLocationCollapsed, contracts.DiagnosticDispositionObservation) || len(client.requests) != 0 {
|
||||||
t.Fatalf("result = %#v, want evidence-anchored IDs and exact duplicate warning", result)
|
t.Fatalf("result = %#v, want evidence-anchored IDs and exact duplicate warning", result)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -108,8 +107,8 @@ func TestNormalizeAppliesSafeAliasGroupAndUsesContextualInputs(t *testing.T) {
|
|||||||
t.Fatalf("Normalize() = %#v, %v", result, err)
|
t.Fatalf("Normalize() = %#v, %v", result, err)
|
||||||
}
|
}
|
||||||
merged := result.Value.Locations[0]
|
merged := result.Value.Locations[0]
|
||||||
if merged.Name != "the Greencloak's refuge" || merged.ID != identity.DeriveID(merged.Name, merged.SourceRefs) || len(merged.SourceRefs) != 2 || !hasWarning(result.Warnings, ReasonCodeDuplicateLocationCollapsed) {
|
if merged.Name != "the Greencloak's refuge" || merged.ID != identity.DeriveID(merged.Name, merged.SourceRefs) || len(merged.SourceRefs) != 2 || !hasDiagnostic(result.Diagnostics, ReasonCodeDuplicateLocationCollapsed, contracts.DiagnosticDispositionObservation) {
|
||||||
t.Fatalf("merged location = %#v, warnings = %#v", merged, result.Warnings)
|
t.Fatalf("merged location = %#v, diagnostics = %#v", merged, result.Diagnostics)
|
||||||
}
|
}
|
||||||
encoded := string(client.requests[0].Inputs["candidates"].Content) + string(client.requests[0].Inputs["transcript"].Content)
|
encoded := string(client.requests[0].Inputs["candidates"].Content) + string(client.requests[0].Inputs["transcript"].Content)
|
||||||
if strings.Contains(encoded, doc.ID) || strings.Contains(encoded, "candidate-") || strings.Contains(encoded, merged.ID) || !strings.Contains(encoded, `"source_refs"`) {
|
if strings.Contains(encoded, doc.ID) || strings.Contains(encoded, "candidate-") || strings.Contains(encoded, merged.ID) || !strings.Contains(encoded, `"source_refs"`) {
|
||||||
@@ -129,8 +128,8 @@ func TestNormalizeRejectsUnsafeAndOverlappingGroupsWithoutLosingCandidates(t *te
|
|||||||
if err != nil || result.Retry == nil || len(result.Value.Locations) != 3 || !strings.Contains(result.Retry.Message, "overlapping_member") {
|
if err != nil || result.Retry == nil || len(result.Value.Locations) != 3 || !strings.Contains(result.Retry.Message, "overlapping_member") {
|
||||||
t.Fatalf("Normalize() = %#v, %v; want safe retry fallback", result, err)
|
t.Fatalf("Normalize() = %#v, %v; want safe retry fallback", result, err)
|
||||||
}
|
}
|
||||||
if len(result.Retry.FallbackWarnings) != 1 || !strings.Contains(result.Retry.FallbackWarnings[0].Message, "2 proposal group") {
|
if len(result.Retry.FallbackDiagnostics) != 1 || !strings.Contains(result.Retry.FallbackDiagnostics[0].Samples[0].Message, "2 proposal group") {
|
||||||
t.Fatalf("fallback warnings = %#v", result.Retry.FallbackWarnings)
|
t.Fatalf("fallback diagnostics = %#v", result.Retry.FallbackDiagnostics)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -177,8 +176,8 @@ func TestNormalizeLimitSkipDoesNotCallLLMAndAddsBoundedFallbackWarning(t *testin
|
|||||||
if len(client.requests) != 0 || len(result.Value.Locations) != limit+1 {
|
if len(client.requests) != 0 || len(result.Value.Locations) != limit+1 {
|
||||||
t.Fatalf("completion calls = %d, locations = %d; want no call and all records", len(client.requests), len(result.Value.Locations))
|
t.Fatalf("completion calls = %d, locations = %d; want no call and all records", len(client.requests), len(result.Value.Locations))
|
||||||
}
|
}
|
||||||
if !hasWarning(result.Warnings, ReasonCodeLocationSemanticReconciliationExhausted) || len(result.Warnings) > diagnostics.MaxWarnings {
|
if !hasDiagnostic(result.Diagnostics, ReasonCodeLocationSemanticReconciliationExhausted, contracts.DiagnosticDispositionWarning) {
|
||||||
t.Fatalf("warnings = %#v, want bounded reconciliation fallback", result.Warnings)
|
t.Fatalf("diagnostics = %#v, want reconciliation fallback warning", result.Diagnostics)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -193,11 +192,6 @@ func TestNormalizeHandlesRetryFallbackAndErrors(t *testing.T) {
|
|||||||
if err == nil || !strings.Contains(err.Error(), "provider unavailable") {
|
if err == nil || !strings.Contains(err.Error(), "provider unavailable") {
|
||||||
t.Fatalf("provider error = %v", err)
|
t.Fatalf("provider error = %v", err)
|
||||||
}
|
}
|
||||||
warnings := make([]contracts.Warning, 25)
|
|
||||||
bounded := limitWarningsForRetry(warnings)
|
|
||||||
if len(bounded) != 19 || bounded[len(bounded)-1].ReasonCode != ReasonCodeLocationNormalizationWarningsOmitted {
|
|
||||||
t.Fatalf("retry warning limit = %#v", bounded)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestNormalizeOrdersEvidenceAndIsIdempotent(t *testing.T) {
|
func TestNormalizeOrdersEvidenceAndIsIdempotent(t *testing.T) {
|
||||||
|
|||||||
@@ -59,9 +59,9 @@ func normalizeRequestWithSource(value dnd.LocationRegistry, doc *source.SourceDo
|
|||||||
func semanticDocument() *source.SourceDocument {
|
func semanticDocument() *source.SourceDocument {
|
||||||
return &source.SourceDocument{ID: "location-session", Units: []source.SourceUnit{{ID: 10, Kind: "speech", Text: "The old mill is the Greencloak's refuge."}, {ID: 20, Kind: "speech", Text: "The mill stands on the northern road."}, {ID: 30, Kind: "speech", Text: "The tavern is beside the mill."}, {ID: 40, Kind: "speech", Text: "The mill's cellar is flooded."}}}
|
return &source.SourceDocument{ID: "location-session", Units: []source.SourceUnit{{ID: 10, Kind: "speech", Text: "The old mill is the Greencloak's refuge."}, {ID: 20, Kind: "speech", Text: "The mill stands on the northern road."}, {ID: 30, Kind: "speech", Text: "The tavern is beside the mill."}, {ID: 40, Kind: "speech", Text: "The mill's cellar is flooded."}}}
|
||||||
}
|
}
|
||||||
func hasWarning(warnings []contracts.Warning, reason string) bool {
|
func hasDiagnostic(diagnostics []contracts.ProducerDiagnostic, reason string, disposition contracts.DiagnosticDisposition) bool {
|
||||||
for _, warning := range warnings {
|
for _, diagnostic := range diagnostics {
|
||||||
if warning.ReasonCode == reason {
|
if diagnostic.ReasonCode == reason && diagnostic.Disposition == disposition {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -32,7 +32,6 @@ const (
|
|||||||
ReasonCodeDuplicateNPCCollapsed = "duplicate_npc_collapsed"
|
ReasonCodeDuplicateNPCCollapsed = "duplicate_npc_collapsed"
|
||||||
ReasonCodeNPCSemanticProposalInvalid = "npc_semantic_proposal_invalid"
|
ReasonCodeNPCSemanticProposalInvalid = "npc_semantic_proposal_invalid"
|
||||||
ReasonCodeNPCSemanticReconciliationExhausted = "npc_semantic_reconciliation_exhausted"
|
ReasonCodeNPCSemanticReconciliationExhausted = "npc_semantic_reconciliation_exhausted"
|
||||||
ReasonCodeNPCNormalizationWarningsOmitted = "npc_normalization_warnings_omitted"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
var requiredCapabilities = []string{"merged"}
|
var requiredCapabilities = []string{"merged"}
|
||||||
@@ -107,10 +106,10 @@ func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalize
|
|||||||
}
|
}
|
||||||
|
|
||||||
order := shared.NewSourceRefOrder(req.Source)
|
order := shared.NewSourceRefOrder(req.Source)
|
||||||
records, warnings := preprocessRecords(req.MergeOutput.Value, order)
|
records, findings := preprocessRecords(req.MergeOutput.Value, order)
|
||||||
deterministic := recordList(records)
|
deterministic := recordList(records)
|
||||||
if len(records) < 2 {
|
if len(records) < 2 {
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil
|
return normalizationResult(deterministic, findings, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
candidates, envelopes, err := reconciliationInputs(records)
|
candidates, envelopes, err := reconciliationInputs(records)
|
||||||
@@ -127,53 +126,36 @@ func (n *Normalizer) Normalize(ctx context.Context, req contracts.TypedNormalize
|
|||||||
|
|
||||||
switch reconciliation.Disposition() {
|
switch reconciliation.Disposition() {
|
||||||
case semanticreconcile.SkippedInsufficientCandidates:
|
case semanticreconcile.SkippedInsufficientCandidates:
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{Value: deterministic, Warnings: limitWarnings(warnings)}, nil
|
return normalizationResult(deterministic, findings, nil)
|
||||||
case semanticreconcile.SkippedLimitExceeded:
|
case semanticreconcile.SkippedLimitExceeded:
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{Value: deterministic, Warnings: limitWarningsWithSemanticFallback(warnings)}, nil
|
return fallbackResult(deterministic, findings, semanticFallbackFinding(-1))
|
||||||
case semanticreconcile.RetryableInvalidStructuredOutput:
|
case semanticreconcile.RetryableInvalidStructuredOutput:
|
||||||
return n.invalidStructuredResult(deterministic, warnings), nil
|
return n.invalidStructuredResult(deterministic, findings)
|
||||||
case semanticreconcile.Complete, semanticreconcile.RetryableDiscardedProposalGroups:
|
case semanticreconcile.Complete, semanticreconcile.RetryableDiscardedProposalGroups:
|
||||||
default:
|
default:
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, normalizerErrorf("unknown semantic reconciliation disposition %d", reconciliation.Disposition())
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, normalizerErrorf("unknown semantic reconciliation disposition %d", reconciliation.Disposition())
|
||||||
}
|
}
|
||||||
|
|
||||||
applied, semanticWarnings, err := applyReconciliationPlan(reconciliation.Plan(), records, envelopes, order)
|
applied, semanticFindings, err := applyReconciliationPlan(reconciliation.Plan(), records, envelopes, order)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, normalizerErrorf("apply semantic reconciliation plan: %w", err)
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, normalizerErrorf("apply semantic reconciliation plan: %w", err)
|
||||||
}
|
}
|
||||||
warnings = append(warnings, semanticWarnings...)
|
findings = append(findings, semanticFindings...)
|
||||||
if reconciliation.Disposition() == semanticreconcile.Complete {
|
if reconciliation.Disposition() == semanticreconcile.Complete {
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{Value: recordList(applied), Warnings: limitWarnings(warnings), ModelCandidate: reconciliation.ModelCandidate()}, nil
|
return normalizationResult(recordList(applied), findings, reconciliation.ModelCandidate())
|
||||||
}
|
}
|
||||||
return retryResult(recordList(applied), warnings, reconciliation), nil
|
return retryResult(recordList(applied), findings, reconciliation)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (n *Normalizer) invalidStructuredResult(value dnd.NPCRegistry, warnings []contracts.Warning) contracts.TypedNormalizeResult[dnd.NPCRegistry] {
|
func (n *Normalizer) invalidStructuredResult(value dnd.NPCRegistry, findings []contracts.Warning) (contracts.TypedNormalizeResult[dnd.NPCRegistry], error) {
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{
|
return retryResultWithFallback(value, findings, nil, "semantic proposal requires retry: invalid structured output", semanticFallbackFinding(-1))
|
||||||
Value: value,
|
|
||||||
Warnings: limitWarningsForRetry(warnings),
|
|
||||||
Retry: &contracts.NormalizeRetry{
|
|
||||||
ReasonCode: ReasonCodeNPCSemanticProposalInvalid,
|
|
||||||
Message: "semantic proposal requires retry: invalid structured output",
|
|
||||||
FallbackWarnings: []contracts.Warning{semanticFallbackWarning(-1)},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func retryResult(value dnd.NPCRegistry, warnings []contracts.Warning, reconciliation semanticreconcile.Result) contracts.TypedNormalizeResult[dnd.NPCRegistry] {
|
func retryResult(value dnd.NPCRegistry, findings []contracts.Warning, reconciliation semanticreconcile.Result) (contracts.TypedNormalizeResult[dnd.NPCRegistry], error) {
|
||||||
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{
|
return retryResultWithFallback(value, findings, reconciliation.ModelCandidate(), diagnostics.Aggregate("semantic proposal requires retry", semanticreconcile.IssueDetails(reconciliation.Issues())), semanticFallbackFinding(reconciliation.DiscardedGroupCount()))
|
||||||
Value: value,
|
|
||||||
Warnings: limitWarningsForRetry(warnings),
|
|
||||||
ModelCandidate: reconciliation.ModelCandidate(),
|
|
||||||
Retry: &contracts.NormalizeRetry{
|
|
||||||
ReasonCode: ReasonCodeNPCSemanticProposalInvalid,
|
|
||||||
Message: diagnostics.Aggregate("semantic proposal requires retry", semanticreconcile.IssueDetails(reconciliation.Issues())),
|
|
||||||
FallbackWarnings: []contracts.Warning{semanticFallbackWarning(reconciliation.DiscardedGroupCount())},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func semanticFallbackWarning(discardedGroups int) contracts.Warning {
|
func semanticFallbackFinding(discardedGroups int) contracts.Warning {
|
||||||
message := "semantic proposal could not be applied"
|
message := "semantic proposal could not be applied"
|
||||||
if discardedGroups >= 0 {
|
if discardedGroups >= 0 {
|
||||||
message = fmt.Sprintf("%d proposal group(s) omitted after semantic proposal retry exhaustion", discardedGroups)
|
message = fmt.Sprintf("%d proposal group(s) omitted after semantic proposal retry exhaustion", discardedGroups)
|
||||||
@@ -181,27 +163,38 @@ func semanticFallbackWarning(discardedGroups int) contracts.Warning {
|
|||||||
return contracts.Warning{Scope: "npcs", ReasonCode: ReasonCodeNPCSemanticReconciliationExhausted, Message: message}
|
return contracts.Warning{Scope: "npcs", ReasonCode: ReasonCodeNPCSemanticReconciliationExhausted, Message: message}
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarnings(warnings []contracts.Warning) []contracts.Warning {
|
func normalizationResult(value dnd.NPCRegistry, findings []contracts.Warning, candidate *contracts.ModelCandidate) (contracts.TypedNormalizeResult[dnd.NPCRegistry], error) {
|
||||||
return diagnostics.LimitWarnings(warnings, "npcs", ReasonCodeNPCNormalizationWarningsOmitted)
|
diagnosticGroups, err := diagnostics.Collect(findings, contracts.DiagnosticDispositionObservation, contracts.DiagnosticCategoryNormalization)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, normalizerErrorf("collect normalization diagnostics: %w", err)
|
||||||
|
}
|
||||||
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{Value: value, Diagnostics: diagnosticGroups, ModelCandidate: candidate}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarningsForRetry(warnings []contracts.Warning) []contracts.Warning {
|
func fallbackResult(value dnd.NPCRegistry, findings []contracts.Warning, fallback contracts.Warning) (contracts.TypedNormalizeResult[dnd.NPCRegistry], error) {
|
||||||
if warnings == nil {
|
result, err := normalizationResult(value, findings, nil)
|
||||||
return nil
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, err
|
||||||
}
|
}
|
||||||
if len(warnings) < diagnostics.MaxWarnings {
|
fallbackGroups, err := diagnostics.Collect([]contracts.Warning{fallback}, contracts.DiagnosticDispositionWarning, contracts.DiagnosticCategoryFallback)
|
||||||
return append([]contracts.Warning(nil), warnings...)
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, normalizerErrorf("collect fallback diagnostic: %w", err)
|
||||||
}
|
}
|
||||||
displayed := diagnostics.MaxWarnings - 2
|
result.Diagnostics = append(result.Diagnostics, fallbackGroups...)
|
||||||
bounded := append([]contracts.Warning(nil), warnings[:displayed]...)
|
return result, nil
|
||||||
return append(bounded, contracts.Warning{
|
|
||||||
Scope: "npcs", ReasonCode: ReasonCodeNPCNormalizationWarningsOmitted,
|
|
||||||
Message: fmt.Sprintf("%d additional warning(s) omitted", len(warnings)-displayed),
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func limitWarningsWithSemanticFallback(warnings []contracts.Warning) []contracts.Warning {
|
func retryResultWithFallback(value dnd.NPCRegistry, findings []contracts.Warning, candidate *contracts.ModelCandidate, message string, fallback contracts.Warning) (contracts.TypedNormalizeResult[dnd.NPCRegistry], error) {
|
||||||
return append(limitWarningsForRetry(warnings), semanticFallbackWarning(-1))
|
result, err := normalizationResult(value, findings, candidate)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, err
|
||||||
|
}
|
||||||
|
fallbackGroups, err := diagnostics.Collect([]contracts.Warning{fallback}, contracts.DiagnosticDispositionWarning, contracts.DiagnosticCategoryFallback)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.TypedNormalizeResult[dnd.NPCRegistry]{}, normalizerErrorf("collect fallback diagnostic: %w", err)
|
||||||
|
}
|
||||||
|
result.Retry = &contracts.NormalizeRetry{ReasonCode: ReasonCodeNPCSemanticProposalInvalid, Message: message, FallbackDiagnostics: fallbackGroups}
|
||||||
|
return result, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
type normalizedRecord struct {
|
type normalizedRecord struct {
|
||||||
|
|||||||
@@ -65,8 +65,8 @@ func TestNormalizeNamesEvidenceAndIDs(t *testing.T) {
|
|||||||
t.Fatalf("NPC = %#v, want %#v", result.Value.NPCs[0], want)
|
t.Fatalf("NPC = %#v, want %#v", result.Value.NPCs[0], want)
|
||||||
}
|
}
|
||||||
for _, reason := range []string{ReasonCodeNPCFieldsNormalized, ReasonCodeSourceReferencesNormalized, ReasonCodeNPCIDRecomputed} {
|
for _, reason := range []string{ReasonCodeNPCFieldsNormalized, ReasonCodeSourceReferencesNormalized, ReasonCodeNPCIDRecomputed} {
|
||||||
if !hasWarning(result.Warnings, reason, "npcs[0]") {
|
if !hasDiagnostic(result.Diagnostics, reason, "npcs[0]", contracts.DiagnosticDispositionObservation) {
|
||||||
t.Fatalf("warnings = %#v, want %s", result.Warnings, reason)
|
t.Fatalf("diagnostics = %#v, want %s", result.Diagnostics, reason)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -104,8 +104,8 @@ func TestNormalizeConsolidatesCanonicalNamesOnlyAndUnionsEvidence(t *testing.T)
|
|||||||
if refs := result.Value.NPCs[0].SourceRefs; len(refs) != 2 || refs[0].SourceID != "a" || refs[1].SourceID != "b" {
|
if refs := result.Value.NPCs[0].SourceRefs; len(refs) != 2 || refs[0].SourceID != "a" || refs[1].SourceID != "b" {
|
||||||
t.Fatalf("source refs = %#v, want evidence union", refs)
|
t.Fatalf("source refs = %#v, want evidence union", refs)
|
||||||
}
|
}
|
||||||
if !hasWarning(result.Warnings, ReasonCodeDuplicateNPCCollapsed, "npcs[0]") {
|
if !hasDiagnostic(result.Diagnostics, ReasonCodeDuplicateNPCCollapsed, "npcs[0]", contracts.DiagnosticDispositionObservation) {
|
||||||
t.Fatalf("warnings = %#v, want duplicate collapse", result.Warnings)
|
t.Fatalf("diagnostics = %#v, want duplicate collapse", result.Diagnostics)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -197,12 +197,17 @@ func normalizeRequestWithSource(value dnd.NPCRegistry, doc *source.SourceDocumen
|
|||||||
return request
|
return request
|
||||||
}
|
}
|
||||||
|
|
||||||
func hasWarning(warnings []contracts.Warning, reason, scope string) bool {
|
func hasDiagnostic(diagnostics []contracts.ProducerDiagnostic, reason, scope string, disposition contracts.DiagnosticDisposition) bool {
|
||||||
for _, warning := range warnings {
|
for _, diagnostic := range diagnostics {
|
||||||
if warning.ReasonCode == reason && warning.Scope == scope {
|
if diagnostic.ReasonCode != reason || diagnostic.Disposition != disposition {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, sample := range diagnostic.Samples {
|
||||||
|
if sample.Scope == scope {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -15,7 +15,6 @@ import (
|
|||||||
"gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile"
|
"gitea.maximumdirect.net/eric/notarius/internal/framework/semanticreconcile"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd"
|
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/npcs/identity"
|
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/npcs/identity"
|
||||||
"gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared/diagnostics"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestNormalizeSkipsSemanticCompletionWithoutTwoEligibleCandidates(t *testing.T) {
|
func TestNormalizeSkipsSemanticCompletionWithoutTwoEligibleCandidates(t *testing.T) {
|
||||||
@@ -67,8 +66,8 @@ func TestNormalizeAppliesSafeProposalAndUsesPrivateInputs(t *testing.T) {
|
|||||||
if merged.ID != identity.DeriveID("Mira Thorn") || !reflect.DeepEqual(merged.SourceRefs, []source.SourceRef{{SourceID: doc.ID, StartUnitID: 10, EndUnitID: 10}, {SourceID: doc.ID, StartUnitID: 20, EndUnitID: 20}}) {
|
if merged.ID != identity.DeriveID("Mira Thorn") || !reflect.DeepEqual(merged.SourceRefs, []source.SourceRef{{SourceID: doc.ID, StartUnitID: 10, EndUnitID: 10}, {SourceID: doc.ID, StartUnitID: 20, EndUnitID: 20}}) {
|
||||||
t.Fatalf("merged NPC = %#v, want canonical ID and original evidence union", merged)
|
t.Fatalf("merged NPC = %#v, want canonical ID and original evidence union", merged)
|
||||||
}
|
}
|
||||||
if !hasWarning(result.Warnings, ReasonCodeDuplicateNPCCollapsed, "npcs[0]") {
|
if !hasDiagnostic(result.Diagnostics, ReasonCodeDuplicateNPCCollapsed, "npcs[0]", contracts.DiagnosticDispositionObservation) {
|
||||||
t.Fatalf("warnings = %#v, want semantic collapse warning", result.Warnings)
|
t.Fatalf("diagnostics = %#v, want semantic collapse observation", result.Diagnostics)
|
||||||
}
|
}
|
||||||
if len(client.requests) != 1 {
|
if len(client.requests) != 1 {
|
||||||
t.Fatalf("completion calls = %d, want one", len(client.requests))
|
t.Fatalf("completion calls = %d, want one", len(client.requests))
|
||||||
@@ -119,8 +118,8 @@ func TestNormalizeUnsafeProposalReturnsSafeRetryFallback(t *testing.T) {
|
|||||||
if len(result.Value.NPCs) != 2 || result.Value.NPCs[0].Name != "Mira Thorn" || result.Value.NPCs[1].Name != "Captain Vale" {
|
if len(result.Value.NPCs) != 2 || result.Value.NPCs[0].Name != "Mira Thorn" || result.Value.NPCs[1].Name != "Captain Vale" {
|
||||||
t.Fatalf("fallback NPCs = %#v, want independently safe group applied", result.Value.NPCs)
|
t.Fatalf("fallback NPCs = %#v, want independently safe group applied", result.Value.NPCs)
|
||||||
}
|
}
|
||||||
if len(result.Retry.FallbackWarnings) != 1 || result.Retry.FallbackWarnings[0].ReasonCode != ReasonCodeNPCSemanticReconciliationExhausted || !strings.Contains(result.Retry.FallbackWarnings[0].Message, "1 proposal group") {
|
if len(result.Retry.FallbackDiagnostics) != 1 || result.Retry.FallbackDiagnostics[0].ReasonCode != ReasonCodeNPCSemanticReconciliationExhausted || !strings.Contains(result.Retry.FallbackDiagnostics[0].Samples[0].Message, "1 proposal group") {
|
||||||
t.Fatalf("fallback warnings = %#v, want exact omitted-group warning", result.Retry.FallbackWarnings)
|
t.Fatalf("fallback diagnostics = %#v, want exact omitted-group warning", result.Retry.FallbackDiagnostics)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -264,19 +263,8 @@ func TestNormalizeLimitSkipDoesNotCallLLMAndAddsBoundedFallbackWarning(t *testin
|
|||||||
if len(client.requests) != 0 || len(result.Value.NPCs) != limit+1 {
|
if len(client.requests) != 0 || len(result.Value.NPCs) != limit+1 {
|
||||||
t.Fatalf("completion calls = %d, NPCs = %d; want no call and all records", len(client.requests), len(result.Value.NPCs))
|
t.Fatalf("completion calls = %d, NPCs = %d; want no call and all records", len(client.requests), len(result.Value.NPCs))
|
||||||
}
|
}
|
||||||
if !hasWarning(result.Warnings, ReasonCodeNPCSemanticReconciliationExhausted, "npcs") || len(result.Warnings) > diagnostics.MaxWarnings {
|
if !hasDiagnostic(result.Diagnostics, ReasonCodeNPCSemanticReconciliationExhausted, "npcs", contracts.DiagnosticDispositionWarning) {
|
||||||
t.Fatalf("warnings = %#v, want bounded reconciliation fallback", result.Warnings)
|
t.Fatalf("diagnostics = %#v, want reconciliation fallback warning", result.Diagnostics)
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestRetryWarningLimitReservesExhaustionWarningPosition(t *testing.T) {
|
|
||||||
warnings := make([]contracts.Warning, 0, 25)
|
|
||||||
for index := 0; index < 25; index++ {
|
|
||||||
warnings = append(warnings, contracts.Warning{Scope: "npcs", ReasonCode: "test", Message: "warning"})
|
|
||||||
}
|
|
||||||
bounded := limitWarningsForRetry(warnings)
|
|
||||||
if len(bounded) != 19 || bounded[len(bounded)-1].ReasonCode != ReasonCodeNPCNormalizationWarningsOmitted || !strings.Contains(bounded[len(bounded)-1].Message, "7 additional") {
|
|
||||||
t.Fatalf("retry warnings = %#v, want 18 warnings plus accurate omission summary", bounded)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -22,19 +22,29 @@ const (
|
|||||||
// locally grouped advisories. These findings do not indicate process
|
// locally grouped advisories. These findings do not indicate process
|
||||||
// degradation.
|
// degradation.
|
||||||
func DataQualityResult(findings []contracts.Warning) (contracts.ValidationResult, error) {
|
func DataQualityResult(findings []contracts.Warning) (contracts.ValidationResult, error) {
|
||||||
|
diagnostics, err := Collect(findings, contracts.DiagnosticDispositionAdvisory, contracts.DiagnosticCategoryDataQuality)
|
||||||
|
if err != nil {
|
||||||
|
return contracts.ValidationResult{}, err
|
||||||
|
}
|
||||||
|
return contracts.ValidationResult{Approved: true, Diagnostics: diagnostics}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Collect converts individual deterministic findings into bounded,
|
||||||
|
// occurrence-grouped producer diagnostics with one consistent classification.
|
||||||
|
func Collect(findings []contracts.Warning, disposition contracts.DiagnosticDisposition, category contracts.DiagnosticCategory) ([]contracts.ProducerDiagnostic, error) {
|
||||||
collector := frameworkdiagnostics.NewCollector()
|
collector := frameworkdiagnostics.NewCollector()
|
||||||
for _, finding := range findings {
|
for _, finding := range findings {
|
||||||
if err := collector.Add(contracts.ProducerDiagnostic{
|
if err := collector.Add(contracts.ProducerDiagnostic{
|
||||||
Disposition: contracts.DiagnosticDispositionAdvisory,
|
Disposition: disposition,
|
||||||
Category: contracts.DiagnosticCategoryDataQuality,
|
Category: category,
|
||||||
ReasonCode: finding.ReasonCode,
|
ReasonCode: finding.ReasonCode,
|
||||||
OccurrenceCount: 1,
|
OccurrenceCount: 1,
|
||||||
Samples: []contracts.DiagnosticSample{{Scope: finding.Scope, Message: finding.Message}},
|
Samples: []contracts.DiagnosticSample{{Scope: finding.Scope, Message: finding.Message}},
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
return contracts.ValidationResult{}, fmt.Errorf("collect data-quality diagnostic: %w", err)
|
return nil, fmt.Errorf("collect diagnostic: %w", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return contracts.ValidationResult{Approved: true, Diagnostics: collector.Diagnostics()}, nil
|
return collector.Diagnostics(), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func Truncate(value string) string {
|
func Truncate(value string) string {
|
||||||
|
|||||||
@@ -69,3 +69,26 @@ func TestDataQualityResultGroupsFindingsWithoutWarnings(t *testing.T) {
|
|||||||
t.Fatalf("diagnostic = %#v", diagnostic)
|
t.Fatalf("diagnostic = %#v", diagnostic)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestCollectGroupsNormalizationFindingsWithBoundedSamples(t *testing.T) {
|
||||||
|
findings := make([]contracts.Warning, 5)
|
||||||
|
for index := range findings {
|
||||||
|
findings[index] = contracts.Warning{
|
||||||
|
Scope: fmt.Sprintf("records[%d]", index),
|
||||||
|
ReasonCode: "record_normalized",
|
||||||
|
Message: fmt.Sprintf("record %d normalized", index),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
groups, err := Collect(findings, contracts.DiagnosticDispositionObservation, contracts.DiagnosticCategoryNormalization)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(groups) != 1 {
|
||||||
|
t.Fatalf("groups = %#v, want one grouped observation", groups)
|
||||||
|
}
|
||||||
|
diagnostic := groups[0]
|
||||||
|
if diagnostic.Disposition != contracts.DiagnosticDispositionObservation || diagnostic.Category != contracts.DiagnosticCategoryNormalization || diagnostic.OccurrenceCount != 5 || len(diagnostic.Samples) != contracts.MaxDiagnosticSamples || diagnostic.OmittedSampleCount != 2 {
|
||||||
|
t.Fatalf("diagnostic = %#v", diagnostic)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user