diff --git a/internal/api/graphql/graph/query/recipe_mapping_test.go b/internal/api/graphql/graph/query/recipe_mapping_test.go index f0d451ca..0584f6ee 100644 --- a/internal/api/graphql/graph/query/recipe_mapping_test.go +++ b/internal/api/graphql/graph/query/recipe_mapping_test.go @@ -124,3 +124,33 @@ func TestGraphQLRecipeResolutionUsesScopedCatalogBeforeSemanticTyping(t *testing t.Fatalf("scoped GraphQL resolution = options:%#v bundle:%#v", got, resolved) } } + +func TestGraphQLRecipeResolutionLowersCatalogConceptSelection(t *testing.T) { + var got catalog.PopulatedFieldOptions + service := NewService(Config{DiscoverFields: func(_ context.Context, options catalog.PopulatedFieldOptions) ([]catalog.PopulatedField, error) { + got = options + return []catalog.PopulatedField{{ResourceType: "Patient", Path: "birthDate", DocCount: 2, SemanticObservations: []catalog.SemanticObservation{{ + Source: catalog.SemanticObservationSource{Canonical: "Patient.birthDate", Type: "Patient", Path: "birthDate"}, + Value: catalog.SemanticObservationValue{Selector: "birthDate", Type: "date"}, Cardinality: "single", Population: 2, + RuleHint: "future.date.v1", RuleVersion: "1", + }}}}, nil + }}) + result := semantic.DiscoverCatalog([]catalog.PopulatedField{{ResourceType: "Patient", Path: "birthDate", DocCount: 2, SemanticObservations: []catalog.SemanticObservation{{ + Source: catalog.SemanticObservationSource{Canonical: "Patient.birthDate", Type: "Patient", Path: "birthDate"}, Value: catalog.SemanticObservationValue{Selector: "birthDate", Type: "date"}, Cardinality: "single", Population: 2, RuleHint: "future.date.v1", RuleVersion: "1", + }}}}, semantic.CatalogOptions{Project: "P1"}) + concept := result.Resources[0].Families[0].Concepts[0] + bundle := recipe.Bundle{RecipeSchemaVersion: recipe.CurrentSchemaVersion, Name: "semantic", TranslationVersion: "v1", Outputs: []recipe.Output{{Name: "patients", RootResourceType: "Patient", RowGrain: "patient", ConceptSelections: []recipe.ConceptSelection{{ConceptID: concept.ID, RuleID: concept.RuleID, ColumnName: "birth_date"}}}}} + resolved, err := service.resolveRecipeBundle(context.Background(), bundle, recipe.RuntimeBindings{Project: "P1", DatasetGeneration: "g1", AuthResourcePaths: []string{"/authorized"}}) + if err != nil { + t.Fatal(err) + } + if got.ResourceType != "Patient" || len(resolved.Outputs[0].ConceptSelections) != 0 || len(resolved.Outputs[0].Fields) != 1 { + t.Fatalf("concept resolution did not produce transient executable copy: options:%#v bundle:%#v", got, resolved) + } + if resolved.Outputs[0].Fields[0].ConceptID != concept.ID || resolved.Outputs[0].Fields[0].RuleID != concept.RuleID { + t.Fatalf("concept provenance lost: %#v", resolved.Outputs[0].Fields) + } + if _, err := semantic.BuildRecipePlan(resolved, recipe.RuntimeBindings{Project: "P1", DatasetGeneration: "g1"}); err != nil { + t.Fatalf("lowered GraphQL recipe does not type-check: %v", err) + } +} diff --git a/internal/api/graphql/graph/query/recipe_resolve.go b/internal/api/graphql/graph/query/recipe_resolve.go index 5aa13453..7893e245 100644 --- a/internal/api/graphql/graph/query/recipe_resolve.go +++ b/internal/api/graphql/graph/query/recipe_resolve.go @@ -2,9 +2,13 @@ package queryapi import ( "context" + "strings" + "github.com/calypr/loom/internal/authscope" + "github.com/calypr/loom/internal/catalog" "github.com/calypr/loom/internal/dataframe/recipe" "github.com/calypr/loom/internal/dataframe/recipe/schema" + "github.com/calypr/loom/internal/dataframe/semantic" ) func (s *Service) resolveRecipeBundle(ctx context.Context, bundle recipe.Bundle, bindings recipe.RuntimeBindings) (recipe.Bundle, error) { @@ -15,5 +19,74 @@ func (s *Service) resolveRecipeBundle(ctx context.Context, bundle recipe.Bundle, if err != nil { return recipe.Bundle{}, queryInvalidErrorOrBackend(err) } - return resolved.Bundle, nil + if !bundleHasConceptSelections(resolved.Bundle) { + return resolved.Bundle, nil + } + fields := make([]catalog.PopulatedField, 0) + for _, resourceType := range conceptResourceTypes(resolved.Bundle) { + part, discoverErr := s.discoverFields(ctx, catalog.PopulatedFieldOptions{ + Project: bindings.Project, DatasetGeneration: bindings.DatasetGeneration, ResourceType: resourceType, + AuthResourcePaths: append([]string(nil), bindings.AuthResourcePaths...), + AuthResourcePathsUnrestricted: authScopeUnrestricted(bindings), + }) + if discoverErr != nil { + return recipe.Bundle{}, queryBackend(discoverErr) + } + fields = append(fields, part...) + } + fields = aggregatePopulatedFields(fields) + concepts := semantic.DiscoverCatalog(fields, semantic.CatalogOptions{Project: bindings.Project, SourceGeneration: bindings.DatasetGeneration}) + lowered, err := semantic.LowerBundleConceptSelections(resolved.Bundle, concepts.ResultsByResource()) + if err != nil { + return recipe.Bundle{}, queryInvalidErrorOrBackend(err) + } + // ConceptSelections are authoring references and remain persisted in the + // immutable draft/revision. The execution planner consumes the concrete + // fields/maps produced above; removing the references from this transient + // resolved copy prevents the legacy planner from attempting to interpret + // them a second time. + for index := range lowered.Outputs { + lowered.Outputs[index].ConceptSelections = nil + } + return lowered, nil +} + +func bundleHasConceptSelections(bundle recipe.Bundle) bool { + for _, output := range bundle.Outputs { + if len(output.ConceptSelections) != 0 { + return true + } + } + return false +} + +func conceptResourceTypes(bundle recipe.Bundle) []string { + seen := make(map[string]struct{}) + result := make([]string, 0) + for _, output := range bundle.Outputs { + if len(output.ConceptSelections) == 0 { + continue + } + resourceType := strings.TrimSpace(output.RootResourceType) + if resourceType == "" { + continue + } + if _, ok := seen[resourceType]; ok { + continue + } + seen[resourceType] = struct{}{} + result = append(result, resourceType) + } + return result +} + +func authScopeUnrestricted(bindings recipe.RuntimeBindings) *bool { + value := false + if bindings.AuthScopeMode == authscope.ReadScopeUnrestricted { + value = true + } + if bindings.AuthScopeMode == "" && len(bindings.AuthResourcePaths) == 0 { + value = true + } + return &value } diff --git a/internal/dataframe/recipe/document_types.go b/internal/dataframe/recipe/document_types.go index 5754dcbc..8b1c4351 100644 --- a/internal/dataframe/recipe/document_types.go +++ b/internal/dataframe/recipe/document_types.go @@ -45,9 +45,20 @@ type Output struct { DynamicColumns []DynamicColumn `json:"dynamicColumns,omitempty"` ExtensionColumns []ExtensionColumn `json:"extensionColumns,omitempty"` CatalogProjections []CatalogProjection `json:"catalogProjections,omitempty"` + ConceptSelections []ConceptSelection `json:"conceptSelections,omitempty"` CollisionPolicy string `json:"collisionPolicy,omitempty"` } +// ConceptSelection is an authored reference to a catalog concept. ColumnName +// is a stable output identity; Label is presentation metadata and must never +// be used to regenerate the column name. +type ConceptSelection struct { + ConceptID string `json:"conceptId"` + RuleID string `json:"ruleId"` + ColumnName string `json:"columnName"` + Label string `json:"label,omitempty"` +} + // Field projects one named semantic value into an output row. type Field struct { Name string `json:"name"` @@ -56,6 +67,9 @@ type Field struct { Fallbacks []Expression `json:"fallbacks,omitempty"` ValueMode ValueMode `json:"valueMode,omitempty"` Discovered bool `json:"-"` + ConceptID string `json:"conceptId,omitempty"` + RuleID string `json:"ruleId,omitempty"` + Label string `json:"label,omitempty"` } // ValueMode controls how a checked selector contributes values to one row. diff --git a/internal/dataframe/recipe/document_validation.go b/internal/dataframe/recipe/document_validation.go index b546322f..650fa9cf 100644 --- a/internal/dataframe/recipe/document_validation.go +++ b/internal/dataframe/recipe/document_validation.go @@ -48,6 +48,28 @@ func (b Bundle) Validate() error { if output.CollisionPolicy != "" && output.CollisionPolicy != "error" && output.CollisionPolicy != "overwrite" && output.CollisionPolicy != "coalesce" { return validationError("invalid_collision_policy", path+".collisionPolicy", "must be error, overwrite, or coalesce") } + seenConcepts := map[string]bool{} + seenConceptColumns := map[string]bool{} + for index, selection := range output.ConceptSelections { + selectionPath := fmt.Sprintf("%s.conceptSelections[%d]", path, index) + if strings.TrimSpace(selection.ConceptID) == "" { + return validationError("required", selectionPath+".conceptId", "conceptId is required") + } + if strings.TrimSpace(selection.RuleID) == "" { + return validationError("required", selectionPath+".ruleId", "ruleId is required") + } + if err := validateRecipeName(selection.ColumnName, selectionPath+".columnName"); err != nil { + return err + } + if seenConcepts[selection.ConceptID] { + return validationError("duplicate_name", selectionPath+".conceptId", "conceptId is selected more than once") + } + if seenConceptColumns[selection.ColumnName] { + return validationError("duplicate_name", selectionPath+".columnName", "columnName is selected more than once") + } + seenConcepts[selection.ConceptID] = true + seenConceptColumns[selection.ColumnName] = true + } budget := 0 if err := validateNodeShape(output.Fields, output.Filters, output.Pivots, output.Aggregates, output.Slices, path, &budget); err != nil { return err diff --git a/internal/dataframe/recipe/revisions_test.go b/internal/dataframe/recipe/revisions_test.go index 292b5e00..2b1a736c 100644 --- a/internal/dataframe/recipe/revisions_test.go +++ b/internal/dataframe/recipe/revisions_test.go @@ -30,7 +30,7 @@ func TestMemoryRevisionStoreRegistersImmutableProjectRevision(t *testing.T) { func TestProjectRecipeNameAndDraftCAS(t *testing.T) { store := NewMemoryDraftStore() - bundle := Bundle{RecipeSchemaVersion: CurrentSchemaVersion, Outputs: []Output{{Name: "Patients", RootResourceType: "Patient", RowGrain: "patient"}}} + bundle := Bundle{RecipeSchemaVersion: CurrentSchemaVersion, Outputs: []Output{{Name: "Patients", RootResourceType: "Patient", RowGrain: "patient", ConceptSelections: []ConceptSelection{{ConceptID: "concept_stable", RuleID: "future.rule.v2", ColumnName: "research_birth"}}}}} draft, err := store.SaveDraft(context.Background(), RecipeDraft{Project: "acme/cancer-study", Document: bundle}, 0) if err != nil { t.Fatal(err) @@ -45,6 +45,10 @@ func TestProjectRecipeNameAndDraftCAS(t *testing.T) { if err != nil || loaded.AuthoringDigest != draft.AuthoringDigest { t.Fatalf("loaded=%#v err=%v", loaded, err) } + selection := loaded.Document.Outputs[0].ConceptSelections + if len(selection) != 1 || selection[0].ConceptID != "concept_stable" || selection[0].RuleID != "future.rule.v2" || selection[0].ColumnName != "research_birth" { + t.Fatalf("concept selection was not persisted immutably: %#v", selection) + } } func TestProjectDataName(t *testing.T) { diff --git a/internal/dataframe/semantic/catalog.go b/internal/dataframe/semantic/catalog.go index 7576b369..0a20fdd1 100644 --- a/internal/dataframe/semantic/catalog.go +++ b/internal/dataframe/semantic/catalog.go @@ -54,6 +54,37 @@ type CatalogResult struct { Diagnostics []CatalogDiagnostic } +// ResultsByResource exposes the catalog in the transport-neutral shape used +// by recipe lowering. Catalog discovery intentionally groups concepts for +// GraphQL presentation (resource -> family -> concept), while lowering needs +// an exact identity index. Keeping this adapter here makes the catalog result +// directly consumable by the authoring/execution boundary and avoids callers +// rebuilding Concepts from labels or GraphQL models. +func (c CatalogResult) ResultsByResource() map[string]Result { + out := make(map[string]Result, len(c.Resources)) + for _, resource := range c.Resources { + result := Result{ResourceType: resource.ResourceType, Families: append([]Family(nil), resource.Families...), Concepts: []Concept{}, Diagnostics: []Diagnostic{}} + for _, family := range resource.Families { + result.Concepts = append(result.Concepts, family.Concepts...) + } + // Family concepts are the canonical catalog set. Deduplicate IDs in + // case a producer includes a concept in more than one presentation + // family, while preserving deterministic ordering from discovery. + seen := make(map[string]struct{}, len(result.Concepts)) + unique := result.Concepts[:0] + for _, concept := range result.Concepts { + if _, ok := seen[concept.ID]; ok { + continue + } + seen[concept.ID] = struct{}{} + unique = append(unique, concept) + } + result.Concepts = unique + out[resource.ResourceType] = result + } + return out +} + // DiscoverCatalog adapts authenticated, persisted catalog rows into semantic // concepts. Rows are expected to have already been filtered by project, // generation, and authorization scope; this function never broadens that diff --git a/internal/dataframe/semantic/catalog_test.go b/internal/dataframe/semantic/catalog_test.go index 2b13f431..182a584c 100644 --- a/internal/dataframe/semantic/catalog_test.go +++ b/internal/dataframe/semantic/catalog_test.go @@ -4,6 +4,7 @@ import ( "testing" "github.com/calypr/loom/internal/catalog" + "github.com/calypr/loom/internal/dataframe/recipe" ) func TestDiscoverCatalogEmptyIsValidAndDistinctFromPartial(t *testing.T) { @@ -24,6 +25,40 @@ func TestDiscoverCatalogEmptyIsValidAndDistinctFromPartial(t *testing.T) { } } +func TestCatalogResultFeedsConceptLoweringAndPlanWithoutTranslation(t *testing.T) { + fields := []catalog.PopulatedField{{ + ResourceType: "Patient", Path: "identifier[]", DocCount: 4, + SemanticObservations: []catalog.SemanticObservation{{ + Source: catalog.SemanticObservationSource{Canonical: "Patient.identifier[]", Type: "Patient", Path: "identifier[]"}, + Key: catalog.SemanticObservationKey{Selector: "identifier[].system", System: "urn:mrn"}, + Value: catalog.SemanticObservationValue{Selector: "identifier[].value", Type: "string"}, + Cardinality: CardinalityRepeated, Population: 3, Examples: []string{"123"}, + RuleHint: "IDENTIFIER_SYSTEM_VALUE", RuleVersion: "2", + }}, + }} + catalogResult := DiscoverCatalog(fields, CatalogOptions{Project: "P1", SourceGeneration: "g1"}) + byResource := catalogResult.ResultsByResource() + concepts := byResource["Patient"].Concepts + if len(concepts) != 1 { + t.Fatalf("adapter concepts = %#v", concepts) + } + bundle := recipe.Bundle{RecipeSchemaVersion: recipe.CurrentSchemaVersion, Name: "semantic", TranslationVersion: "v1", Outputs: []recipe.Output{{Name: "patients", RootResourceType: "Patient", RowGrain: "patient", ConceptSelections: []recipe.ConceptSelection{{ConceptID: concepts[0].ID, RuleID: concepts[0].RuleID, ColumnName: "mrn"}}}}} + plan, err := BuildRecipePlanWithConcepts(bundle, recipe.RuntimeBindings{Project: "P1", DatasetGeneration: "g1"}, byResource) + if err != nil { + t.Fatalf("catalog -> concept lowering -> plan: %v", err) + } + output := plan.Outputs[0] + if len(output.DynamicMaps) != 1 || output.DynamicMaps[0].Source.Expression.Selector.Path != "identifier[]" { + t.Fatalf("dynamic map unexpectedly shaped: %#v", output.DynamicMaps) + } + if output.Unnest != nil { + t.Fatalf("dynamic concept introduced row expansion: %#v", output.Unnest) + } + if len(output.ConceptColumns) != 1 || output.ConceptColumns[0].ConceptID != concepts[0].ID || !output.ConceptColumns[0].Repeated { + t.Fatalf("concept audit metadata lost: %#v", output.ConceptColumns) + } +} + func TestDiscoverCatalogMergesPartitionsAndPreservesStableIdentity(t *testing.T) { obs := catalog.SemanticObservation{SchemaVersion: 1, Source: catalog.SemanticObservationSource{Canonical: "Observation.code", Type: "Observation", Path: "code"}, Key: catalog.SemanticObservationKey{Selector: "code.coding[]", System: "s", Code: "c"}, Value: catalog.SemanticObservationValue{Selector: "valueQuantity.value", Type: "number"}, Cardinality: "repeated", Population: 3, Examples: []string{"2"}, RuleHint: "OBSERVATION_CODE_VALUE", RuleVersion: "1"} left := catalog.PopulatedField{ResourceType: "Observation", Path: "code", DocCount: 4, SemanticObservations: []catalog.SemanticObservation{obs}} diff --git a/internal/dataframe/semantic/concept_lowering.go b/internal/dataframe/semantic/concept_lowering.go new file mode 100644 index 00000000..cf94e97f --- /dev/null +++ b/internal/dataframe/semantic/concept_lowering.go @@ -0,0 +1,409 @@ +package semantic + +// This file is the authoring-to-recipe boundary for semantic concepts. A +// client supplies only stable concept identity and a stable column name; +// selector paths and row-shaping constructs always come from the producer's +// resolved Concept metadata. + +import ( + "fmt" + "regexp" + "sort" + "strings" + + "github.com/calypr/loom/internal/dataframe/recipe" + "github.com/calypr/loom/internal/dataframe/spec" +) + +var conceptColumnNamePattern = regexp.MustCompile(`^[A-Za-z_][A-Za-z0-9_]*$`) + +// ConceptSelectionLowering is the executable recipe fragment plus its audit +// records. The slices are in authored selection order and are never rebuilt +// from labels. +type ConceptSelectionLowering struct { + Fields []recipe.Field + Pivots []recipe.Pivot + DynamicColumns []recipe.DynamicColumn + Columns []ConceptColumn +} + +// LowerBundleConceptSelections returns a copy of bundle with authored concept +// selections translated into ordinary fields, pivots, and dynamic maps. The +// authored selections remain on the copy for publication/audit consumers. +func LowerBundleConceptSelections(bundle recipe.Bundle, catalogByResource map[string]Result) (recipe.Bundle, error) { + copyBundle := bundle + copyBundle.Outputs = append([]recipe.Output(nil), bundle.Outputs...) + for index := range copyBundle.Outputs { + output := copyBundle.Outputs[index] + if len(output.ConceptSelections) == 0 { + continue + } + catalog, ok := catalogByResource[output.RootResourceType] + if !ok { + return recipe.Bundle{}, conceptDiagnostic(fmt.Sprintf("outputs[%d].conceptSelections", index), "CONCEPT_CATALOG_MISSING", fmt.Sprintf("no producer catalog was supplied for resource %q", output.RootResourceType)) + } + lowered, err := LowerConceptSelections(output.RootResourceType, output.ConceptSelections, catalog) + if err != nil { + return recipe.Bundle{}, fmt.Errorf("outputs[%d] %w", index, err) + } + used := map[string]struct{}{} + for _, field := range output.Fields { + used[field.Name] = struct{}{} + } + for _, field := range lowered.Fields { + if _, exists := used[field.Name]; exists { + return recipe.Bundle{}, conceptDiagnostic(fmt.Sprintf("outputs[%d].conceptSelections", index), "COLUMN_NAME_COLLISION", fmt.Sprintf("concept column %q collides with an explicit field", field.Name)) + } + } + output.Fields = append(lowered.Fields, output.Fields...) + output.Pivots = append(lowered.Pivots, output.Pivots...) + output.DynamicColumns = append(lowered.DynamicColumns, output.DynamicColumns...) + copyBundle.Outputs[index] = output + } + return copyBundle, nil +} + +// BuildRecipePlanWithConcepts is the catalog-aware counterpart to +// BuildRecipePlan. It is intentionally additive so legacy recipe callers do +// not need to provide a catalog. +func BuildRecipePlanWithConcepts(bundle recipe.Bundle, bindings recipe.RuntimeBindings, catalogByResource map[string]Result) (RecipePlan, error) { + lowered, err := LowerBundleConceptSelections(bundle, catalogByResource) + if err != nil { + return RecipePlan{}, err + } + // The authored selections remain in the persisted bundle returned by + // LowerBundleConceptSelections. The executable plan carries their identity + // on each generated field, while the ordinary recipe planner consumes only + // concrete recipe constructs. + for index := range lowered.Outputs { + lowered.Outputs[index].ConceptSelections = nil + } + plan, err := BuildRecipePlan(lowered, bindings) + if err != nil { + return RecipePlan{}, err + } + for index := range plan.Outputs { + if len(bundle.Outputs[index].ConceptSelections) == 0 { + continue + } + metadata, metadataErr := LowerConceptSelections(bundle.Outputs[index].RootResourceType, bundle.Outputs[index].ConceptSelections, catalogByResource[bundle.Outputs[index].RootResourceType]) + if metadataErr != nil { + return RecipePlan{}, metadataErr + } + plan.Outputs[index].ConceptColumns = metadata.Columns + } + return plan, nil +} + +// LowerConceptSelections verifies authored identity/rule pairs against a +// producer-resolved catalog and lowers them to ordinary recipe constructs. +// Unknown rule/source strings are accepted when the resolved concept carries +// executable selector metadata; they are producer metadata, not enums. +func LowerConceptSelections(resourceType string, authored []recipe.ConceptSelection, catalog Result) (ConceptSelectionLowering, error) { + resourceType = strings.TrimSpace(resourceType) + if resourceType == "" { + return ConceptSelectionLowering{}, fmt.Errorf("concept selections: resource type is required") + } + byID := make(map[string][]Concept, len(catalog.Concepts)) + for _, concept := range catalog.Concepts { + byID[strings.TrimSpace(concept.ID)] = append(byID[strings.TrimSpace(concept.ID)], concept) + } + for _, family := range catalog.Families { + for _, concept := range family.Concepts { + if _, exists := byID[concept.ID]; !exists { + byID[concept.ID] = append(byID[concept.ID], concept) + } + } + } + + out := ConceptSelectionLowering{ + Fields: make([]recipe.Field, 0, len(authored)), + Pivots: make([]recipe.Pivot, 0), DynamicColumns: make([]recipe.DynamicColumn, 0), + Columns: make([]ConceptColumn, 0, len(authored)), + } + seenColumns := map[string]struct{}{} + seenConceptIDs := map[string]struct{}{} + for index, selection := range authored { + path := fmt.Sprintf("conceptSelections[%d]", index) + id := strings.TrimSpace(selection.ConceptID) + ruleID := strings.TrimSpace(selection.RuleID) + column := strings.TrimSpace(selection.ColumnName) + if id == "" { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_ID_REQUIRED", "conceptId is required") + } + if _, exists := seenConceptIDs[id]; exists { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_AMBIGUOUS", fmt.Sprintf("concept %q is selected more than once", id)) + } + seenConceptIDs[id] = struct{}{} + if ruleID == "" { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_RULE_REQUIRED", "ruleId is required") + } + if !conceptColumnNamePattern.MatchString(column) { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "COLUMN_NAME_INVALID", fmt.Sprintf("columnName %q is not a safe stable column name", column)) + } + if _, exists := seenColumns[column]; exists { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "COLUMN_NAME_COLLISION", fmt.Sprintf("columnName %q is selected more than once", column)) + } + seenColumns[column] = struct{}{} + candidates := byID[id] + if len(candidates) == 0 { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_NOT_FOUND", fmt.Sprintf("concept %q is stale or absent from the producer catalog", id)) + } + if len(candidates) > 1 { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_AMBIGUOUS", fmt.Sprintf("concept %q resolves to %d producer definitions", id, len(candidates))) + } + concept := candidates[0] + if strings.TrimSpace(concept.RuleID) != ruleID { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_RULE_MISMATCH", fmt.Sprintf("concept %q was authored with ruleId %q but catalog resolves ruleId %q", id, ruleID, concept.RuleID)) + } + if strings.TrimSpace(concept.Source.ResourceType) != "" && strings.TrimSpace(concept.Source.ResourceType) != resourceType { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_RESOURCE_MISMATCH", fmt.Sprintf("concept %q belongs to resource %q, not %q", id, concept.Source.ResourceType, resourceType)) + } + fragment, err := lowerResolvedConcept(resourceType, selection, concept, path) + if err != nil { + return ConceptSelectionLowering{}, err + } + out.Fields = append(out.Fields, fragment.Fields...) + out.Pivots = append(out.Pivots, fragment.Pivots...) + out.DynamicColumns = append(out.DynamicColumns, fragment.DynamicColumns...) + out.Columns = append(out.Columns, fragment.Columns...) + } + return out, nil +} + +func conceptDiagnostic(path, code, message string) error { + return fmt.Errorf("%s [%s]: %s", path, code, message) +} + +func lowerResolvedConcept(resourceType string, authored recipe.ConceptSelection, concept Concept, path string) (ConceptSelectionLowering, error) { + s := concept.Output.Selection + mode := strings.TrimSpace(s.Mode) + if mode == "" { + mode = strings.TrimSpace(concept.Output.Mode) + } + if mode == "" { + mode = OutputScalar + } + column := ConceptColumn{ + Name: authored.ColumnName, ConceptID: authored.ConceptID, RuleID: authored.RuleID, + Label: firstNonEmpty(authored.Label, concept.Label), Selector: s, + LogicalType: firstNonEmpty(concept.Output.ValueType, concept.Source.Primitive, "unknown"), + Repeated: concept.Output.Cardinality == CardinalityRepeated || concept.Output.Cardinality == CardinalityPivoted || concept.Source.Repeated, + } + result := ConceptSelectionLowering{Columns: []ConceptColumn{column}} + sourcePath := conceptSelectorPath(resourceType, firstNonEmpty(s.SourcePath, concept.Source.Path)) + if sourcePath == "" && strings.TrimSpace(s.ItemSource) == "" && strings.TrimSpace(s.ValueSelector) == "" { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_SELECTOR_MISSING", "resolved concept has no executable sourcePath or itemSource") + } + + switch mode { + case OutputDynamicFamily: + dynamic, err := lowerDynamicConcept(resourceType, authored, concept, path) + if err != nil { + return ConceptSelectionLowering{}, err + } + result.DynamicColumns = append(result.DynamicColumns, dynamic) + return result, nil + case OutputMeasurement: + if concept.Output.Cardinality == CardinalityPivoted && strings.TrimSpace(s.ItemSource) != "" && len(conceptExamples(concept)) > 0 { + pivot, err := lowerPivotConcept(resourceType, authored, concept, path) + if err != nil { + return ConceptSelectionLowering{}, err + } + result.Pivots = append(result.Pivots, pivot) + return result, nil + } + // A measurement without bounded pivot keys is still an executable + // scalar/array projection. This preserves repeated values as arrays. + } + + valuePath := joinConceptPath(sourcePath, s.ValueSelector) + if valuePath == "" { + valuePath = sourcePath + } + if valuePath == "" { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_SELECTOR_UNSUPPORTED", "resolved concept selector has no executable value path") + } + if err := validateConceptExecutableSelector(resourceType, valuePath); err != nil { + return ConceptSelectionLowering{}, conceptDiagnostic(path, "CONCEPT_SELECTOR_UNSUPPORTED", err.Error()) + } + field := recipe.Field{ + Name: authored.ColumnName, FieldRef: authored.ConceptID, + Expr: recipe.Expression{Select: "root." + valuePath}, ConceptID: authored.ConceptID, + RuleID: authored.RuleID, Label: firstNonEmpty(authored.Label, concept.Label), + } + for _, fallback := range s.ValueFallbacks { + fallbackPath := joinConceptPath(sourcePath, fallback) + if fallbackPath != "" { + field.Fallbacks = append(field.Fallbacks, recipe.Expression{Select: "root." + fallbackPath}) + } + } + if column.Repeated { + field.ValueMode = recipe.ValueModeAll + } + result.Fields = append(result.Fields, field) + return result, nil +} + +func lowerDynamicConcept(resourceType string, authored recipe.ConceptSelection, concept Concept, path string) (recipe.DynamicColumn, error) { + s := concept.Output.Selection + itemSource := conceptSelectorPath(resourceType, firstNonEmpty(s.ItemSource, s.SourcePath, concept.Source.Path)) + if itemSource == "" { + return recipe.DynamicColumn{}, conceptDiagnostic(path, "CONCEPT_DYNAMIC_SOURCE_MISSING", "dynamic concept has no executable itemSource") + } + if !strings.Contains(itemSource, "[]") { + itemSource += "[]" + } + if err := validateConceptExecutableSelector(resourceType, itemSource); err != nil { + return recipe.DynamicColumn{}, conceptDiagnostic(path, "CONCEPT_DYNAMIC_SOURCE_UNSUPPORTED", err.Error()) + } + if strings.TrimSpace(s.KeySelector) == "" || strings.TrimSpace(s.ValueSelector) == "" { + return recipe.DynamicColumn{}, conceptDiagnostic(path, "CONCEPT_DYNAMIC_SELECTOR_MISSING", "dynamic concept requires keySelector and valueSelector") + } + keySelector := selectorRelativeToItem(itemSource, s.KeySelector) + valueSelector := selectorRelativeToItem(itemSource, s.ValueSelector) + key := recipe.Expression{Select: "item." + keySelector} + value := recipe.Expression{Select: "item." + valueSelector} + dynamic := recipe.DynamicColumn{ + Name: authored.ColumnName, Source: recipe.Expression{Select: "root." + itemSource}, + Key: &key, Value: &value, MaxColumns: 256, Discovered: true, + } + if values := conceptExamples(concept); len(values) > 0 { + dynamic.Columns = values + } + return dynamic, nil +} + +func lowerPivotConcept(resourceType string, authored recipe.ConceptSelection, concept Concept, path string) (recipe.Pivot, error) { + s := concept.Output.Selection + itemSource := conceptSelectorPath(resourceType, s.ItemSource) + if itemSource == "" { + return recipe.Pivot{}, conceptDiagnostic(path, "CONCEPT_PIVOT_SOURCE_MISSING", "pivot concept has no executable itemSource") + } + if !strings.Contains(itemSource, "[]") { + itemSource += "[]" + } + if err := validateConceptExecutableSelector(resourceType, itemSource); err != nil { + return recipe.Pivot{}, conceptDiagnostic(path, "CONCEPT_PIVOT_SOURCE_UNSUPPORTED", err.Error()) + } + columns := conceptExamples(concept) + if len(columns) == 0 { + return recipe.Pivot{}, conceptDiagnostic(path, "CONCEPT_PIVOT_COLUMNS_MISSING", "pivot concept requires bounded key examples or catalog columns") + } + columnSelector := joinConceptPath(itemSource, selectorRelativeToItem(itemSource, s.KeySelector)) + valueSelector := joinConceptPath(itemSource, selectorRelativeToItem(itemSource, s.ValueSelector)) + pivot := recipe.Pivot{ + Name: authored.ColumnName, FieldRef: authored.ConceptID, + ColumnExpr: recipe.Expression{Select: "root." + columnSelector}, + ValueExpr: recipe.Expression{Select: "root." + valueSelector}, + ItemSource: recipe.Expression{Select: "root." + itemSource}, + ItemResourceType: resourceType, Columns: columns, Discovered: true, + } + for _, fallback := range s.ValueFallbacks { + fallbackPath := joinConceptPath(itemSource, fallback) + if fallbackPath != "" { + pivot.ValueFallbacks = append(pivot.ValueFallbacks, recipe.Expression{Select: "root." + fallbackPath}) + } + } + return pivot, nil +} + +func conceptExamples(concept Concept) []string { + values := append([]string(nil), concept.Examples.Values...) + if len(values) == 0 { + for _, example := range concept.Source.Examples { + if example.Safe && strings.TrimSpace(example.Value) != "" { + values = append(values, strings.TrimSpace(example.Value)) + } + } + } + values = uniqueNonEmpty(values) + sort.Strings(values) + return values +} + +func conceptSelectorPath(resourceType, path string) string { + path = strings.TrimSpace(strings.TrimPrefix(path, "root.")) + path = strings.TrimPrefix(path, ".") + resourcePrefix := strings.TrimSpace(resourceType) + "." + if strings.HasPrefix(path, resourcePrefix) { + path = strings.TrimPrefix(path, resourcePrefix) + } + if path == strings.TrimSpace(resourceType) { + return "" + } + return path +} + +func joinConceptPath(source, relative string) string { + source = strings.TrimSpace(strings.TrimPrefix(source, "root.")) + relative = strings.TrimSpace(strings.TrimPrefix(relative, ".")) + if relative == "" { + return source + } + if source == "" { + return relative + } + if relative == source || strings.HasPrefix(relative, source+".") { + return relative + } + return source + "." + relative +} + +// selectorRelativeToItem converts catalog selectors that are recorded from +// the document root (for example identifier[].value) into selectors relative +// to the repeated item binding used by dynamic maps and pivots (value). Older +// producer catalogs may already contain item-relative selectors, which are +// retained unchanged. +func selectorRelativeToItem(itemSource, selector string) string { + itemSource = strings.TrimSpace(strings.TrimPrefix(itemSource, "root.")) + selector = strings.TrimSpace(strings.TrimPrefix(selector, "root.")) + selector = strings.TrimPrefix(selector, ".") + if itemSource == "" { + return selector + } + if selector == itemSource { + return "" + } + if strings.HasPrefix(selector, itemSource+".") { + return strings.TrimPrefix(selector, itemSource+".") + } + return selector +} + +func firstNonEmpty(values ...string) string { + for _, value := range values { + if strings.TrimSpace(value) != "" { + return strings.TrimSpace(value) + } + } + return "" +} + +func uniqueNonEmpty(values []string) []string { + seen := map[string]struct{}{} + out := make([]string, 0, len(values)) + for _, value := range values { + value = strings.TrimSpace(value) + if value == "" { + continue + } + if _, ok := seen[value]; ok { + continue + } + seen[value] = struct{}{} + out = append(out, value) + } + return out +} + +func validateConceptExecutableSelector(resourceType, path string) error { + selector, err := spec.ParseSelector(path) + if err != nil { + return err + } + if _, _, err := spec.SelectorCardinality(resourceType, selector); err != nil { + return err + } + return nil +} diff --git a/internal/dataframe/semantic/concept_lowering_test.go b/internal/dataframe/semantic/concept_lowering_test.go new file mode 100644 index 00000000..49f70576 --- /dev/null +++ b/internal/dataframe/semantic/concept_lowering_test.go @@ -0,0 +1,79 @@ +package semantic + +import ( + "strings" + "testing" + + "github.com/calypr/loom/internal/dataframe/recipe" +) + +func TestLowerConceptSelectionsPreservesIdentityAndRepeatedArrays(t *testing.T) { + catalog := Result{ResourceType: "Patient", Concepts: []Concept{ + {ID: "patient.birth", RuleID: "future.direct.v9", Label: "Birth date", Source: SourceDescriptor{ResourceType: "Patient", Path: "birthDate"}, Output: OutputDescriptor{Mode: OutputScalar, ValueType: "date", Selection: Selection{Mode: OutputScalar, SourcePath: "Patient", ValueSelector: "birthDate"}}}, + {ID: "patient.identifiers", RuleID: "future.identifier.v2", Label: "Identifiers", Source: SourceDescriptor{ResourceType: "Patient", Path: "identifier[]", Repeated: true}, Output: OutputDescriptor{Mode: OutputScalar, ValueType: "string", Cardinality: CardinalityRepeated, Selection: Selection{Mode: OutputScalar, SourcePath: "Patient", ValueSelector: "identifier[].value"}}}, + }} + got, err := LowerConceptSelections("Patient", []recipe.ConceptSelection{ + {ConceptID: "patient.birth", RuleID: "future.direct.v9", ColumnName: "birth_date", Label: "Research birth date"}, + {ConceptID: "patient.identifiers", RuleID: "future.identifier.v2", ColumnName: "identifiers"}, + }, catalog) + if err != nil { + t.Fatal(err) + } + if len(got.Fields) != 2 || got.Fields[0].Expr.Select != "root.birthDate" { + t.Fatalf("unexpected lowered fields: %#v", got.Fields) + } + if got.Fields[0].Label != "Research birth date" || got.Fields[0].ConceptID != "patient.birth" { + t.Fatalf("authored metadata was not retained: %#v", got.Fields[0]) + } + if got.Fields[1].ValueMode != recipe.ValueModeAll || got.Fields[1].Expr.Select != "root.identifier[].value" { + t.Fatalf("repeated concept was not lowered as an array: %#v", got.Fields[1]) + } +} + +func TestLowerConceptSelectionsRejectsStaleAndMismatchedIdentity(t *testing.T) { + catalog := Result{ResourceType: "Patient", Concepts: []Concept{{ID: "patient.birth", RuleID: "direct.v1", Source: SourceDescriptor{ResourceType: "Patient", Path: "birthDate"}, Output: OutputDescriptor{Selection: Selection{SourcePath: "Patient", ValueSelector: "birthDate"}}}}} + _, err := LowerConceptSelections("Patient", []recipe.ConceptSelection{{ConceptID: "patient.missing", RuleID: "direct.v1", ColumnName: "missing"}}, catalog) + if err == nil || !strings.Contains(err.Error(), "CONCEPT_NOT_FOUND") { + t.Fatalf("stale concept error=%v", err) + } + _, err = LowerConceptSelections("Patient", []recipe.ConceptSelection{{ConceptID: "patient.birth", RuleID: "old.v1", ColumnName: "birth"}}, catalog) + if err == nil || !strings.Contains(err.Error(), "CONCEPT_RULE_MISMATCH") { + t.Fatalf("rule mismatch error=%v", err) + } +} + +func TestLowerBundleConceptSelectionsUsesUnknownRuleMetadata(t *testing.T) { + bundle := recipe.Bundle{RecipeSchemaVersion: recipe.CurrentSchemaVersion, Name: "semantic", TranslationVersion: "v1", Outputs: []recipe.Output{{Name: "patients", RootResourceType: "Patient", RowGrain: "patient", ConceptSelections: []recipe.ConceptSelection{{ConceptID: "patient.future", RuleID: "FutureClinicalSource.rule.v8", ColumnName: "future_score"}}}}} + catalog := map[string]Result{"Patient": {ResourceType: "Patient", Concepts: []Concept{{ID: "patient.future", RuleID: "FutureClinicalSource.rule.v8", Source: SourceDescriptor{ResourceType: "Patient", Path: "birthDate"}, Output: OutputDescriptor{Mode: "future_value_family", ValueType: "futureDecimal128", Selection: Selection{Mode: "future_value_family", SourcePath: "Patient", ValueSelector: "birthDate"}}}}}} + lowered, err := LowerBundleConceptSelections(bundle, catalog) + if err != nil { + t.Fatal(err) + } + if len(lowered.Outputs[0].Fields) != 1 || lowered.Outputs[0].Fields[0].RuleID != "FutureClinicalSource.rule.v8" { + t.Fatalf("unknown future rule was not retained: %#v", lowered.Outputs[0].Fields) + } + if len(lowered.Outputs[0].ConceptSelections) != 1 { + t.Fatalf("authored concept selection was dropped") + } + plan, err := BuildRecipePlanWithConcepts(bundle, recipe.RuntimeBindings{}, catalog) + if err != nil { + t.Fatalf("build lowered concept plan: %v", err) + } + if len(plan.Outputs) != 1 || len(plan.Outputs[0].ConceptColumns) != 1 || plan.Outputs[0].ConceptColumns[0].ConceptID != "patient.future" || plan.Outputs[0].ConceptColumns[0].Selector.ValueSelector != "birthDate" { + t.Fatalf("concept audit metadata did not reach semantic plan: %#v", plan.Outputs[0].ConceptColumns) + } +} + +func TestLowerDynamicConceptUsesItemSelectorsWithoutRowExpansion(t *testing.T) { + catalog := Result{ResourceType: "Patient", Concepts: []Concept{{ID: "patient.identifiers", RuleID: "identifier.v1", Source: SourceDescriptor{ResourceType: "Patient", Path: "identifier[]", Repeated: true}, Output: OutputDescriptor{Mode: OutputDynamicFamily, Cardinality: CardinalityRepeated, Selection: Selection{Mode: OutputDynamicFamily, ItemSource: "identifier[]", KeySelector: "system", ValueSelector: "value"}}}}} + got, err := LowerConceptSelections("Patient", []recipe.ConceptSelection{{ConceptID: "patient.identifiers", RuleID: "identifier.v1", ColumnName: "identifier_values"}}, catalog) + if err != nil { + t.Fatal(err) + } + if len(got.DynamicColumns) != 1 || got.DynamicColumns[0].Source.Select != "root.identifier[]" || got.DynamicColumns[0].Key.Select != "item.system" || got.DynamicColumns[0].Value.Select != "item.value" { + t.Fatalf("unexpected dynamic lowering: %#v", got.DynamicColumns) + } + if len(got.Fields) != 0 { + t.Fatalf("dynamic concept unexpectedly expanded into fields") + } +} diff --git a/internal/dataframe/semantic/recipe_plan_output.go b/internal/dataframe/semantic/recipe_plan_output.go index 905eb4dd..562c660b 100644 --- a/internal/dataframe/semantic/recipe_plan_output.go +++ b/internal/dataframe/semantic/recipe_plan_output.go @@ -16,6 +16,9 @@ import ( ) func buildRecipeOutput(output recipe.Output, bindings recipe.RuntimeBindings) (OutputPlan, error) { + if len(output.ConceptSelections) != 0 { + return OutputPlan{}, fmt.Errorf("conceptSelections require producer catalog resolution before recipe planning") + } if !fhirschema.HasResource(output.RootResourceType) { return OutputPlan{}, fmt.Errorf("root resource type %q is not represented by the active generated FHIR schema", output.RootResourceType) } @@ -110,6 +113,13 @@ func finishRecipeOutput(plan OutputPlan, output recipe.Output, scope scopeFrame) plan.Fields = append(plan.Fields, normalized.projection) plan.DeclaredOrder = append(plan.DeclaredOrder, field.Name) plan.Root.Fields = append(plan.Root.Fields, normalized.field) + if field.ConceptID != "" { + plan.ConceptColumns = append(plan.ConceptColumns, ConceptColumn{ + Name: field.Name, ConceptID: field.ConceptID, RuleID: field.RuleID, + Label: field.Label, LogicalType: string(normalized.projection.Expr.Type.Kind), + Repeated: normalized.projection.Expr.Type.Cardinality == expression.Many, + }) + } } for index, traversal := range output.Traversals { child, err := buildRecipeTraversal(traversal, scope, fmt.Sprintf("traversals[%d]", index)) diff --git a/internal/dataframe/semantic/recipe_plan_types.go b/internal/dataframe/semantic/recipe_plan_types.go index 8fdbc72a..2f9eb1a6 100644 --- a/internal/dataframe/semantic/recipe_plan_types.go +++ b/internal/dataframe/semantic/recipe_plan_types.go @@ -44,6 +44,20 @@ type OutputPlan struct { CatalogProjections []string Collision string DeclaredOrder []string + ConceptColumns []ConceptColumn +} + +// ConceptColumn is the immutable audit record carried with a lowered output +// column. It keeps authored identity and producer rule metadata available to +// validation, resolution, and publication consumers. +type ConceptColumn struct { + Name string + ConceptID string + RuleID string + Label string + Selector Selection + LogicalType string + Repeated bool } // SemanticExpression keeps the checked typed AST together with the logical @@ -62,6 +76,9 @@ type SemanticProjection struct { ValueMode string Expr SemanticExpression Discovered bool + ConceptID string + RuleID string + Label string } // UnnestJoinMode makes null/empty collection behavior explicit at the diff --git a/internal/dataframe/semantic/recipe_projection.go b/internal/dataframe/semantic/recipe_projection.go index 2ed5883d..f9beb5f6 100644 --- a/internal/dataframe/semantic/recipe_projection.go +++ b/internal/dataframe/semantic/recipe_projection.go @@ -73,6 +73,7 @@ func normalizeRecipeProjection(field recipe.Field, scope scopeFrame, path string Name: field.Name, FieldRef: field.FieldRef, Selector: selector, Fallbacks: fallbackSelectors, ValueMode: string(field.ValueMode), Expr: &primary.Expression, ExprType: primary.Type, SourcePath: primary.SourcePath, Discovered: field.Discovered, + ConceptID: field.ConceptID, RuleID: field.RuleID, Label: field.Label, } if err := validateRecipeProjectionTypes(primary.Type, fallbacks); err != nil { return normalizedRecipeProjection{}, fmt.Errorf("%s: %w", path, err) @@ -102,6 +103,7 @@ func normalizeRecipeProjection(field recipe.Field, scope scopeFrame, path string fieldSemantic.ExprType = projected.Type return normalizedRecipeProjection{projection: SemanticProjection{ Name: field.Name, FieldRef: field.FieldRef, ValueMode: string(field.ValueMode), Expr: projected, Discovered: field.Discovered, + ConceptID: field.ConceptID, RuleID: field.RuleID, Label: field.Label, }, field: fieldSemantic}, nil } @@ -114,6 +116,7 @@ func normalizeRecipeProjection(field recipe.Field, scope scopeFrame, path string } return normalizedRecipeProjection{projection: SemanticProjection{Name: field.Name, FieldRef: field.FieldRef, ValueMode: string(field.ValueMode), Expr: projected}, field: SemanticField{ Name: field.Name, FieldRef: field.FieldRef, ValueMode: string(field.ValueMode), Expr: &projected.Expression, ExprType: projected.Type, SourcePath: projected.SourcePath, Discovered: field.Discovered, + ConceptID: field.ConceptID, RuleID: field.RuleID, Label: field.Label, }}, nil } diff --git a/internal/dataframe/semantic/resolved_plan.go b/internal/dataframe/semantic/resolved_plan.go index 83035d9f..57e51558 100644 --- a/internal/dataframe/semantic/resolved_plan.go +++ b/internal/dataframe/semantic/resolved_plan.go @@ -22,8 +22,12 @@ type ResolvedColumn struct { // production execution/materialization adapter. Stored recipe data and // request-scoped discovery are never mutated in place. type ResolvedRecipePlan struct { - SemanticPlan RecipePlan - ResolvedColumns map[string][]ResolvedColumn + SemanticPlan RecipePlan + ResolvedColumns map[string][]ResolvedColumn + // ConceptColumns is carried independently of dynamic discovery so + // publication adapters can retain authored concept identity even when a + // release has no dynamic families. + ConceptColumns map[string][]ConceptColumn ResolvedSchemaDigest string ScopeDigest string SourceGeneration string @@ -51,10 +55,14 @@ func ResolveRecipePlan(plan RecipePlan, scopeDigest, sourceGeneration string) (R resolved := ResolvedRecipePlan{ SemanticPlan: plan, ResolvedColumns: make(map[string][]ResolvedColumn), + ConceptColumns: make(map[string][]ConceptColumn), ScopeDigest: scopeDigest, SourceGeneration: sourceGeneration, } for _, output := range plan.Outputs { + if len(output.ConceptColumns) > 0 { + resolved.ConceptColumns[output.Name] = append([]ConceptColumn(nil), output.ConceptColumns...) + } var walk func(SemanticNode) error walk = func(node SemanticNode) error { for _, dynamic := range node.DynamicMaps { @@ -121,7 +129,8 @@ func ResolveRecipePlan(plan RecipePlan, scopeDigest, sourceGeneration string) (R canonical, err := json.Marshal(struct { RecipeDigest, ScopeDigest, Generation string Columns any `json:"columns"` - }{plan.RecipeDigest, scopeDigest, sourceGeneration, ordered}) + ConceptColumns any `json:"conceptColumns"` + }{plan.RecipeDigest, scopeDigest, sourceGeneration, ordered, resolved.ConceptColumns}) if err != nil { return ResolvedRecipePlan{}, fmt.Errorf("resolved schema digest: %w", err) } diff --git a/internal/dataframe/semantic/semantic_plan.go b/internal/dataframe/semantic/semantic_plan.go index dba76395..c69852a4 100644 --- a/internal/dataframe/semantic/semantic_plan.go +++ b/internal/dataframe/semantic/semantic_plan.go @@ -45,6 +45,9 @@ type SemanticField struct { ExprType expression.Type SourcePath string Discovered bool + ConceptID string + RuleID string + Label string } type SemanticPivot struct { diff --git a/internal/server/recipe_discovery.go b/internal/server/recipe_discovery.go index 6e47e187..7b6527ba 100644 --- a/internal/server/recipe_discovery.go +++ b/internal/server/recipe_discovery.go @@ -3,11 +3,13 @@ package server import ( "context" "fmt" + "strings" "github.com/calypr/loom/internal/authscope" "github.com/calypr/loom/internal/catalog" "github.com/calypr/loom/internal/dataframe/recipe" "github.com/calypr/loom/internal/dataframe/recipe/schema" + "github.com/calypr/loom/internal/dataframe/semantic" ) type recipeCatalogDiscovery struct { @@ -86,6 +88,66 @@ func recipeSchemaResolver(read func(context.Context, catalog.PopulatedFieldOptio if err != nil { return recipe.Bundle{}, err } - return resolved.Bundle, nil + if !bundleHasConceptSelections(resolved.Bundle) { + return resolved.Bundle, nil + } + fields := make([]catalog.PopulatedField, 0) + for _, resourceType := range conceptResourceTypes(resolved.Bundle) { + part, discoverErr := read(ctx, catalog.PopulatedFieldOptions{ + Project: bindings.Project, DatasetGeneration: bindings.DatasetGeneration, ResourceType: resourceType, + AuthResourcePaths: append([]string(nil), bindings.AuthResourcePaths...), + AuthResourcePathsUnrestricted: unrestrictedAuthScope(bindings), + }) + if discoverErr != nil { + return recipe.Bundle{}, discoverErr + } + fields = append(fields, part...) + } + concepts := semantic.DiscoverCatalog(fields, semantic.CatalogOptions{Project: bindings.Project, SourceGeneration: bindings.DatasetGeneration}) + lowered, err := semantic.LowerBundleConceptSelections(resolved.Bundle, concepts.ResultsByResource()) + if err != nil { + return recipe.Bundle{}, err + } + // Keep authored selections in the durable draft/revision; this transient + // execution copy contains only concrete recipe constructs for the legacy + // planner and still carries concept identity on generated columns. + for index := range lowered.Outputs { + lowered.Outputs[index].ConceptSelections = nil + } + return lowered, nil } } + +func bundleHasConceptSelections(bundle recipe.Bundle) bool { + for _, output := range bundle.Outputs { + if len(output.ConceptSelections) != 0 { + return true + } + } + return false +} + +func conceptResourceTypes(bundle recipe.Bundle) []string { + seen := make(map[string]struct{}) + result := make([]string, 0) + for _, output := range bundle.Outputs { + if len(output.ConceptSelections) == 0 { + continue + } + resourceType := strings.TrimSpace(output.RootResourceType) + if resourceType == "" { + continue + } + if _, ok := seen[resourceType]; ok { + continue + } + seen[resourceType] = struct{}{} + result = append(result, resourceType) + } + return result +} + +func unrestrictedAuthScope(bindings recipe.RuntimeBindings) *bool { + value := bindings.AuthScopeMode == authscope.ReadScopeUnrestricted || (bindings.AuthScopeMode == "" && len(bindings.AuthResourcePaths) == 0) + return &value +}