Skip to content
Draft
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
100 changes: 82 additions & 18 deletions ckanext/validate/actions/action.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,9 @@

from ckanext.validate import helpers as h
from ckanext.validate.model import Validation
from ckanext.validate.model.validation_configuration import (
ValidationConfiguration,
)
from ckanext.validate.model.validation_jobs import JobStatus, ValidationJob
from ckanext.validate.resource_hooks import is_csv_resource
from ckanext.validate.detector import ValidateDetector
Expand All @@ -15,7 +18,7 @@
log = logging.getLogger(__name__)


def get_validation_report(source, format):
def get_validation_report(source, resource_format, schema=None):
"""Main function to execute and and get the validation report.

The validation process is opinionated to reflect end users requirements
Expand All @@ -36,19 +39,40 @@ def get_validation_report(source, format):
Returns:
A Frictionless report.
"""
_detector = ValidateDetector(
field_missing_values=["null", "NULL", "None"],
field_confidence=0.5
)

detector = ValidateDetector(
field_missing_values=[
"null",
"NULL",
"None",
],
field_confidence=0.5,
)

with system.use_context(trusted=True):
res = Resource(source, format=format, detector=_detector)
res.infer()
for field in res.schema.fields:
field.constraints = {"required": True}
report = res.validate()
resource = Resource(
source,
format=resource_format,
detector=detector,
)

resource.infer()

if schema is not None:
resource.schema = h.merge_validation_schema(
resource.schema,
schema,
)
else:
for field in resource.schema.fields:
constraints = dict(
field.constraints or {}
)

constraints["required"] = True
field.constraints = constraints

return report
return resource.validate()


def resource_validate(context, data_dict):
Expand All @@ -60,6 +84,7 @@ def resource_validate(context, data_dict):
:returns: the resource dict
:rtype: dict
"""

resource_id = toolkit.get_or_bust(data_dict, "id")
toolkit.check_access("resource_update", context, {"id": resource_id})
resource = toolkit.get_action("resource_show")(context, {"id": resource_id})
Expand All @@ -69,22 +94,58 @@ def resource_validate(context, data_dict):
{"format": [toolkit._("Only CSV resources can be validated.")]}
)

configuration_id = data_dict.get("validation_configuration_id")

if configuration_id:
configuration = ValidationConfiguration.get_active(configuration_id)

if configuration is None:
raise toolkit.ValidationError(
{
"validation_configuration_id": [
toolkit._(
"The selected validation configuration "
"does not exist or is inactive."
)
]
}
)
else:
configuration = h.get_configuration_for_resource(resource)

is_uploaded = resource.get("url_type") == "upload"
fmt_lower = h.normalize_format(resource)
resource_format = h.normalize_format(resource)
if is_uploaded:
# TODO: Refactor to new file API when migrating to CKAN 2.12.
upload = uploader.get_resource_uploader(resource)
source = "file://" + upload.get_path(resource["id"])
source = "file://" + upload.get_path(resource_id)
else:
source = resource["url"]

log.info(
"Starting validation for resource %s (format=%s, uploaded=%s, source=%s)",
resource_id, fmt_lower, is_uploaded, source,
"Starting validation for resource %s "
"(format=%s, uploaded=%s, source=%s, "
"configuration_id=%s, configuration_name=%s)",
resource_id,
resource_format,
is_uploaded,
source,
configuration.id if configuration else None,
configuration.name if configuration else None,
)

try:
report = get_validation_report(source, fmt_lower)
configured_schema = (
configuration.get_schema()
if configuration
else None
)

report = get_validation_report(
source,
resource_format,
schema=configured_schema,
)
except Exception as exc:
log.exception("Frictionless raised an exception for resource %s", resource_id)
raise toolkit.ValidationError(
Expand All @@ -108,8 +169,11 @@ def resource_validate(context, data_dict):
)

log.info(
"Resource %s validation finished: status=%s errors=%d",
resource_id, status, error_count,
"Resource %s validation finished: status=%s errors=%d configuration_id=%s",
resource_id,
status,
error_count,
configuration.id if configuration else None,
)

return resource
Expand Down
5 changes: 4 additions & 1 deletion ckanext/validate/blueprints/resource.py
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,10 @@ def validate_test_file_view():
try:
uploaded_file.save(tmp_path)

report = get_validation_report("file://" + tmp_path, format="csv")
report = get_validation_report(
"file://" + tmp_path,
resource_format="csv",
)

report_valid = report.valid
errors = h.collect_report_errors(report)
Expand Down
161 changes: 161 additions & 0 deletions ckanext/validate/helpers.py
Original file line number Diff line number Diff line change
@@ -1,16 +1,25 @@
from datetime import datetime
import logging

from collections import OrderedDict

import ckan.plugins.toolkit as toolkit

from ckanext.validate.model.validation import Validation
from ckanext.validate.model.validation_configuration import ValidationConfigurationAssignment, ValidationConfiguration
from ckanext.validate.model.validation_jobs import JobStatus, ValidationJob

from copy import deepcopy

from frictionless import Schema


MAX_ERROR_ROWS_PER_GROUP = 20


log = logging.getLogger(__name__)


def collect_report_errors(report):
descriptor = report.to_descriptor() or {}
errors = []
Expand Down Expand Up @@ -311,3 +320,155 @@ def validation_error_message(error):
messages.append(str(value))

return "; ".join(messages) or str(error)


def get_configuration_for_resource(resource):
"""Return the active validation configuration for a given resource."""
resource_id = resource["id"]
package_id = resource["package_id"]

assignment = (
ValidationConfigurationAssignment
.get_for_resource(resource_id)
)

if assignment:
return ValidationConfiguration.get_active(
assignment.configuration_id
)

assignment = (
ValidationConfigurationAssignment
.get_for_package(package_id)
)

if assignment:
return ValidationConfiguration.get_active(
assignment.configuration_id
)

assignment = (
ValidationConfigurationAssignment
.get_global()
)

if assignment:
return ValidationConfiguration.get_active(
assignment.configuration_id
)

active_configurations = ValidationConfiguration.get_all(
active=True
)

if len(active_configurations) == 1:
configuration = active_configurations[0]
log.info(
"Using single active validation configuration by default: %s",
configuration.id,
)
return configuration

return None


DEFAULT_MISSING_VALUES = [
"",
"null",
"NULL",
"None",
]


def merge_validation_schema(
inferred_schema,
configured_schema,
):
"""Apply configured rules over the complete inferred schema.

The visual editor can define rules for only some CSV columns.
Therefore, the configured schema must not replace the complete inferred
schema because unrelated CSV columns would be reported as extra labels.
"""
inferred_descriptor = inferred_schema.to_descriptor()
configured_descriptor = configured_schema.to_descriptor()

configured_fields = {
field["name"]: field
for field in configured_descriptor.get(
"fields",
[],
)
}

merged_fields = []
processed_fields = set()

for inferred_field in inferred_descriptor.get(
"fields",
[],
):
field_name = inferred_field.get("name")
configured_field = configured_fields.get(
field_name
)

if configured_field is None:
merged_fields.append(
deepcopy(inferred_field)
)
continue

merged_field = deepcopy(inferred_field)

# The explicitly configured type, format and constraints
# take precedence over inferred values.
merged_field.update(
deepcopy(configured_field)
)

merged_fields.append(merged_field)
processed_fields.add(field_name)

# Keep configured fields that do not exist in the CSV.
# Frictionless can then report them as missing labels.
for field_name, configured_field in configured_fields.items():
if field_name in processed_fields:
continue

merged_fields.append(
deepcopy(configured_field)
)

merged_descriptor = deepcopy(
inferred_descriptor
)

# Preserve any supported schema-level properties configured in
# Frictionless, except fields and missingValues, which are merged below.
for key, value in configured_descriptor.items():
if key in {"fields", "missingValues"}:
continue

merged_descriptor[key] = deepcopy(value)

merged_descriptor["fields"] = merged_fields

missing_values = list(
configured_descriptor.get(
"missingValues",
[],
)
)

for value in DEFAULT_MISSING_VALUES:
if value not in missing_values:
missing_values.append(value)

merged_descriptor["missingValues"] = (
missing_values
)

return Schema.from_descriptor(
merged_descriptor
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
"""Add ValidationConfigurationAssignment

Revision ID: 0ea099fdb60f
Revises: 76437dc3696c
Create Date: 2026-07-03 16:18:17.205720

"""
from alembic import op
import sqlalchemy as sa


# revision identifiers, used by Alembic.
revision = '0ea099fdb60f'
down_revision = '76437dc3696c'
branch_labels = None
depends_on = None


def upgrade():
op.create_table(
"validate_validation_configuration_assignment",
sa.Column("id", sa.Integer(), nullable=False, autoincrement=True),
sa.Column("configuration_id", sa.UnicodeText(), nullable=False),
sa.Column("target_type", sa.UnicodeText(), nullable=False),
sa.Column("target_id", sa.UnicodeText(), nullable=False),
sa.PrimaryKeyConstraint("id"),
sa.UniqueConstraint(
"target_type",
"target_id",
name="uq_validation_configuration_target",
),
)

op.create_index(
"ix_validation_configuration_assignment_configuration",
"validate_validation_configuration_assignment",
["configuration_id"],
unique=False,
)


def downgrade():
op.drop_index(
"ix_validation_configuration_assignment_configuration",
table_name="validate_validation_configuration_assignment",
)

op.drop_table("validate_validation_configuration_assignment")
Loading
Loading