diff --git a/database/functions.sql b/database/functions.sql index 3301f7a..82e5671 100644 --- a/database/functions.sql +++ b/database/functions.sql @@ -184,118 +184,120 @@ CREATE AGGREGATE dataflow.jsonb_concat_obj(JSONB) ( DROP FUNCTION IF EXISTS apply_transformations(TEXT, INTEGER[]); CREATE OR REPLACE FUNCTION apply_transformations( p_source_name TEXT, - p_record_ids INTEGER[] DEFAULT NULL, -- NULL = all eligible records - p_overwrite BOOLEAN DEFAULT FALSE -- FALSE = skip already-transformed, TRUE = overwrite all + p_record_ids INTEGER[] DEFAULT NULL, -- NULL = all eligible records + p_overwrite BOOLEAN DEFAULT FALSE -- FALSE = skip already-transformed, TRUE = overwrite all ) RETURNS JSON AS $$ -WITH --- All records to process -qualifying AS ( - SELECT id, data +DECLARE + v_seq INT; + v_count INT := 0; +BEGIN + -- Accumulator: one row per qualifying record, additions built up across sequence steps. + -- Each sequence step reads from data || additions so later rules can reference earlier outputs. + CREATE TEMP TABLE _xform_acc ON COMMIT DROP AS + SELECT id, data, '{}'::jsonb AS additions FROM dataflow.records WHERE source_name = p_source_name AND (p_overwrite OR transformed IS NULL) - AND (p_record_ids IS NULL OR id = ANY(p_record_ids)) -), --- Mirror TPS rx: fan out one row per regex match, drive from rules → records -rx AS ( - SELECT - q.id, - r.name AS rule_name, - r.sequence, - r.output_field, - r.retain, - r.function_type, - COALESCE(mt.rn, rp.rn, 1) AS result_number, - -- extract: build map_val and retain_val per match (mirrors TPS) - CASE WHEN array_length(mt.mt, 1) = 1 THEN to_jsonb(mt.mt[1]) ELSE to_jsonb(mt.mt) END AS match_val, - to_jsonb(rp.rp) AS replace_val - FROM dataflow.rules r - INNER JOIN qualifying q ON q.data ? r.field - LEFT JOIN LATERAL regexp_matches(q.data ->> r.field, r.pattern, r.flags) - WITH ORDINALITY AS mt(mt, rn) ON r.function_type = 'extract' - LEFT JOIN LATERAL regexp_replace(q.data ->> r.field, r.pattern, r.replace_value, r.flags) - WITH ORDINALITY AS rp(rp, rn) ON r.function_type = 'replace' - WHERE r.source_name = p_source_name - AND r.enabled = true -), --- Aggregate match rows back into one value per (record, rule) — mirrors TPS agg_to_target_items -agg_matches AS ( - SELECT - id, - rule_name, - sequence, - output_field, - retain, - function_type, - CASE function_type - WHEN 'replace' THEN jsonb_agg(replace_val) -> 0 - ELSE - CASE WHEN max(result_number) = 1 - THEN jsonb_agg(match_val ORDER BY result_number) -> 0 - ELSE jsonb_agg(match_val ORDER BY result_number) - END - END AS extracted - FROM rx - GROUP BY id, rule_name, sequence, output_field, retain, function_type -), --- Join with mappings to find mapped output — mirrors TPS link_map -linked AS ( - SELECT - a.id, - a.sequence, - a.output_field, - a.retain, - a.extracted, - m.output AS mapped - FROM agg_matches a - LEFT JOIN dataflow.mappings m ON - m.source_name = p_source_name - AND m.rule_name = a.rule_name - AND m.input_value = a.extracted - WHERE a.extracted IS NOT NULL -), --- Build per-rule output JSONB: --- mapped → use mapping output; also write output_field if retain = true --- no map → write extracted value to output_field -rule_output AS ( - SELECT - id, - sequence, - CASE - WHEN mapped IS NOT NULL THEN - mapped || - CASE WHEN retain - THEN jsonb_build_object(output_field, extracted) - ELSE '{}'::jsonb - END - ELSE - jsonb_build_object(output_field, extracted) - END AS output - FROM linked -), --- Merge all rule outputs per record in sequence order — mirrors TPS agg_to_id -record_additions AS ( - SELECT - id, - dataflow.jsonb_concat_obj(output ORDER BY sequence) AS additions - FROM rule_output - GROUP BY id -), --- Update all qualifying records; records with no rule matches get transformed = data -updated AS ( - UPDATE dataflow.records rec - SET transformed = rec.data || COALESCE(ra.additions, '{}'::jsonb) || COALESCE(rec.overrides, '{}'::jsonb), - transformed_at = CURRENT_TIMESTAMP - FROM qualifying q - LEFT JOIN record_additions ra ON ra.id = q.id - WHERE rec.id = q.id - RETURNING rec.id -) -SELECT json_build_object('success', true, 'transformed', count(*)) -FROM updated -$$ LANGUAGE sql; + AND (p_record_ids IS NULL OR id = ANY(p_record_ids)); -COMMENT ON FUNCTION apply_transformations IS 'Apply transformation rules and mappings to records (set-based CTE)'; + -- Process one sequence group at a time, in order. + -- Rules at sequence N can read fields written by rules at sequence < N. + FOR v_seq IN + SELECT DISTINCT sequence + FROM dataflow.rules + WHERE source_name = p_source_name AND enabled = true + ORDER BY sequence + LOOP + WITH + -- Current view of each record: original data merged with accumulated outputs so far + current AS ( + SELECT id, data || additions AS current_data + FROM _xform_acc + ), + -- Fan out one row per regex match for rules at this sequence level + rx AS ( + SELECT + c.id, + r.name AS rule_name, + r.sequence, + r.output_field, + r.retain, + r.function_type, + COALESCE(mt.rn, rp.rn, 1) AS result_number, + CASE WHEN array_length(mt.mt, 1) = 1 THEN to_jsonb(mt.mt[1]) ELSE to_jsonb(mt.mt) END AS match_val, + to_jsonb(rp.rp) AS replace_val + FROM dataflow.rules r + INNER JOIN current c ON (c.current_data ? r.field) + LEFT JOIN LATERAL regexp_matches(c.current_data ->> r.field, r.pattern, r.flags) + WITH ORDINALITY AS mt(mt, rn) ON r.function_type = 'extract' + LEFT JOIN LATERAL regexp_replace(c.current_data ->> r.field, r.pattern, r.replace_value, r.flags) + WITH ORDINALITY AS rp(rp, rn) ON r.function_type = 'replace' + WHERE r.source_name = p_source_name + AND r.sequence = v_seq + AND r.enabled = true + ), + agg_matches AS ( + SELECT + id, rule_name, sequence, output_field, retain, function_type, + CASE function_type + WHEN 'replace' THEN jsonb_agg(replace_val) -> 0 + ELSE + CASE WHEN max(result_number) = 1 + THEN jsonb_agg(match_val ORDER BY result_number) -> 0 + ELSE jsonb_agg(match_val ORDER BY result_number) + END + END AS extracted + FROM rx + GROUP BY id, rule_name, sequence, output_field, retain, function_type + ), + linked AS ( + SELECT + a.id, a.sequence, a.output_field, a.retain, a.extracted, + m.output AS mapped + FROM agg_matches a + LEFT JOIN dataflow.mappings m ON + m.source_name = p_source_name + AND m.rule_name = a.rule_name + AND m.input_value = a.extracted + WHERE a.extracted IS NOT NULL + ), + rule_output AS ( + SELECT id, sequence, + CASE + WHEN mapped IS NOT NULL THEN + mapped || CASE WHEN retain THEN jsonb_build_object(output_field, extracted) ELSE '{}'::jsonb END + ELSE + jsonb_build_object(output_field, extracted) + END AS output + FROM linked + ), + seq_additions AS ( + SELECT id, dataflow.jsonb_concat_obj(output ORDER BY sequence) AS additions + FROM rule_output + GROUP BY id + ) + UPDATE _xform_acc acc + SET additions = additions || COALESCE(sa.additions, '{}'::jsonb) + FROM seq_additions sa + WHERE acc.id = sa.id; + END LOOP; + + -- Write final result: original data + all accumulated rule outputs + any manual overrides + WITH updated AS ( + UPDATE dataflow.records rec + SET transformed = rec.data || acc.additions || COALESCE(rec.overrides, '{}'::jsonb), + transformed_at = CURRENT_TIMESTAMP + FROM _xform_acc acc + WHERE rec.id = acc.id + RETURNING rec.id + ) + SELECT count(*) INTO v_count FROM updated; + + RETURN json_build_object('success', true, 'transformed', v_count); +END; +$$ LANGUAGE plpgsql; + +COMMENT ON FUNCTION apply_transformations IS 'Apply transformation rules and mappings to records. Rules are processed in sequence order — a rule at sequence N can read fields written by rules at sequence < N (chaining).'; ------------------------------------------------------ -- Function: get_all_values diff --git a/database/queries/rules.sql b/database/queries/rules.sql index e79e925..7cd48cb 100644 --- a/database/queries/rules.sql +++ b/database/queries/rules.sql @@ -86,21 +86,27 @@ CREATE OR REPLACE FUNCTION preview_rule( p_limit INT DEFAULT 20 ) RETURNS TABLE (id INT, raw_value TEXT, extracted_value JSONB) AS $$ +-- Field is resolved from data first, then transformed (supports chained rules whose +-- input field was produced by an earlier-sequence rule rather than the raw import). BEGIN IF p_function_type = 'replace' THEN RETURN QUERY SELECT r.id, - r.data ->> p_field, - to_jsonb(regexp_replace(r.data ->> p_field, p_pattern, p_replace_value, p_flags)) + COALESCE(r.data ->> p_field, r.transformed ->> p_field), + to_jsonb(regexp_replace( + COALESCE(r.data ->> p_field, r.transformed ->> p_field), + p_pattern, p_replace_value, p_flags + )) FROM dataflow.records r - WHERE source_name = p_source AND data ? p_field + WHERE source_name = p_source + AND (data ? p_field OR transformed ? p_field) ORDER BY r.id DESC LIMIT p_limit; ELSE RETURN QUERY SELECT r.id, - r.data ->> p_field, + COALESCE(r.data ->> p_field, r.transformed ->> p_field), CASE WHEN agg.match_count = 0 THEN NULL WHEN agg.match_count = 1 THEN agg.matches -> 0 @@ -114,10 +120,14 @@ BEGIN ORDER BY rn ) AS matches, count(*)::int AS match_count - FROM regexp_matches(r.data ->> p_field, p_pattern, p_flags) + FROM regexp_matches( + COALESCE(r.data ->> p_field, r.transformed ->> p_field), + p_pattern, p_flags + ) WITH ORDINALITY AS m(mt, rn) ) agg - WHERE r.source_name = p_source AND r.data ? p_field + WHERE r.source_name = p_source + AND (r.data ? p_field OR r.transformed ? p_field) ORDER BY r.id DESC LIMIT p_limit; END IF; END; diff --git a/ui/src/pages/Rules.jsx b/ui/src/pages/Rules.jsx index b9359a2..7b2641e 100644 --- a/ui/src/pages/Rules.jsx +++ b/ui/src/pages/Rules.jsx @@ -42,7 +42,7 @@ function PreviewModal({ rows, onClose }) { ) } -function FormPanel({ form, setForm, editing, error, loading, fields, source, onSubmit, onCancel }) { +function FormPanel({ form, setForm, editing, error, loading, fields, rules, source, onSubmit, onCancel }) { const [preview, setPreview] = useState([]) const [previewing, setPreviewing] = useState(false) const [modalOpen, setModalOpen] = useState(false) @@ -91,15 +91,28 @@ function FormPanel({ form, setForm, editing, error, loading, fields, source, onS
- {fields.length > 0 ? ( - - ) : ( + {fields.length > 0 ? (() => { + // Output fields from rules at a lower sequence — available as chained inputs + const chainedFields = [...new Set( + (rules || []) + .filter(r => r.sequence < form.sequence && r.output_field && (!editing || r.id !== editing)) + .map(r => r.output_field) + )].filter(f => !fields.includes(f)) + return ( + + ) + })() : ( setForm(f => ({ ...f, field: e.target.value }))} @@ -322,7 +335,7 @@ export default function Rules({ source }) { {creating && ( { setCreating(false); setError('') }} /> @@ -377,8 +390,8 @@ export default function Rules({ source }) {
handleSubmit(e, rule.id)} onCancel={() => { setEditing(null); setExpanded(null) }} />