Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions internal/api/graphql/graph/query/recipe_mapping_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
75 changes: 74 additions & 1 deletion internal/api/graphql/graph/query/recipe_resolve.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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
}
14 changes: 14 additions & 0 deletions internal/dataframe/recipe/document_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand All @@ -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.
Expand Down
22 changes: 22 additions & 0 deletions internal/dataframe/recipe/document_validation.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 5 additions & 1 deletion internal/dataframe/recipe/revisions_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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) {
Expand Down
31 changes: 31 additions & 0 deletions internal/dataframe/semantic/catalog.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
35 changes: 35 additions & 0 deletions internal/dataframe/semantic/catalog_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"testing"

"github.com/calypr/loom/internal/catalog"
"github.com/calypr/loom/internal/dataframe/recipe"
)

func TestDiscoverCatalogEmptyIsValidAndDistinctFromPartial(t *testing.T) {
Expand All @@ -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}}
Expand Down
Loading
Loading