|
| 1 | +from collections import defaultdict |
| 2 | + |
| 3 | +from cognite.neat._data_model._fix import FixAction |
| 4 | +from cognite.neat._data_model.deployer.data_classes import ( |
| 5 | + AddedField, |
| 6 | + ChangedField, |
| 7 | + FieldChange, |
| 8 | + PrimitiveField, |
| 9 | + RemovedField, |
| 10 | +) |
| 11 | +from cognite.neat._data_model.models.dms import ( |
| 12 | + ContainerReference, |
| 13 | + DataModelReference, |
| 14 | + DataModelResource, |
| 15 | + SchemaResourceId, |
| 16 | + SpaceReference, |
| 17 | + ViewReference, |
| 18 | +) |
| 19 | +from cognite.neat._data_model.models.dms._schema import RequestSchema |
| 20 | +from cognite.neat._data_model.transformers._base import Transformer |
| 21 | + |
| 22 | + |
| 23 | +class FixApplicator(Transformer): |
| 24 | + """Applies the changes in FixAction objects to a schema.""" |
| 25 | + |
| 26 | + def __init__(self, fix_actions: list[FixAction]) -> None: |
| 27 | + self._fix_actions = fix_actions |
| 28 | + |
| 29 | + def transform(self, data_model: RequestSchema) -> RequestSchema: |
| 30 | + """Apply fix actions and return the fixed schema (a deep copy of the original).""" |
| 31 | + result = data_model.model_copy(deep=True) |
| 32 | + |
| 33 | + if not self._fix_actions: |
| 34 | + return result |
| 35 | + |
| 36 | + fix_by_resource_id: dict[SchemaResourceId, list[FixAction]] = defaultdict(list) |
| 37 | + for action in self._fix_actions: |
| 38 | + fix_by_resource_id[action.resource_id].append(action) |
| 39 | + |
| 40 | + resources_list_lookup: dict[type, dict[SchemaResourceId, DataModelResource]] = { |
| 41 | + ViewReference: {view.as_reference(): view for view in result.views}, |
| 42 | + ContainerReference: {container.as_reference(): container for container in result.containers}, |
| 43 | + SpaceReference: {space.as_reference(): space for space in result.spaces}, |
| 44 | + DataModelReference: {result.data_model.as_reference(): result.data_model}, |
| 45 | + } |
| 46 | + |
| 47 | + for resource_id, actions in fix_by_resource_id.items(): |
| 48 | + resource_lookup = resources_list_lookup.get(type(resource_id)) |
| 49 | + if resource_lookup is None: |
| 50 | + raise RuntimeError( |
| 51 | + f"{type(self).__name__}: Unsupported resource type {type(resource_id)}. This is a bug in NEAT." |
| 52 | + ) |
| 53 | + resource = resource_lookup.get(resource_id) |
| 54 | + if resource is None: |
| 55 | + raise RuntimeError( |
| 56 | + f"{type(self).__name__}: Resource {resource_id} not found in schema. This is a bug in NEAT." |
| 57 | + ) |
| 58 | + |
| 59 | + all_changes_for_resource = [change for action in actions for change in action.changes] |
| 60 | + self._check_no_field_path_conflicts(all_changes_for_resource) |
| 61 | + self._apply_changes_to_resource(resource, all_changes_for_resource) |
| 62 | + |
| 63 | + return result |
| 64 | + |
| 65 | + def _apply_changes_to_resource(self, resource: DataModelResource, changes: list[FieldChange]) -> None: |
| 66 | + """Apply field changes to the resource in place.""" |
| 67 | + for change in changes: |
| 68 | + if not isinstance(change, PrimitiveField): |
| 69 | + raise RuntimeError( |
| 70 | + f"{type(self).__name__}: Only primitive field changes are supported, " |
| 71 | + f"got {type(change).__name__}. This is a bug in NEAT." |
| 72 | + ) |
| 73 | + if "." not in change.field_path: |
| 74 | + raise RuntimeError( |
| 75 | + f"{type(self).__name__}: Invalid field_path '{change.field_path}' " |
| 76 | + "(expected 'field_name.identifier' format). This is a bug in NEAT." |
| 77 | + ) |
| 78 | + field_name, identifier = change.field_path.split(".", maxsplit=1) |
| 79 | + field_map = getattr(resource, field_name, None) |
| 80 | + if field_map is None: |
| 81 | + field_map = {} |
| 82 | + setattr(resource, field_name, field_map) |
| 83 | + if isinstance(change, RemovedField): |
| 84 | + field_map.pop(identifier, None) |
| 85 | + elif isinstance(change, AddedField | ChangedField): |
| 86 | + field_map[identifier] = change.new_value |
| 87 | + if not field_map: |
| 88 | + setattr(resource, field_name, None) |
| 89 | + |
| 90 | + def _check_no_field_path_conflicts(self, changes: list[FieldChange]) -> None: |
| 91 | + """Raise if any changes touch a field_path already modified by a previous change.""" |
| 92 | + seen_paths: set[str] = set() |
| 93 | + for change in changes: |
| 94 | + if change.field_path in seen_paths: |
| 95 | + raise RuntimeError( |
| 96 | + f"{type(self).__name__}: Conflicting fixes — multiple changes " |
| 97 | + f"to '{change.field_path}'. This is a bug in NEAT." |
| 98 | + ) |
| 99 | + seen_paths.add(change.field_path) |
0 commit comments