From 2ea2548715abf8b7c3c095c22d8103c18650ce66 Mon Sep 17 00:00:00 2001 From: Paul Trowbridge Date: Sun, 26 Jul 2026 21:55:37 -0400 Subject: [PATCH] Flatten database/queries into database/ and fix five stale functions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit database/functions.sql held a full pg_dump appended onto the original hand-written file, duplicating ~50 functions that had since been split into database/queries/. Nothing deployed it, but CLAUDE.md and the tutorial still told you to psql it, which would have reverted the split versions. The reverse had also happened: five functions in queries/ were behind the live database, all of them undoing the May 2026 split of the transformed column. preview_rule lost its data -> transformed fallback for chained rules; set_/clear_/bulk_set_record_overrides wrote overrides back into transformed and returned the wrong type (which would have made the redeploy error outright); generate_source_view read only transformed instead of merging all three layers. Those are corrected here from the live definitions. generate_source_view additionally regains the _overridden column that queries/ had and live lacked — Records.jsx reads row._overridden to highlight manually edited rows, so that indicator had been dead. The seven functions that existed only in functions.sql move to two new files, import.sql (import + audit trail) and transform.sql (the rule/mapping engine), leaving database/ flat: schema.sql plus one file per route. The four already applied migrate_*.sql scripts are removed. manage.py picks up the new files in QUERY_FILES, and its DB_ACTIONS set now keys off the action functions rather than duplicated label strings that no longer matched any menu entry, so the "into database X" hint renders again. uninstall.sh is folded into manage.py as menu option 10. Beyond what the script did, it stops/disables/removes the systemd unit, removes the nginx site with an nginx -t check before reloading, and deletes public/. Co-Authored-By: Claude Opus 5 --- database/functions.sql | 1515 -------------------- database/import.sql | 159 ++ database/{queries => }/mappings.sql | 0 database/migrate_input_value_jsonb.sql | 22 - database/migrate_overrides_column.sql | 40 - database/migrate_pivot_layouts_drop_fk.sql | 4 - database/migrate_tps.sql | 121 -- database/{queries => }/records.sql | 48 +- database/{queries => }/rules.sql | 22 +- database/{queries => }/sources.sql | 18 +- database/{queries => }/stacks.sql | 0 database/{queries => }/status.sql | 0 database/transform.sql | 156 ++ manage.py | 170 ++- uninstall.sh | 85 -- 15 files changed, 513 insertions(+), 1847 deletions(-) delete mode 100644 database/functions.sql create mode 100644 database/import.sql rename database/{queries => }/mappings.sql (100%) delete mode 100644 database/migrate_input_value_jsonb.sql delete mode 100644 database/migrate_overrides_column.sql delete mode 100644 database/migrate_pivot_layouts_drop_fk.sql delete mode 100644 database/migrate_tps.sql rename database/{queries => }/records.sql (67%) rename database/{queries => }/rules.sql (87%) rename database/{queries => }/sources.sql (92%) rename database/{queries => }/stacks.sql (100%) rename database/{queries => }/status.sql (100%) create mode 100644 database/transform.sql delete mode 100755 uninstall.sh diff --git a/database/functions.sql b/database/functions.sql deleted file mode 100644 index 18bedde..0000000 --- a/database/functions.sql +++ /dev/null @@ -1,1515 +0,0 @@ --- --- Dataflow Functions --- Simple, clear functions for import and transformation --- - -SET search_path TO dataflow, public; - ------------------------------------------------------- --- Function: import_records --- Import data with automatic deduplication ------------------------------------------------------- -CREATE OR REPLACE FUNCTION import_records( - p_source_name TEXT, - p_data JSONB -- Array of records -) RETURNS JSON AS $$ -DECLARE - v_constraint_fields TEXT[]; - v_inserted INTEGER; - v_duplicates INTEGER; - v_log_id INTEGER; -BEGIN - SELECT constraint_fields INTO v_constraint_fields - FROM dataflow.sources - WHERE name = p_source_name; - - IF v_constraint_fields IS NULL THEN - RETURN json_build_object( - 'success', false, - 'error', 'Source not found: ' || p_source_name - ); - END IF; - - WITH - -- All incoming records with their constraint keys - pending AS ( - SELECT - rec.value AS data, - rec.ordinality AS seq, - (SELECT jsonb_object_agg(f, rec.value->>f) - FROM unnest(v_constraint_fields) AS f) AS constraint_key - FROM jsonb_array_elements(p_data) WITH ORDINALITY AS rec - ), - -- Keys already in the database (excluded) - existing AS ( - SELECT DISTINCT r.constraint_key - FROM dataflow.records r - INNER JOIN pending p ON p.constraint_key = r.constraint_key - WHERE r.source_name = p_source_name - ), - -- Rows whose constraint key is not yet in the database - new_records AS ( - SELECT p.data, p.constraint_key, p.seq - FROM pending p - WHERE NOT EXISTS (SELECT 1 FROM existing e WHERE e.constraint_key = p.constraint_key) - ), - -- Write the log entry - log_entry AS ( - INSERT INTO dataflow.import_log (source_name, records_imported, records_duplicate, info) - VALUES ( - p_source_name, - (SELECT count(*) FROM new_records), - (SELECT count(*) FROM pending) - (SELECT count(*) FROM new_records), - jsonb_build_object( - 'total', jsonb_array_length(p_data), - 'inserted_keys', (SELECT jsonb_agg(constraint_key ORDER BY constraint_key) FROM new_records), - 'excluded_keys', (SELECT jsonb_agg(constraint_key) FROM existing) - ) - ) - RETURNING id, records_imported, records_duplicate - ), - -- Insert new records - inserted AS ( - INSERT INTO dataflow.records (source_name, data, constraint_key, import_id) - SELECT p_source_name, nr.data, nr.constraint_key, (SELECT id FROM log_entry) - FROM new_records nr - ORDER BY nr.seq - RETURNING id - ) - SELECT le.id, le.records_imported, le.records_duplicate - INTO v_log_id, v_inserted, v_duplicates - FROM log_entry le; - - RETURN json_build_object( - 'success', true, - 'imported', v_inserted, - 'duplicates', v_duplicates, - 'log_id', v_log_id - ); -END; -$$ LANGUAGE plpgsql; - -COMMENT ON FUNCTION import_records IS 'Import records with automatic deduplication'; - ------------------------------------------------------- --- Function: get_import_log --- Return import history for a source ------------------------------------------------------- -CREATE OR REPLACE FUNCTION get_import_log(p_source_name TEXT) -RETURNS TABLE ( - id INTEGER, - source_name TEXT, - records_imported INTEGER, - records_duplicate INTEGER, - imported_at TIMESTAMPTZ, - info JSONB -) AS $$ - SELECT id, source_name, records_imported, records_duplicate, imported_at, info - FROM dataflow.import_log - WHERE source_name = p_source_name - ORDER BY imported_at DESC; -$$ LANGUAGE sql; - -COMMENT ON FUNCTION get_import_log IS 'Return import history for a source, newest first, including inserted/excluded key lists'; - ------------------------------------------------------- --- Function: get_all_import_logs --- Return import history across all sources ------------------------------------------------------- -CREATE OR REPLACE FUNCTION get_all_import_logs() -RETURNS TABLE ( - id INTEGER, - source_name TEXT, - records_imported INTEGER, - records_duplicate INTEGER, - imported_at TIMESTAMPTZ, - info JSONB -) AS $$ - SELECT id, source_name, records_imported, records_duplicate, imported_at, info - FROM dataflow.import_log - ORDER BY imported_at DESC; -$$ LANGUAGE sql; - -COMMENT ON FUNCTION get_all_import_logs IS 'Return import history across all sources, newest first'; - ------------------------------------------------------- --- Function: delete_import --- Delete all records from a specific import and remove the log entry ------------------------------------------------------- -CREATE OR REPLACE FUNCTION delete_import(p_log_id INTEGER) -RETURNS JSON AS $$ -DECLARE - v_deleted INTEGER; -BEGIN - IF NOT EXISTS (SELECT 1 FROM dataflow.import_log WHERE id = p_log_id) THEN - RETURN json_build_object('success', false, 'error', 'Import log entry not found'); - END IF; - - SELECT count(*) INTO v_deleted FROM dataflow.records WHERE import_id = p_log_id; - - -- Cascade handles deleting records via FK ON DELETE CASCADE - DELETE FROM dataflow.import_log WHERE id = p_log_id; - - RETURN json_build_object( - 'success', true, - 'records_deleted', v_deleted, - 'log_id', p_log_id - ); -END; -$$ LANGUAGE plpgsql; - -COMMENT ON FUNCTION delete_import IS 'Delete all records belonging to an import batch and remove the log entry'; - ------------------------------------------------------- --- Aggregate: jsonb_concat_obj --- Merge JSONB objects across rows (later rows win on key conflicts) --- Usage: jsonb_concat_obj(col ORDER BY sequence) ------------------------------------------------------- -CREATE OR REPLACE FUNCTION dataflow.jsonb_merge(a JSONB, b JSONB) -RETURNS JSONB AS $$ - SELECT COALESCE(a, '{}') || COALESCE(b, '{}') -$$ LANGUAGE sql IMMUTABLE; - -DROP AGGREGATE IF EXISTS dataflow.jsonb_concat_obj(JSONB); -CREATE AGGREGATE dataflow.jsonb_concat_obj(JSONB) ( - sfunc = dataflow.jsonb_merge, - stype = JSONB, - initcond = '{}' -); - ------------------------------------------------------- --- Function: apply_transformations --- Apply all transformation rules to records (set-based) ------------------------------------------------------- -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 -) RETURNS JSON AS $$ -WITH --- All records to process -qualifying AS ( - SELECT id, data - 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 = COALESCE(ra.additions, '{}'::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; - -COMMENT ON FUNCTION apply_transformations IS 'Apply transformation rules and mappings to records (set-based CTE)'; - ------------------------------------------------------- --- Function: get_all_values --- All extracted values (mapped + unmapped) with counts and mapping output ------------------------------------------------------- -DROP FUNCTION IF EXISTS get_all_values(TEXT, TEXT); -CREATE FUNCTION get_all_values( - p_source_name TEXT, - p_rule_name TEXT DEFAULT NULL -) RETURNS TABLE ( - rule_name TEXT, - output_field TEXT, - source_field TEXT, - extracted_value JSONB, - record_count BIGINT, - sample JSONB, - mapping_id INTEGER, - output JSONB, - is_mapped BOOLEAN -) AS $$ -BEGIN - RETURN QUERY - WITH extracted AS ( - SELECT - r.name AS rule_name, - r.output_field, - r.field AS source_field, - rec.transformed->r.output_field AS extracted_value, - rec.data AS record_data, - row_number() OVER ( - PARTITION BY r.name, rec.transformed->r.output_field - ORDER BY rec.id - ) AS rn - FROM dataflow.records rec - CROSS JOIN dataflow.rules r - WHERE - rec.source_name = p_source_name - AND r.source_name = p_source_name - AND rec.transformed IS NOT NULL - AND rec.transformed ? r.output_field - AND (p_rule_name IS NULL OR r.name = p_rule_name) - AND rec.data ? r.field - ), - aggregated AS ( - SELECT - e.rule_name, - e.output_field, - e.source_field, - e.extracted_value, - count(*) AS record_count, - jsonb_agg(e.record_data ORDER BY e.rn) FILTER (WHERE e.rn <= 5) AS sample - FROM extracted e - GROUP BY e.rule_name, e.output_field, e.source_field, e.extracted_value - ) - SELECT - a.rule_name, - a.output_field, - a.source_field, - a.extracted_value, - a.record_count, - a.sample, - m.id AS mapping_id, - m.output, - (m.id IS NOT NULL) AS is_mapped - FROM aggregated 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_value - ORDER BY a.record_count DESC; -END; -$$ LANGUAGE plpgsql; - -COMMENT ON FUNCTION get_all_values IS 'All extracted values with record counts and mapping output (single query for All tab)'; - ------------------------------------------------------- --- Function: get_unmapped_values --- Find extracted values that need mappings ------------------------------------------------------- -DROP FUNCTION IF EXISTS get_unmapped_values(TEXT, TEXT); -CREATE FUNCTION get_unmapped_values( - p_source_name TEXT, - p_rule_name TEXT DEFAULT NULL -) RETURNS TABLE ( - rule_name TEXT, - output_field TEXT, - source_field TEXT, - extracted_value JSONB, - record_count BIGINT, - sample JSONB -) AS $$ -BEGIN - RETURN QUERY - WITH extracted AS ( - SELECT - r.name AS rule_name, - r.output_field, - r.field AS source_field, - rec.transformed->r.output_field AS extracted_value, - rec.data AS record_data, - row_number() OVER ( - PARTITION BY r.name, rec.transformed->r.output_field - ORDER BY rec.id - ) AS rn - FROM - dataflow.records rec - CROSS JOIN dataflow.rules r - WHERE - rec.source_name = p_source_name - AND r.source_name = p_source_name - AND rec.transformed IS NOT NULL - AND rec.transformed ? r.output_field - AND (p_rule_name IS NULL OR r.name = p_rule_name) - AND rec.data ? r.field - ) - SELECT - e.rule_name, - e.output_field, - e.source_field, - e.extracted_value, - count(*) AS record_count, - jsonb_agg(e.record_data ORDER BY e.rn) FILTER (WHERE e.rn <= 5) AS sample - FROM extracted e - WHERE NOT EXISTS ( - SELECT 1 FROM dataflow.mappings m - WHERE m.source_name = p_source_name - AND m.rule_name = e.rule_name - AND m.input_value = e.extracted_value - ) - GROUP BY e.rule_name, e.output_field, e.source_field, e.extracted_value - ORDER BY count(*) DESC; -END; -$$ LANGUAGE plpgsql; - -COMMENT ON FUNCTION get_unmapped_values IS 'Find extracted values that need mappings defined'; - ------------------------------------------------------- --- Function: reprocess_records --- Clear and reapply transformations ------------------------------------------------------- -CREATE OR REPLACE FUNCTION reprocess_records(p_source_name TEXT) -RETURNS JSON AS $$ - -- Overwrite all records directly — no clear step, mirrors TPS srce_map_overwrite - SELECT dataflow.apply_transformations(p_source_name, NULL, TRUE) -$$ LANGUAGE sql; - -COMMENT ON FUNCTION reprocess_records IS 'Clear and reapply all transformations for a source'; - ------------------------------------------------------- --- Function: generate_source_view --- Build a typed flat view in dfv schema ------------------------------------------------------- -CREATE OR REPLACE FUNCTION generate_source_view(p_source_name TEXT) -RETURNS JSON AS $$ -DECLARE - v_config JSONB; - v_fields JSONB; - v_field JSONB; - v_cols TEXT := ''; - v_sql TEXT; - v_view TEXT; -BEGIN - SELECT config INTO v_config - FROM dataflow.sources - WHERE name = p_source_name; - - IF v_config IS NULL OR NOT (v_config ? 'fields') OR jsonb_array_length(v_config->'fields') = 0 THEN - RETURN json_build_object('success', false, 'error', 'No schema fields defined for this source'); - END IF; - - v_fields := v_config->'fields'; - - FOR v_field IN SELECT * FROM jsonb_array_elements(v_fields) - LOOP - IF v_cols != '' THEN v_cols := v_cols || ', '; END IF; - - IF v_field->>'expression' IS NOT NULL THEN - DECLARE - v_expr TEXT := v_field->>'expression'; - v_ref TEXT; - v_cast TEXT := COALESCE(NULLIF(v_field->>'type', ''), 'numeric'); - BEGIN - WHILE v_expr ~ '\{[^}]+\}' LOOP - v_ref := substring(v_expr FROM '\{([^}]+)\}'); - v_expr := replace(v_expr, '{' || v_ref || '}', - format('(r->>%L)::numeric', v_ref)); - END LOOP; - v_cols := v_cols || format('%s AS %I', v_expr, v_field->>'name'); - END; - ELSE - CASE v_field->>'type' - WHEN 'date' THEN - v_cols := v_cols || format('(r->>%L)::date AS %I', - v_field->>'name', v_field->>'name'); - WHEN 'numeric' THEN - v_cols := v_cols || format('(r->>%L)::numeric AS %I', - v_field->>'name', v_field->>'name'); - ELSE - v_cols := v_cols || format('r->>%L AS %I', - v_field->>'name', v_field->>'name'); - END CASE; - END IF; - END LOOP; - - CREATE SCHEMA IF NOT EXISTS dfv; - - v_view := 'dfv.' || quote_ident(p_source_name); - - EXECUTE format('DROP VIEW IF EXISTS %s CASCADE', v_view); - - v_sql := format( - 'CREATE VIEW %s AS SELECT id, %s FROM (SELECT id, data || COALESCE(transformed, ''{}''::jsonb) || COALESCE(overrides, ''{}''::jsonb) AS r FROM dataflow.records WHERE source_name = %L AND transformed IS NOT NULL) rec', - v_view, v_cols, p_source_name - ); - - EXECUTE v_sql; - - RETURN json_build_object('success', true, 'view', v_view, 'sql', v_sql); -END; -$$ LANGUAGE plpgsql; - -COMMENT ON FUNCTION generate_source_view IS 'Generate a typed flat view in dfv schema from source config.fields'; - ------------------------------------------------------- --- Function: set_record_overrides --- Save override values for a single record ------------------------------------------------------- -DROP FUNCTION IF EXISTS set_record_overrides(INTEGER, JSONB); -CREATE OR REPLACE FUNCTION set_record_overrides(p_id INTEGER, p_overrides JSONB) -RETURNS JSON AS $$ - WITH updated AS ( - UPDATE dataflow.records - SET overrides = CASE WHEN p_overrides = '{}'::jsonb THEN NULL ELSE p_overrides END - WHERE id = p_id - RETURNING * - ) - SELECT row_to_json(updated) FROM updated; -$$ LANGUAGE sql; - ------------------------------------------------------- --- Function: clear_record_overrides --- Remove all overrides for a single record ------------------------------------------------------- -DROP FUNCTION IF EXISTS clear_record_overrides(INTEGER); -CREATE OR REPLACE FUNCTION clear_record_overrides(p_id INTEGER) -RETURNS JSON AS $$ - WITH updated AS ( - UPDATE dataflow.records - SET overrides = NULL - WHERE id = p_id - RETURNING * - ) - SELECT row_to_json(updated) FROM updated; -$$ LANGUAGE sql; - ------------------------------------------------------- --- Function: bulk_set_record_overrides --- Apply override values to multiple records ------------------------------------------------------- -DROP FUNCTION IF EXISTS bulk_set_record_overrides(TEXT, INTEGER[], JSONB); -CREATE OR REPLACE FUNCTION bulk_set_record_overrides(p_source_name TEXT, p_ids INTEGER[], p_overrides JSONB) -RETURNS JSON AS $$ - WITH updated AS ( - UPDATE dataflow.records - SET overrides = COALESCE(overrides, '{}'::jsonb) || p_overrides - WHERE id = ANY(p_ids) - AND source_name = p_source_name - RETURNING id - ) - SELECT json_build_object('updated', count(*)) FROM updated; -$$ LANGUAGE sql; - -CREATE OR REPLACE FUNCTION dataflow.calibrate_balance(p_stack_name text, p_source_name text, p_as_of_date date, p_known_balance numeric) - RETURNS json - LANGUAGE plpgsql - STABLE -AS $function$ -DECLARE - v_src dataflow.stack_sources%ROWTYPE; - v_running NUMERIC; - v_sql TEXT; -BEGIN - SELECT * INTO v_src - FROM dataflow.stack_sources - WHERE stack_name = p_stack_name AND source_name = p_source_name; - - IF NOT FOUND THEN - RETURN json_build_object('success', false, 'error', 'Source not in stack'); - END IF; - IF v_src.amount_field IS NULL OR v_src.date_field IS NULL THEN - RETURN json_build_object('success', false, 'error', 'Set amount and date fields on this source first'); - END IF; - - BEGIN - IF p_as_of_date IS NULL THEN - v_sql := format( - 'SELECT COALESCE(SUM(%I * %s), 0) FROM dfv.%I', - v_src.amount_field, v_src.amount_sign, p_source_name - ); - ELSE - v_sql := format( - 'SELECT COALESCE(SUM(%I * %s), 0) FROM dfv.%I WHERE %I <= %L::date', - v_src.amount_field, v_src.amount_sign, p_source_name, v_src.date_field, p_as_of_date - ); - END IF; - EXECUTE v_sql INTO v_running; - EXCEPTION WHEN undefined_table THEN - RETURN json_build_object('success', false, 'error', 'Source view not found — generate the source view first'); - END; - - RETURN json_build_object( - 'success', true, - 'source', p_source_name, - 'as_of_date', p_as_of_date, - 'known_balance', p_known_balance, - 'computed_sum', v_running, - 'suggested_offset', p_known_balance - v_running - ); -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.create_mapping(p_source_name text, p_rule_name text, p_input_value jsonb, p_output jsonb) - RETURNS dataflow.mappings - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.mappings (source_name, rule_name, input_value, output) - VALUES (p_source_name, p_rule_name, p_input_value, p_output) - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.create_rule(p_source_name text, p_name text, p_field text, p_pattern text, p_output_field text, p_function_type text DEFAULT 'extract'::text, p_flags text DEFAULT ''::text, p_replace_value text DEFAULT ''::text, p_enabled boolean DEFAULT true, p_retain boolean DEFAULT false, p_sequence integer DEFAULT 0) - RETURNS dataflow.rules - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.rules - (source_name, name, field, pattern, output_field, function_type, flags, replace_value, enabled, retain, sequence) - VALUES - (p_source_name, p_name, p_field, p_pattern, p_output_field, p_function_type, p_flags, p_replace_value, p_enabled, p_retain, p_sequence) - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.create_source(p_name text, p_constraint_fields text[], p_config jsonb DEFAULT '{}'::jsonb, p_global_picklist boolean DEFAULT true) - RETURNS dataflow.sources - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.sources (name, constraint_fields, config, global_picklist) - VALUES (p_name, p_constraint_fields, p_config, p_global_picklist) - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.create_source(p_name text, p_constraint_fields text[], p_config jsonb DEFAULT '{}'::jsonb) - RETURNS dataflow.sources - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.sources (name, constraint_fields, config) - VALUES (p_name, p_constraint_fields, p_config) - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.create_stack(p_name text, p_label text DEFAULT NULL::text, p_fields jsonb DEFAULT '[]'::jsonb, p_amount_field text DEFAULT NULL::text, p_date_field text DEFAULT NULL::text, p_balance_offset numeric DEFAULT 0) - RETURNS dataflow.stacks - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.stacks (name, label, fields, amount_field, date_field, balance_offset) - VALUES (p_name, p_label, p_fields, p_amount_field, p_date_field, p_balance_offset) - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.delete_mapping(p_id integer) - RETURNS TABLE(id integer) - LANGUAGE sql -AS $function$ - DELETE FROM dataflow.mappings WHERE id = p_id RETURNING id; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.delete_pivot_layout(p_id integer) - RETURNS TABLE(id integer) - LANGUAGE sql -AS $function$ - DELETE FROM dataflow.pivot_layouts WHERE id = p_id RETURNING id; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.delete_record(p_id bigint) - RETURNS TABLE(id bigint) - LANGUAGE sql -AS $function$ - DELETE FROM dataflow.records WHERE id = p_id RETURNING id; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.delete_rule(p_id integer) - RETURNS TABLE(id integer, name text) - LANGUAGE sql -AS $function$ - DELETE FROM dataflow.rules WHERE id = p_id RETURNING id, name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.delete_source(p_name text) - RETURNS text - LANGUAGE sql -AS $function$ - DELETE FROM dataflow.sources WHERE name = p_name RETURNING name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.delete_source_records(p_source_name text) - RETURNS TABLE(deleted_count bigint) - LANGUAGE sql -AS $function$ - WITH deleted AS ( - DELETE FROM dataflow.records WHERE source_name = p_source_name RETURNING id - ) - SELECT count(*) AS deleted_count FROM deleted; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.delete_stack(p_name text) - RETURNS TABLE(name text) - LANGUAGE sql -AS $function$ - DELETE FROM dataflow.stacks WHERE name = p_name RETURNING name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.generate_stack_view(p_stack_name text, p_dry_run boolean DEFAULT false) - RETURNS json - LANGUAGE plpgsql -AS $function$ -DECLARE - v_stack dataflow.stacks%ROWTYPE; - v_src dataflow.stack_sources%ROWTYPE; - v_field JSONB; - v_ctes TEXT[] := '{}'; - v_cte_names TEXT[] := '{}'; - v_select TEXT; - v_col TEXT; - v_src_field TEXT; - v_amt_src TEXT; - v_date_src TEXT; - v_view TEXT; - v_sql TEXT; - v_has_bal BOOLEAN; - v_canon_cols TEXT; - v_src_bal_cols TEXT; - v_total_offset NUMERIC := 0; - v_cascade_stale TEXT[]; -BEGIN - SELECT * INTO v_stack FROM dataflow.stacks WHERE name = p_stack_name; - IF NOT FOUND THEN - RETURN json_build_object('success', false, 'error', 'Stack not found'); - END IF; - - v_has_bal := v_stack.amount_field IS NOT NULL AND v_stack.date_field IS NOT NULL; - - -- Build one CTE per source querying dfv.{source} directly - FOR v_src IN - SELECT * FROM dataflow.stack_sources WHERE stack_name = p_stack_name ORDER BY seq, id - LOOP - v_select := format('SELECT %L AS _source, id AS _id', v_src.source_name); - - FOR v_field IN SELECT * FROM jsonb_array_elements(v_stack.fields) - LOOP - v_col := v_field->>'name'; - - IF v_has_bal AND v_col = v_stack.amount_field THEN - -- Use per-source amount_field with sign applied - IF v_src.amount_field IS NULL THEN - v_select := v_select || format(', NULL::%s AS %I', v_field->>'type', v_col); - ELSE - v_select := v_select || format(', %I * %s AS %I', v_src.amount_field, v_src.amount_sign, v_col); - END IF; - ELSIF v_has_bal AND v_col = v_stack.date_field THEN - -- Use per-source date_field - IF v_src.date_field IS NULL THEN - v_select := v_select || format(', NULL::date AS %I', v_col); - ELSE - v_select := v_select || format(', %I AS %I', v_src.date_field, v_col); - END IF; - ELSE - -- Other canonical fields: use field_map or same name, NULL if column doesn't exist - v_src_field := COALESCE(v_src.field_map->>v_col, v_col); - IF EXISTS ( - SELECT 1 FROM information_schema.columns - WHERE table_schema = 'dfv' - AND table_name = v_src.source_name - AND column_name = v_src_field - ) THEN - v_select := v_select || format(', %I AS %I', v_src_field, v_col); - ELSE - v_select := v_select || format(', NULL::text AS %I', v_col); - END IF; - END IF; - END LOOP; - - v_select := v_select || format(' FROM dfv.%I', v_src.source_name); - - v_ctes := v_ctes || format('%I AS (%s)', v_src.source_name, v_select); - v_cte_names := v_cte_names || quote_ident(v_src.source_name); - - -- Accumulate carried-forward source balance column and total offset - IF v_has_bal THEN - IF v_src_bal_cols IS NOT NULL THEN v_src_bal_cols := v_src_bal_cols || ', '; END IF; - v_src_bal_cols := COALESCE(v_src_bal_cols, '') || format( - 'SUM(CASE WHEN _source = %L THEN %I END) OVER (ORDER BY %I ASC, _id ASC) + %s AS %I', - v_src.source_name, v_stack.amount_field, v_stack.date_field, - v_src.balance_offset, v_src.source_name || '_balance' - ); - v_total_offset := v_total_offset + v_src.balance_offset; - END IF; - END LOOP; - - IF array_length(v_ctes, 1) IS NULL THEN - RETURN json_build_object('success', false, 'error', 'Stack has no sources'); - END IF; - - v_view := 'dfv.' || quote_ident(p_stack_name); - - v_canon_cols := ( - SELECT string_agg(quote_ident(f->>'name'), ', ') - FROM jsonb_array_elements(v_stack.fields) f - ); - - IF v_has_bal THEN - v_sql := format( - 'CREATE VIEW %s AS ' - 'WITH %s, _stacked AS (SELECT * FROM %s) ' - 'SELECT _source, _id, %s, ' - '%s, ' - 'SUM(%I) OVER (ORDER BY %I ASC, _id ASC) + %s AS net_balance ' - 'FROM _stacked ORDER BY %I DESC, _id DESC', - v_view, - array_to_string(v_ctes, ', '), - array_to_string(v_cte_names, ' UNION ALL SELECT * FROM '), - v_canon_cols, - v_src_bal_cols, - v_stack.amount_field, - v_stack.date_field, - v_total_offset, - v_stack.date_field - ); - ELSE - v_sql := format( - 'CREATE VIEW %s AS ' - 'WITH %s, _stacked AS (SELECT * FROM %s) ' - 'SELECT _source, _id, %s FROM _stacked', - v_view, - array_to_string(v_ctes, ', '), - array_to_string(v_cte_names, ' UNION ALL SELECT * FROM '), - v_canon_cols - ); - END IF; - - IF NOT p_dry_run THEN - CREATE SCHEMA IF NOT EXISTS dfv; - EXECUTE format('DROP VIEW IF EXISTS %s CASCADE', v_view); - EXECUTE v_sql; - - -- Detect stacks whose views were dropped by CASCADE and mark them stale - SELECT array_agg(s.name) INTO v_cascade_stale - FROM dataflow.stacks s - WHERE s.name != p_stack_name - AND s.view_generated_at IS NOT NULL - AND NOT EXISTS ( - SELECT 1 FROM pg_views v - WHERE v.schemaname = 'dfv' AND v.viewname = s.name - ); - - UPDATE dataflow.stacks SET view_generated_at = NULL - WHERE name = ANY(v_cascade_stale); - END IF; - - RETURN json_build_object( - 'success', true, - 'view', v_view, - 'sql', v_sql, - 'cascade_stale', COALESCE(to_json(v_cascade_stale), '[]'::json) - ); -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_global_output_values() - RETURNS TABLE(col text, val text) - LANGUAGE sql - STABLE -AS $function$ - SELECT DISTINCT e.key AS col, e.value AS val - FROM dataflow.mappings m - JOIN dataflow.sources s ON s.name = m.source_name - CROSS JOIN LATERAL jsonb_each_text(m.output) AS e(key, value) - WHERE s.global_picklist = true - AND e.value IS NOT NULL - AND e.value <> '' - ORDER BY e.key, e.value; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_mapping(p_id integer) - RETURNS dataflow.mappings - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.mappings WHERE id = p_id; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_mapping_counts(p_source_name text, p_rule_name text DEFAULT NULL::text) - RETURNS TABLE(rule_name text, input_value jsonb, record_count bigint) - LANGUAGE sql - STABLE -AS $function$ - SELECT - m.rule_name, - m.input_value, - COUNT(rec.id) AS record_count - FROM dataflow.mappings m - JOIN dataflow.rules r ON r.source_name = m.source_name AND r.name = m.rule_name - LEFT JOIN dataflow.records rec ON - rec.source_name = m.source_name - AND rec.transformed ? r.output_field - AND rec.transformed -> r.output_field = m.input_value - WHERE m.source_name = p_source_name - AND (p_rule_name IS NULL OR m.rule_name = p_rule_name) - GROUP BY m.rule_name, m.input_value; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_mappings_by_output_field(p_col text, p_val text) - RETURNS TABLE(id integer, source_name text, rule_name text, input_value jsonb, output jsonb) - LANGUAGE sql - STABLE -AS $function$ - SELECT m.id, m.source_name, m.rule_name, m.input_value, m.output - FROM dataflow.mappings m - WHERE m.output->>(p_col) = p_val - ORDER BY m.source_name, m.rule_name, m.input_value::text; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_record(p_id bigint) - RETURNS dataflow.records - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.records WHERE id = p_id; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_rule(p_id integer) - RETURNS dataflow.rules - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.rules WHERE id = p_id; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_source(p_name text) - RETURNS dataflow.sources - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.sources WHERE name = p_name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_source_fields(p_source_name text) - RETURNS TABLE(key text, origins text[]) - LANGUAGE sql - STABLE -AS $function$ - SELECT key, array_agg(DISTINCT origin ORDER BY origin) AS origins - FROM ( - SELECT f->>'name' AS key, 'schema' AS origin - FROM dataflow.sources, jsonb_array_elements(config->'fields') f - WHERE name = p_source_name AND config ? 'fields' - UNION ALL - SELECT jsonb_object_keys(data) AS key, 'raw' AS origin - FROM dataflow.records WHERE source_name = p_source_name - UNION ALL - SELECT output_field AS key, 'rule: ' || name AS origin - FROM dataflow.rules WHERE source_name = p_source_name - UNION ALL - SELECT jsonb_object_keys(output) AS key, 'mapping' AS origin - FROM dataflow.mappings WHERE source_name = p_source_name - ) keys - GROUP BY key - ORDER BY key; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_source_stats(p_source_name text) - RETURNS TABLE(total_records bigint, transformed_records bigint, pending_records bigint) - LANGUAGE sql - STABLE -AS $function$ - SELECT - COUNT(*) AS total_records, - COUNT(*) FILTER (WHERE transformed IS NOT NULL) AS transformed_records, - COUNT(*) FILTER (WHERE transformed IS NULL) AS pending_records - FROM dataflow.records - WHERE source_name = p_source_name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_stack(p_name text) - RETURNS TABLE(name text, label text, fields jsonb, amount_field text, date_field text, balance_offset numeric, created_at timestamp with time zone, sources jsonb) - LANGUAGE sql - STABLE -AS $function$ - SELECT - s.name, s.label, s.fields, - s.amount_field, s.date_field, s.balance_offset, - s.created_at, - COALESCE(jsonb_agg( - jsonb_build_object( - 'id', ss.id, - 'source_name', ss.source_name, - 'field_map', ss.field_map, - 'amount_field', ss.amount_field, - 'amount_sign', ss.amount_sign, - 'date_field', ss.date_field, - 'balance_offset', ss.balance_offset, - 'seq', ss.seq - ) ORDER BY ss.seq, ss.id - ) FILTER (WHERE ss.id IS NOT NULL), '[]') - FROM dataflow.stacks s - LEFT JOIN dataflow.stack_sources ss ON ss.stack_name = s.name - WHERE s.name = p_name - GROUP BY s.name, s.label, s.fields, s.amount_field, s.date_field, s.balance_offset, s.created_at; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_stack_balance(p_stack_name text) - RETURNS json - LANGUAGE plpgsql - STABLE -AS $function$ -DECLARE - v_stack dataflow.stacks%ROWTYPE; - v_balance NUMERIC; - v_view TEXT; - v_sql TEXT; -BEGIN - SELECT * INTO v_stack FROM dataflow.stacks WHERE name = p_stack_name; - IF NOT FOUND THEN - RETURN json_build_object('success', false, 'error', 'Stack not found'); - END IF; - IF v_stack.amount_field IS NULL OR v_stack.date_field IS NULL THEN - RETURN json_build_object('success', false, 'error', 'amount_field and date_field must be set'); - END IF; - - v_view := 'dfv.' || quote_ident(p_stack_name); - - BEGIN - v_sql := format( - 'SELECT net_balance FROM %s ORDER BY %I DESC, _id DESC LIMIT 1', - v_view, v_stack.date_field - ); - EXECUTE v_sql INTO v_balance; - EXCEPTION WHEN undefined_table THEN - RETURN json_build_object('success', false, 'error', 'View not generated yet — click Generate first'); - END; - - RETURN json_build_object('success', true, 'balance', v_balance); -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_status() - RETURNS json - LANGUAGE plpgsql - STABLE -AS $function$ -DECLARE - v_sources JSON; - v_stacks JSON; -BEGIN - SELECT COALESCE(json_agg(json_build_object('name', name, 'view_generated_at', view_generated_at) ORDER BY name), '[]'::json) - INTO v_sources - FROM dataflow.sources - WHERE view_generated_at IS NULL; - - SELECT COALESCE(json_agg(json_build_object('name', name, 'view_generated_at', view_generated_at) ORDER BY name), '[]'::json) - INTO v_stacks - FROM dataflow.stacks - WHERE view_generated_at IS NULL; - - RETURN json_build_object('stale_sources', v_sources, 'stale_stacks', v_stacks); -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.get_view_data(p_source_name text, p_limit integer DEFAULT 100, p_offset integer DEFAULT 0, p_sort_col text DEFAULT NULL::text, p_sort_dir text DEFAULT 'asc'::text, p_filters jsonb DEFAULT NULL::jsonb) - RETURNS json - LANGUAGE plpgsql - STABLE -AS $function$ -DECLARE - v_exists BOOLEAN; - v_where TEXT := ''; - v_order TEXT := ''; - v_rows JSON; - v_filter JSONB; - v_col TEXT; - v_pattern TEXT; -BEGIN - SELECT EXISTS ( - SELECT 1 FROM information_schema.views - WHERE table_schema = 'dfv' AND table_name = p_source_name - ) INTO v_exists; - - IF NOT v_exists THEN - RETURN json_build_object('exists', FALSE, 'rows', '[]'::json); - END IF; - - -- Build WHERE from filters (validate each column exists in the view) - IF p_filters IS NOT NULL THEN - FOR v_filter IN SELECT value FROM jsonb_array_elements(p_filters) LOOP - v_col := v_filter->>'col'; - v_pattern := v_filter->>'pattern'; - IF v_pattern IS NOT NULL AND v_pattern <> '' AND EXISTS ( - SELECT 1 FROM information_schema.columns - WHERE table_schema = 'dfv' - AND table_name = p_source_name - AND column_name = v_col - ) THEN - v_where := v_where || - CASE WHEN v_where = '' THEN ' WHERE ' ELSE ' AND ' END || - quote_ident(v_col) || '::text ~* ' || quote_literal(v_pattern); - END IF; - END LOOP; - END IF; - - IF p_sort_col IS NOT NULL AND EXISTS ( - SELECT 1 FROM information_schema.columns - WHERE table_schema = 'dfv' - AND table_name = p_source_name - AND column_name = p_sort_col - ) THEN - v_order := ' ORDER BY ' || quote_ident(p_sort_col) - || CASE WHEN lower(p_sort_dir) = 'desc' THEN ' DESC' ELSE ' ASC' END - || ' NULLS LAST'; - END IF; - - EXECUTE format( - 'SELECT COALESCE(json_agg(row_to_json(t)), ''[]''::json) FROM (SELECT * FROM dfv.%I%s%s LIMIT %s OFFSET %s) t', - p_source_name, v_where, v_order, p_limit, p_offset - ) INTO v_rows; - - RETURN json_build_object('exists', TRUE, 'rows', v_rows); -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.list_mappings(p_source_name text, p_rule_name text DEFAULT NULL::text) - RETURNS SETOF dataflow.mappings - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.mappings - WHERE source_name = p_source_name - AND (p_rule_name IS NULL OR rule_name = p_rule_name) - ORDER BY rule_name, input_value::text; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.list_pivot_layouts(p_source_name text) - RETURNS TABLE(id integer, source_name text, layout_name text, config jsonb, created_at timestamp with time zone) - LANGUAGE sql -AS $function$ - SELECT id, source_name, layout_name, config, created_at - FROM dataflow.pivot_layouts - WHERE source_name = p_source_name - ORDER BY layout_name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.list_records(p_source_name text, p_limit integer DEFAULT 100, p_offset integer DEFAULT 0, p_transformed_only boolean DEFAULT false) - RETURNS SETOF dataflow.records - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.records - WHERE source_name = p_source_name - AND (NOT p_transformed_only OR transformed IS NOT NULL) - ORDER BY id DESC - LIMIT p_limit OFFSET p_offset; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.list_rules(p_source_name text) - RETURNS SETOF dataflow.rules - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.rules - WHERE source_name = p_source_name - ORDER BY sequence, name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.list_sources() - RETURNS SETOF dataflow.sources - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.sources ORDER BY name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.list_stacks() - RETURNS TABLE(name text, label text, fields jsonb, amount_field text, date_field text, balance_offset numeric, source_count bigint, created_at timestamp with time zone) - LANGUAGE sql - STABLE -AS $function$ - SELECT - s.name, s.label, s.fields, - s.amount_field, s.date_field, s.balance_offset, - count(ss.id) AS source_count, - s.created_at - FROM dataflow.stacks s - LEFT JOIN dataflow.stack_sources ss ON ss.stack_name = s.name - GROUP BY s.name, s.label, s.fields, s.amount_field, s.date_field, s.balance_offset, s.created_at - ORDER BY s.name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.preview_rule(p_source text, p_field text, p_pattern text, p_flags text DEFAULT ''::text, p_function_type text DEFAULT 'extract'::text, p_replace_value text DEFAULT ''::text, p_limit integer DEFAULT 20) - RETURNS TABLE(id integer, raw_value text, extracted_value jsonb) - LANGUAGE plpgsql - STABLE -AS $function$ --- 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, - 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 OR transformed ? p_field) - ORDER BY r.id DESC LIMIT p_limit; - ELSE - RETURN QUERY - SELECT - r.id, - 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 - ELSE agg.matches - END - FROM dataflow.records r - CROSS JOIN LATERAL ( - SELECT - jsonb_agg( - CASE WHEN array_length(mt, 1) = 1 THEN to_jsonb(mt[1]) ELSE to_jsonb(mt) END - ORDER BY rn - ) AS matches, - count(*)::int AS match_count - 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 OR r.transformed ? p_field) - ORDER BY r.id DESC LIMIT p_limit; - END IF; -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.remap_output_field(p_col text, p_from_val text, p_to_val text) - RETURNS integer - LANGUAGE plpgsql -AS $function$ -DECLARE - updated_count INTEGER; -BEGIN - UPDATE dataflow.mappings - SET output = jsonb_set(output, ARRAY[p_col], to_jsonb(p_to_val)) - WHERE output->>(p_col) = p_from_val; - GET DIAGNOSTICS updated_count = ROW_COUNT; - RETURN updated_count; -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.remove_stack_source(p_stack_name text, p_source_name text) - RETURNS TABLE(source_name text) - LANGUAGE sql -AS $function$ - DELETE FROM dataflow.stack_sources - WHERE stack_name = p_stack_name AND source_name = p_source_name - RETURNING source_name; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.reorder_stack_sources(p_stack_name text, p_source_names text[]) - RETURNS void - LANGUAGE plpgsql -AS $function$ -DECLARE - i INTEGER; -BEGIN - FOR i IN 1..array_length(p_source_names, 1) LOOP - UPDATE dataflow.stack_sources - SET seq = i - WHERE stack_name = p_stack_name AND source_name = p_source_names[i]; - END LOOP; -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.save_pivot_layout(p_source_name text, p_layout_name text, p_config jsonb) - RETURNS TABLE(id integer, source_name text, layout_name text, config jsonb, created_at timestamp with time zone) - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.pivot_layouts (source_name, layout_name, config) - VALUES (p_source_name, p_layout_name, p_config) - ON CONFLICT (source_name, layout_name) DO UPDATE - SET config = EXCLUDED.config - RETURNING id, source_name, layout_name, config, created_at; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.search_mapping_outputs(p_search text) - RETURNS TABLE(col text, val text, mapping_count bigint) - LANGUAGE sql - STABLE -AS $function$ - SELECT e.key AS col, e.value AS val, COUNT(*) AS mapping_count - FROM dataflow.mappings m - CROSS JOIN LATERAL jsonb_each_text(m.output) AS e(key, value) - WHERE e.value ILIKE '%' || p_search || '%' - AND e.value IS NOT NULL - AND e.value <> '' - GROUP BY e.key, e.value - ORDER BY e.key, e.value; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.search_records(p_source_name text, p_query jsonb, p_limit integer DEFAULT 100) - RETURNS SETOF dataflow.records - LANGUAGE sql - STABLE -AS $function$ - SELECT * FROM dataflow.records - WHERE source_name = p_source_name - AND (data @> p_query OR transformed @> p_query) - ORDER BY id DESC - LIMIT p_limit; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.source_config_changed() - RETURNS trigger - LANGUAGE plpgsql -AS $function$ -BEGIN - IF NEW.config IS DISTINCT FROM OLD.config THEN - NEW.view_generated_at := NULL; - END IF; - RETURN NEW; -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.stack_sources_changed() - RETURNS trigger - LANGUAGE plpgsql -AS $function$ -BEGIN - IF TG_OP = 'UPDATE' THEN - IF NEW.field_map IS NOT DISTINCT FROM OLD.field_map AND - NEW.amount_sign IS NOT DISTINCT FROM OLD.amount_sign AND - NEW.balance_offset IS NOT DISTINCT FROM OLD.balance_offset AND - NEW.amount_field IS NOT DISTINCT FROM OLD.amount_field AND - NEW.date_field IS NOT DISTINCT FROM OLD.date_field AND - NEW.seq IS NOT DISTINCT FROM OLD.seq THEN - RETURN NULL; - END IF; - END IF; - UPDATE dataflow.stacks SET view_generated_at = NULL - WHERE name = COALESCE(NEW.stack_name, OLD.stack_name); - RETURN NULL; -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.test_rule(p_rule_id integer, p_limit integer DEFAULT 20) - RETURNS TABLE(rule jsonb, results jsonb) - LANGUAGE plpgsql - STABLE -AS $function$ -DECLARE - v_rule dataflow.rules%ROWTYPE; - v_results JSONB; -BEGIN - SELECT * INTO v_rule FROM dataflow.rules WHERE id = p_rule_id; - IF NOT FOUND THEN RETURN; END IF; - - SELECT jsonb_agg(row_to_json(t)) INTO v_results FROM ( - SELECT - r.id, - r.data ->> v_rule.field AS raw_value, - CASE - WHEN agg.match_count = 0 THEN NULL - WHEN agg.match_count = 1 AND array_length(agg.matches[1], 1) = 1 - THEN to_jsonb(agg.matches[1][1]) - WHEN agg.match_count = 1 - THEN to_jsonb(agg.matches[1]) - WHEN array_length(agg.matches[1], 1) = 1 - THEN (SELECT jsonb_agg(m[1] ORDER BY idx) FROM unnest(agg.matches) WITH ORDINALITY u(m, idx)) - ELSE to_jsonb(agg.matches) - END AS extracted_value - FROM dataflow.records r - CROSS JOIN LATERAL ( - SELECT array_agg(mt ORDER BY rn) AS matches, count(*)::int AS match_count - FROM regexp_matches(r.data ->> v_rule.field, v_rule.pattern, COALESCE(v_rule.flags, '')) - WITH ORDINALITY AS m(mt, rn) - ) agg - WHERE r.source_name = v_rule.source_name AND r.data ? v_rule.field - ORDER BY r.id DESC LIMIT p_limit - ) t; - - RETURN QUERY SELECT - jsonb_build_object( - 'id', v_rule.id, 'name', v_rule.name, - 'field', v_rule.field, 'pattern', v_rule.pattern, - 'output_field', v_rule.output_field - ), - COALESCE(v_results, '[]'::jsonb); -END; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.update_mapping(p_id integer, p_input_value jsonb DEFAULT NULL::jsonb, p_output jsonb DEFAULT NULL::jsonb) - RETURNS dataflow.mappings - LANGUAGE sql -AS $function$ - UPDATE dataflow.mappings SET - input_value = COALESCE(p_input_value, input_value), - output = COALESCE(p_output, output) - WHERE id = p_id - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.update_rule(p_id integer, p_name text DEFAULT NULL::text, p_field text DEFAULT NULL::text, p_pattern text DEFAULT NULL::text, p_output_field text DEFAULT NULL::text, p_function_type text DEFAULT NULL::text, p_flags text DEFAULT NULL::text, p_replace_value text DEFAULT NULL::text, p_enabled boolean DEFAULT NULL::boolean, p_retain boolean DEFAULT NULL::boolean, p_sequence integer DEFAULT NULL::integer) - RETURNS dataflow.rules - LANGUAGE sql -AS $function$ - UPDATE dataflow.rules SET - name = COALESCE(p_name, name), - field = COALESCE(p_field, field), - pattern = COALESCE(p_pattern, pattern), - output_field = COALESCE(p_output_field, output_field), - function_type = COALESCE(p_function_type, function_type), - flags = COALESCE(p_flags, flags), - replace_value = COALESCE(p_replace_value, replace_value), - enabled = COALESCE(p_enabled, enabled), - retain = COALESCE(p_retain, retain), - sequence = COALESCE(p_sequence, sequence) - WHERE id = p_id - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.update_source(p_name text, p_constraint_fields text[] DEFAULT NULL::text[], p_config jsonb DEFAULT NULL::jsonb, p_global_picklist boolean DEFAULT NULL::boolean) - RETURNS dataflow.sources - LANGUAGE sql -AS $function$ - UPDATE dataflow.sources - SET constraint_fields = COALESCE(p_constraint_fields, constraint_fields), - config = COALESCE(p_config, config), - global_picklist = COALESCE(p_global_picklist, global_picklist), - updated_at = CURRENT_TIMESTAMP - WHERE name = p_name - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.update_source(p_name text, p_constraint_fields text[] DEFAULT NULL::text[], p_config jsonb DEFAULT NULL::jsonb) - RETURNS dataflow.sources - LANGUAGE sql -AS $function$ - UPDATE dataflow.sources - SET constraint_fields = COALESCE(p_constraint_fields, constraint_fields), - config = COALESCE(p_config, config), - updated_at = CURRENT_TIMESTAMP - WHERE name = p_name - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.update_stack(p_name text, p_label text DEFAULT NULL::text, p_fields jsonb DEFAULT NULL::jsonb, p_amount_field text DEFAULT NULL::text, p_date_field text DEFAULT NULL::text, p_balance_offset numeric DEFAULT NULL::numeric) - RETURNS dataflow.stacks - LANGUAGE sql -AS $function$ - UPDATE dataflow.stacks SET - label = COALESCE(p_label, label), - fields = COALESCE(p_fields, fields), - amount_field = COALESCE(p_amount_field, amount_field), - date_field = COALESCE(p_date_field, date_field), - balance_offset = COALESCE(p_balance_offset, balance_offset) - WHERE name = p_name - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.upsert_mapping(p_source_name text, p_rule_name text, p_input_value jsonb, p_output jsonb) - RETURNS dataflow.mappings - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.mappings (source_name, rule_name, input_value, output) - VALUES (p_source_name, p_rule_name, p_input_value, p_output) - ON CONFLICT (source_name, rule_name, input_value) - DO UPDATE SET output = EXCLUDED.output - RETURNING *; -$function$ -; -CREATE OR REPLACE FUNCTION dataflow.upsert_stack_source(p_stack_name text, p_source_name text, p_field_map jsonb DEFAULT '{}'::jsonb, p_amount_sign integer DEFAULT 1, p_balance_offset numeric DEFAULT 0, p_amount_field text DEFAULT NULL::text, p_date_field text DEFAULT NULL::text) - RETURNS dataflow.stack_sources - LANGUAGE sql -AS $function$ - INSERT INTO dataflow.stack_sources (stack_name, source_name, field_map, amount_sign, balance_offset, amount_field, date_field, seq) - VALUES ( - p_stack_name, p_source_name, p_field_map, p_amount_sign, p_balance_offset, p_amount_field, p_date_field, - (SELECT COALESCE(MAX(seq), 0) + 1 FROM dataflow.stack_sources WHERE stack_name = p_stack_name) - ) - ON CONFLICT (stack_name, source_name) DO UPDATE SET - field_map = EXCLUDED.field_map, - amount_sign = EXCLUDED.amount_sign, - balance_offset = EXCLUDED.balance_offset, - amount_field = EXCLUDED.amount_field, - date_field = EXCLUDED.date_field - RETURNING *; -$function$ -; - ------------------------------------------------------- --- Summary ------------------------------------------------------- --- All dataflow functions are defined above. --- Deploy with: psql -d dataflow -f database/functions.sql ------------------------------------------------------- diff --git a/database/import.sql b/database/import.sql new file mode 100644 index 0000000..43a5ad0 --- /dev/null +++ b/database/import.sql @@ -0,0 +1,159 @@ +-- +-- Import queries +-- CSV import and the import audit trail; SQL for the import/log endpoints in +-- api/routes/sources.js +-- + +SET search_path TO dataflow, public; + +-- ── Import ──────────────────────────────────────────────────────────────────── + +-- Import records, skipping any whose constraint key already exists in the table. +-- +-- Dedup is enforced here, in the CTE — there is no unique constraint on +-- constraint_key and ON CONFLICT must never be used. Within one batch every row +-- inserts even if two rows share a constraint key, because banks legitimately send +-- identical-looking transactions (same date, description, amount) on the same day. +-- The key exists only to stop a re-imported overlapping date range from +-- double-counting rows already in the table. +CREATE OR REPLACE FUNCTION import_records( + p_source_name TEXT, + p_data JSONB -- Array of records +) RETURNS JSON AS $$ +DECLARE + v_constraint_fields TEXT[]; + v_inserted INTEGER; + v_duplicates INTEGER; + v_log_id INTEGER; +BEGIN + SELECT constraint_fields INTO v_constraint_fields + FROM dataflow.sources + WHERE name = p_source_name; + + IF v_constraint_fields IS NULL THEN + RETURN json_build_object( + 'success', false, + 'error', 'Source not found: ' || p_source_name + ); + END IF; + + WITH + -- All incoming records with their constraint keys + pending AS ( + SELECT + rec.value AS data, + rec.ordinality AS seq, + (SELECT jsonb_object_agg(f, rec.value->>f) + FROM unnest(v_constraint_fields) AS f) AS constraint_key + FROM jsonb_array_elements(p_data) WITH ORDINALITY AS rec + ), + -- Keys already in the database (excluded) + existing AS ( + SELECT DISTINCT r.constraint_key + FROM dataflow.records r + INNER JOIN pending p ON p.constraint_key = r.constraint_key + WHERE r.source_name = p_source_name + ), + -- Rows whose constraint key is not yet in the database + new_records AS ( + SELECT p.data, p.constraint_key, p.seq + FROM pending p + WHERE NOT EXISTS (SELECT 1 FROM existing e WHERE e.constraint_key = p.constraint_key) + ), + -- Write the log entry + log_entry AS ( + INSERT INTO dataflow.import_log (source_name, records_imported, records_duplicate, info) + VALUES ( + p_source_name, + (SELECT count(*) FROM new_records), + (SELECT count(*) FROM pending) - (SELECT count(*) FROM new_records), + jsonb_build_object( + 'total', jsonb_array_length(p_data), + 'inserted_keys', (SELECT jsonb_agg(constraint_key ORDER BY constraint_key) FROM new_records), + 'excluded_keys', (SELECT jsonb_agg(constraint_key) FROM existing) + ) + ) + RETURNING id, records_imported, records_duplicate + ), + -- Insert new records + inserted AS ( + INSERT INTO dataflow.records (source_name, data, constraint_key, import_id) + SELECT p_source_name, nr.data, nr.constraint_key, (SELECT id FROM log_entry) + FROM new_records nr + ORDER BY nr.seq + RETURNING id + ) + SELECT le.id, le.records_imported, le.records_duplicate + INTO v_log_id, v_inserted, v_duplicates + FROM log_entry le; + + RETURN json_build_object( + 'success', true, + 'imported', v_inserted, + 'duplicates', v_duplicates, + 'log_id', v_log_id + ); +END; +$$ LANGUAGE plpgsql; + +COMMENT ON FUNCTION import_records IS 'Import records with automatic deduplication'; + +-- ── Audit trail ─────────────────────────────────────────────────────────────── + +CREATE OR REPLACE FUNCTION get_import_log(p_source_name TEXT) +RETURNS TABLE ( + id INTEGER, + source_name TEXT, + records_imported INTEGER, + records_duplicate INTEGER, + imported_at TIMESTAMPTZ, + info JSONB +) AS $$ + SELECT id, source_name, records_imported, records_duplicate, imported_at, info + FROM dataflow.import_log + WHERE source_name = p_source_name + ORDER BY imported_at DESC; +$$ LANGUAGE sql; + +COMMENT ON FUNCTION get_import_log IS 'Return import history for a source, newest first, including inserted/excluded key lists'; + +CREATE OR REPLACE FUNCTION get_all_import_logs() +RETURNS TABLE ( + id INTEGER, + source_name TEXT, + records_imported INTEGER, + records_duplicate INTEGER, + imported_at TIMESTAMPTZ, + info JSONB +) AS $$ + SELECT id, source_name, records_imported, records_duplicate, imported_at, info + FROM dataflow.import_log + ORDER BY imported_at DESC; +$$ LANGUAGE sql; + +COMMENT ON FUNCTION get_all_import_logs IS 'Return import history across all sources, newest first'; + +-- Records are removed by the import_id FK's ON DELETE CASCADE +CREATE OR REPLACE FUNCTION delete_import(p_log_id INTEGER) +RETURNS JSON AS $$ +DECLARE + v_deleted INTEGER; +BEGIN + IF NOT EXISTS (SELECT 1 FROM dataflow.import_log WHERE id = p_log_id) THEN + RETURN json_build_object('success', false, 'error', 'Import log entry not found'); + END IF; + + SELECT count(*) INTO v_deleted FROM dataflow.records WHERE import_id = p_log_id; + + -- Cascade handles deleting records via FK ON DELETE CASCADE + DELETE FROM dataflow.import_log WHERE id = p_log_id; + + RETURN json_build_object( + 'success', true, + 'records_deleted', v_deleted, + 'log_id', p_log_id + ); +END; +$$ LANGUAGE plpgsql; + +COMMENT ON FUNCTION delete_import IS 'Delete all records belonging to an import batch and remove the log entry'; diff --git a/database/queries/mappings.sql b/database/mappings.sql similarity index 100% rename from database/queries/mappings.sql rename to database/mappings.sql diff --git a/database/migrate_input_value_jsonb.sql b/database/migrate_input_value_jsonb.sql deleted file mode 100644 index 515770a..0000000 --- a/database/migrate_input_value_jsonb.sql +++ /dev/null @@ -1,22 +0,0 @@ --- --- Migration: Change mappings.input_value from TEXT to JSONB --- Allows multi-capture-group regex results to be used as mapping keys --- - -SET search_path TO dataflow, public; - --- Drop dependent constraint and index first -ALTER TABLE dataflow.mappings DROP CONSTRAINT mappings_source_name_rule_name_input_value_key; -DROP INDEX IF EXISTS dataflow.idx_mappings_input; - --- Convert column: existing TEXT values become JSONB strings e.g. "MEIJER" -ALTER TABLE dataflow.mappings - ALTER COLUMN input_value TYPE JSONB - USING to_jsonb(input_value); - --- Recreate constraint and index -ALTER TABLE dataflow.mappings - ADD CONSTRAINT mappings_source_name_rule_name_input_value_key - UNIQUE (source_name, rule_name, input_value); - -CREATE INDEX idx_mappings_input ON dataflow.mappings(source_name, rule_name, input_value); diff --git a/database/migrate_overrides_column.sql b/database/migrate_overrides_column.sql deleted file mode 100644 index bd0c8f7..0000000 --- a/database/migrate_overrides_column.sql +++ /dev/null @@ -1,40 +0,0 @@ --- --- Migration: add overrides column to records --- --- Separates the three data layers: --- data — original import values, never mutated --- transformed — rule/mapping output fields only (delta) --- overrides — manual user overrides (highest precedence) --- --- Consumers merge as: data || COALESCE(transformed,'{}') || COALESCE(overrides,'{}') --- --- Safe to run multiple times (IF NOT EXISTS guards). --- - -SET search_path TO dataflow, public; - --- 1. Add overrides column -ALTER TABLE dataflow.records - ADD COLUMN IF NOT EXISTS overrides JSONB; - --- 2. Add partial GIN index (only indexes rows that have overrides) -CREATE INDEX IF NOT EXISTS idx_records_overrides - ON dataflow.records USING gin(overrides) - WHERE overrides IS NOT NULL; - --- 3. Redeploy functions (CREATE OR REPLACE — non-destructive) -\i functions.sql - --- 4. Reprocess all sources to strip stale data keys from transformed --- (apply_transformations now writes only rule additions, not data || additions) -DO $$ -DECLARE - src TEXT; - result JSON; -BEGIN - FOR src IN SELECT name FROM dataflow.sources ORDER BY name LOOP - SELECT dataflow.reprocess_records(src) INTO result; - RAISE NOTICE 'Reprocessed %: %', src, result; - END LOOP; -END; -$$; diff --git a/database/migrate_pivot_layouts_drop_fk.sql b/database/migrate_pivot_layouts_drop_fk.sql deleted file mode 100644 index 6735691..0000000 --- a/database/migrate_pivot_layouts_drop_fk.sql +++ /dev/null @@ -1,4 +0,0 @@ --- Drop the foreign key from pivot_layouts.source_name so stack view names can also --- be used as layout keys (stacks are not rows in the sources table). -ALTER TABLE dataflow.pivot_layouts - DROP CONSTRAINT pivot_layouts_source_name_fkey; diff --git a/database/migrate_tps.sql b/database/migrate_tps.sql deleted file mode 100644 index a4adf71..0000000 --- a/database/migrate_tps.sql +++ /dev/null @@ -1,121 +0,0 @@ --- --- TPS → Dataflow Migration --- --- Migrates sources, rules, mappings, and records from the TPS system. --- Run against the dataflow database: --- PGPASSWORD=dataflow psql -U dataflow -d dataflow -h localhost -f database/migrate_tps.sql --- --- Existing rows are skipped (ON CONFLICT DO NOTHING) so the script is safe to re-run. --- NOTE: dcard already configured in dataflow will NOT be overwritten. --- - -SET search_path TO dataflow, public; - -CREATE EXTENSION IF NOT EXISTS dblink; - --- Connection string to the TPS database -\set tps_conn 'host=192.168.1.110 dbname=ubm user=api password=gyaswddh1983' - -\echo '' -\echo '=== 1. Sources ===' - -INSERT INTO dataflow.sources (name, constraint_fields, config) -SELECT - srce AS name, - -- Strip {} wrappers from constraint paths → constraint field names - ARRAY( - SELECT regexp_replace(c, '^\{|\}$', '', 'g') - FROM jsonb_array_elements_text(defn->'constraint') AS c - ) AS constraint_fields, - -- Build config.fields from the first schema (index 0 = "mapped" for dcard, "default" for others) - jsonb_build_object('fields', - (SELECT jsonb_agg( - jsonb_build_object( - 'name', regexp_replace(col->>'path', '^\{|\}$', '', 'g'), - 'type', COALESCE(NULLIF(col->>'type', ''), 'text') - ) ORDER BY ord - ) - FROM jsonb_array_elements(defn->'schemas'->0->'columns') - WITH ORDINALITY AS t(col, ord) - ) - ) AS config -FROM dblink(:'tps_conn', - 'SELECT srce, defn FROM tps.srce' -) AS t(srce TEXT, defn JSONB) -ON CONFLICT (name) DO NOTHING; - -SELECT name, constraint_fields, jsonb_array_length(config->'fields') AS field_count -FROM dataflow.sources ORDER BY name; - -\echo '' -\echo '=== 2. Rules ===' - -INSERT INTO dataflow.rules - (source_name, name, field, pattern, output_field, function_type, flags, replace_value, sequence, enabled, retain) -SELECT - srce AS source_name, - target AS name, - -- Strip {} from the input field key - regexp_replace(regex->'regex'->'defn'->0->>'key', '^\{|\}$', '', 'g') AS field, - regex->'regex'->'defn'->0->>'regex' AS pattern, - regex->'regex'->'defn'->0->>'field' AS output_field, - COALESCE(NULLIF(regex->'regex'->>'function', ''), 'extract') AS function_type, - COALESCE(regex->'regex'->'defn'->0->>'flag', '') AS flags, - '' AS replace_value, - seq AS sequence, - true AS enabled, - (regex->'regex'->'defn'->0->>'retain') = 'y' AS retain -FROM dblink(:'tps_conn', - 'SELECT srce, target, seq, regex FROM tps.map_rm' -) AS t(srce TEXT, target TEXT, seq INT, regex JSONB) -ON CONFLICT (source_name, name) DO NOTHING; - -SELECT source_name, name, field, pattern, output_field, sequence -FROM dataflow.rules ORDER BY source_name, sequence; - -\echo '' -\echo '=== 3. Mappings ===' - -INSERT INTO dataflow.mappings (source_name, rule_name, input_value, output) -SELECT - srce AS source_name, - target AS rule_name, - -- retval is {"f20": ""} — pull out the value as JSONB - (SELECT value FROM jsonb_each(retval) LIMIT 1) AS input_value, - map AS output -FROM dblink(:'tps_conn', - 'SELECT srce, target, retval, map FROM tps.map_rv' -) AS t(srce TEXT, target TEXT, retval JSONB, map JSONB) -ON CONFLICT (source_name, rule_name, input_value) DO NOTHING; - -SELECT source_name, rule_name, COUNT(*) AS mapping_count -FROM dataflow.mappings GROUP BY source_name, rule_name ORDER BY source_name, rule_name; - -\echo '' -\echo '=== 4. Records ===' -\echo ' (13 000+ rows — may take a moment)' - -INSERT INTO dataflow.records (source_name, data, constraint_key, transformed, imported_at, transformed_at) -SELECT - t.srce AS source_name, - t.rec AS data, - (SELECT jsonb_object_agg(f, t.rec->>f) FROM unnest(s.constraint_fields) AS f) AS constraint_key, - t.allj AS transformed, - CURRENT_TIMESTAMP AS imported_at, - CASE WHEN t.allj IS NOT NULL THEN CURRENT_TIMESTAMP END AS transformed_at -FROM dblink(:'tps_conn', - 'SELECT srce, rec, allj FROM tps.trans' -) AS t(srce TEXT, rec JSONB, allj JSONB) -JOIN dataflow.sources s ON s.name = t.srce -ON CONFLICT (source_name, constraint_key) DO NOTHING; - -SELECT source_name, COUNT(*) AS records, COUNT(transformed) AS transformed -FROM dataflow.records GROUP BY source_name ORDER BY source_name; - -\echo '' -\echo '=== Migration complete ===' -SELECT - (SELECT COUNT(*) FROM dataflow.sources) AS sources, - (SELECT COUNT(*) FROM dataflow.rules) AS rules, - (SELECT COUNT(*) FROM dataflow.mappings) AS mappings, - (SELECT COUNT(*) FROM dataflow.records) AS records; diff --git a/database/queries/records.sql b/database/records.sql similarity index 67% rename from database/queries/records.sql rename to database/records.sql index 410b338..9b2dea7 100644 --- a/database/queries/records.sql +++ b/database/records.sql @@ -41,37 +41,45 @@ $$ LANGUAGE sql STABLE; -- ── Overrides ───────────────────────────────────────────────────────────────── --- Store manual overrides and immediately merge into transformed -CREATE OR REPLACE FUNCTION set_record_overrides(p_id INT, p_overrides JSONB) -RETURNS dataflow.records AS $$ - UPDATE dataflow.records - SET overrides = CASE WHEN p_overrides = '{}'::jsonb THEN NULL ELSE p_overrides END, - transformed = COALESCE(transformed, data) || COALESCE(p_overrides, '{}'::jsonb) - WHERE id = p_id - RETURNING *; +-- Store manual overrides. Overrides stay in their own column — never merged into +-- transformed — so reprocessing rules cannot clobber a manual edit. +DROP FUNCTION IF EXISTS set_record_overrides(INTEGER, JSONB); +CREATE OR REPLACE FUNCTION set_record_overrides(p_id INTEGER, p_overrides JSONB) +RETURNS JSON AS $$ + WITH updated AS ( + UPDATE dataflow.records + SET overrides = CASE WHEN p_overrides = '{}'::jsonb THEN NULL ELSE p_overrides END + WHERE id = p_id + RETURNING * + ) + SELECT row_to_json(updated) FROM updated; $$ LANGUAGE sql; -- Merge overrides into multiple records at once; returns actual updated count -CREATE OR REPLACE FUNCTION bulk_set_record_overrides(p_source_name TEXT, p_ids INT[], p_overrides JSONB) -RETURNS BIGINT AS $$ +DROP FUNCTION IF EXISTS bulk_set_record_overrides(TEXT, INTEGER[], JSONB); +CREATE OR REPLACE FUNCTION bulk_set_record_overrides(p_source_name TEXT, p_ids INTEGER[], p_overrides JSONB) +RETURNS JSON AS $$ WITH updated AS ( UPDATE dataflow.records - SET overrides = COALESCE(overrides, '{}'::jsonb) || p_overrides, - transformed = COALESCE(transformed, data) || p_overrides + SET overrides = COALESCE(overrides, '{}'::jsonb) || p_overrides WHERE id = ANY(p_ids) AND source_name = p_source_name RETURNING id ) - SELECT count(*) FROM updated; + SELECT json_build_object('updated', count(*)) FROM updated; $$ LANGUAGE sql; --- Clear overrides; caller should reprocess to restore computed transformed value -CREATE OR REPLACE FUNCTION clear_record_overrides(p_id INT) -RETURNS dataflow.records AS $$ - UPDATE dataflow.records - SET overrides = NULL - WHERE id = p_id - RETURNING *; +-- Clear overrides; the computed values in transformed are untouched +DROP FUNCTION IF EXISTS clear_record_overrides(INTEGER); +CREATE OR REPLACE FUNCTION clear_record_overrides(p_id INTEGER) +RETURNS JSON AS $$ + WITH updated AS ( + UPDATE dataflow.records + SET overrides = NULL + WHERE id = p_id + RETURNING * + ) + SELECT row_to_json(updated) FROM updated; $$ LANGUAGE sql; -- ── Delete ──────────────────────────────────────────────────────────────────── diff --git a/database/queries/rules.sql b/database/rules.sql similarity index 87% rename from database/queries/rules.sql rename to database/rules.sql index e79e925..7cd48cb 100644 --- a/database/queries/rules.sql +++ b/database/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/database/queries/sources.sql b/database/sources.sql similarity index 92% rename from database/queries/sources.sql rename to database/sources.sql index 0c8cc32..cf31cba 100644 --- a/database/queries/sources.sql +++ b/database/sources.sql @@ -40,8 +40,6 @@ RETURNS TEXT AS $$ DELETE FROM dataflow.sources WHERE name = p_name RETURNING name; $$ LANGUAGE sql; --- ── Import log ──────────────────────────────────────────────────────────────── - -- ── Stats ───────────────────────────────────────────────────────────────────── CREATE OR REPLACE FUNCTION get_source_stats(p_source_name TEXT) @@ -161,6 +159,7 @@ BEGIN RETURN json_build_object('success', false, 'error', 'No schema fields defined for this source'); END IF; + -- Columns read from r, the merged data || transformed || overrides object FOR v_field IN SELECT * FROM jsonb_array_elements(v_config->'fields') LOOP IF v_cols != '' THEN v_cols := v_cols || ', '; END IF; @@ -171,24 +170,27 @@ BEGIN BEGIN WHILE v_expr ~ '\{[^}]+\}' LOOP v_ref := substring(v_expr FROM '\{([^}]+)\}'); - v_expr := replace(v_expr, '{' || v_ref || '}', format('(transformed->>%L)::numeric', v_ref)); + v_expr := replace(v_expr, '{' || v_ref || '}', format('(r->>%L)::numeric', v_ref)); END LOOP; v_cols := v_cols || format('%s AS %I', v_expr, v_field->>'name'); END; ELSE CASE v_field->>'type' - WHEN 'date' THEN v_cols := v_cols || format('(transformed->>%L)::date AS %I', v_field->>'name', v_field->>'name'); - WHEN 'numeric' THEN v_cols := v_cols || format('(transformed->>%L)::numeric AS %I', v_field->>'name', v_field->>'name'); - ELSE v_cols := v_cols || format('transformed->>%L AS %I', v_field->>'name', v_field->>'name'); + WHEN 'date' THEN v_cols := v_cols || format('(r->>%L)::date AS %I', v_field->>'name', v_field->>'name'); + WHEN 'numeric' THEN v_cols := v_cols || format('(r->>%L)::numeric AS %I', v_field->>'name', v_field->>'name'); + ELSE v_cols := v_cols || format('r->>%L AS %I', v_field->>'name', v_field->>'name'); END CASE; END IF; END LOOP; CREATE SCHEMA IF NOT EXISTS dfv; v_view := 'dfv.' || quote_ident(p_source_name); - EXECUTE format('DROP VIEW IF EXISTS %s', v_view); + EXECUTE format('DROP VIEW IF EXISTS %s CASCADE', v_view); v_sql := format( - 'CREATE VIEW %s AS SELECT id, overrides IS NOT NULL AS _overridden, %s FROM dataflow.records WHERE source_name = %L AND transformed IS NOT NULL', + 'CREATE VIEW %s AS SELECT id, _overridden, %s FROM (' + || 'SELECT id, overrides IS NOT NULL AS _overridden, ' + || 'data || COALESCE(transformed, ''{}''::jsonb) || COALESCE(overrides, ''{}''::jsonb) AS r ' + || 'FROM dataflow.records WHERE source_name = %L AND transformed IS NOT NULL) rec', v_view, v_cols, p_source_name ); EXECUTE v_sql; diff --git a/database/queries/stacks.sql b/database/stacks.sql similarity index 100% rename from database/queries/stacks.sql rename to database/stacks.sql diff --git a/database/queries/status.sql b/database/status.sql similarity index 100% rename from database/queries/status.sql rename to database/status.sql diff --git a/database/transform.sql b/database/transform.sql new file mode 100644 index 0000000..aa3e4c4 --- /dev/null +++ b/database/transform.sql @@ -0,0 +1,156 @@ +-- +-- Transform queries +-- The rule/mapping engine; SQL for the transform endpoints in api/routes/sources.js, +-- api/routes/rules.js and api/routes/records.js +-- +-- Order matters within this file: the aggregate is used by apply_transformations, +-- which in turn is called by reprocess_records. +-- + +SET search_path TO dataflow, public; + +-- ── Merge aggregate ─────────────────────────────────────────────────────────── + +-- Merge JSONB objects across rows (later rows win on key conflicts) +-- Usage: jsonb_concat_obj(col ORDER BY sequence) +CREATE OR REPLACE FUNCTION dataflow.jsonb_merge(a JSONB, b JSONB) +RETURNS JSONB AS $$ + SELECT COALESCE(a, '{}') || COALESCE(b, '{}') +$$ LANGUAGE sql IMMUTABLE; + +DROP AGGREGATE IF EXISTS dataflow.jsonb_concat_obj(JSONB); +CREATE AGGREGATE dataflow.jsonb_concat_obj(JSONB) ( + sfunc = dataflow.jsonb_merge, + stype = JSONB, + initcond = '{}' +); + +-- ── Apply rules and mappings ────────────────────────────────────────────────── + +-- Writes only the rule/mapping output into records.transformed. Raw values stay in +-- data and manual edits stay in overrides; readers merge the three layers. +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 +) RETURNS JSON AS $$ +WITH +-- All records to process +qualifying AS ( + SELECT id, data + 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 an empty object +updated AS ( + UPDATE dataflow.records rec + SET transformed = COALESCE(ra.additions, '{}'::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; + +COMMENT ON FUNCTION apply_transformations IS 'Apply transformation rules and mappings to records (set-based CTE)'; + +-- ── Reprocess ───────────────────────────────────────────────────────────────── + +CREATE OR REPLACE FUNCTION reprocess_records(p_source_name TEXT) +RETURNS JSON AS $$ + -- Overwrite all records directly — no clear step, mirrors TPS srce_map_overwrite + SELECT dataflow.apply_transformations(p_source_name, NULL, TRUE) +$$ LANGUAGE sql; + +COMMENT ON FUNCTION reprocess_records IS 'Reapply all transformations for a source, overwriting existing values'; diff --git a/manage.py b/manage.py index efa33dd..ccb7268 100755 --- a/manage.py +++ b/manage.py @@ -18,13 +18,16 @@ SERVICE_FILE = Path('/etc/systemd/system/dataflow.service') SERVICE_SRC = ROOT / 'dataflow.service' NGINX_DIR = Path('/etc/nginx/sites-enabled') -# Deployed in order — stacks.sql creates tables that reference sources -QUERIES_DIR = ROOT / 'database' / 'queries' +# Deployed in order — stacks.sql creates tables that reference sources, and +# transform.sql defines an aggregate its own functions depend on +QUERIES_DIR = ROOT / 'database' QUERY_FILES = [ QUERIES_DIR / 'sources.sql', QUERIES_DIR / 'rules.sql', QUERIES_DIR / 'mappings.sql', QUERIES_DIR / 'records.sql', + QUERIES_DIR / 'import.sql', + QUERIES_DIR / 'transform.sql', QUERIES_DIR / 'stacks.sql', QUERIES_DIR / 'status.sql', ] @@ -164,23 +167,31 @@ def ui_build_time(): return datetime.fromtimestamp(ts).strftime('%Y-%m-%d %H:%M') return None -def nginx_domain(port): - """Find nginx site proxying to our port.""" +def nginx_conf_path(port): + """Path of the nginx site proxying to our port, if any.""" if not NGINX_DIR.exists(): return None for f in NGINX_DIR.iterdir(): try: - text = f.read_text() - if f':{port}' in text: - for line in text.splitlines(): - if 'server_name' in line: - parts = line.split() - if len(parts) >= 2: - return parts[1].rstrip(';') + if f':{port}' in f.read_text(): + return f except Exception: pass return None + +def nginx_domain(port): + """server_name of the nginx site proxying to our port.""" + conf = nginx_conf_path(port) + if not conf: + return None + for line in conf.read_text().splitlines(): + if 'server_name' in line: + parts = line.split() + if len(parts) >= 2: + return parts[1].rstrip(';') + return None + def sudo_run(args, **kwargs): return subprocess.run(['sudo'] + args, **kwargs) @@ -433,7 +444,7 @@ def action_deploy_schema(cfg): def action_deploy_functions(cfg): - header('Deploy SQL functions (database/queries/)') + header('Deploy SQL functions (database/*.sql)') if not cfg: err(f'{ENV_FILE} not found — run option 1 to configure the database connection first') return @@ -729,6 +740,116 @@ def action_stop_service(): ok('dataflow.service stopped') +def action_uninstall(cfg): + """Reverse everything this script installs, outside the repo itself.""" + header('Uninstall dataflow') + + port = cfg.get('API_PORT', '3020') if cfg else '3020' + db_name = cfg.get('DB_NAME', 'dataflow') if cfg else 'dataflow' + db_user = cfg.get('DB_USER', 'dataflow') if cfg else 'dataflow' + conf_path = nginx_conf_path(port) + + # Everything that exists right now, in reverse install order + targets = [] + if service_installed(): + targets.append(f'systemd service {SERVICE_FILE}' + + (' (running)' if service_running() else '')) + if conf_path: + targets.append(f'nginx site {conf_path}') + if cfg and can_connect(cfg): + targets.append(f'database "{db_name}" on {cfg["DB_HOST"]}:{cfg["DB_PORT"]} (ALL DATA)') + targets.append(f'database user {db_user}') + if ENV_FILE.exists(): + targets.append(f'config {ENV_FILE}') + if (ROOT / 'public').exists(): + targets.append(f'built UI {ROOT / "public"}') + if (ROOT / 'node_modules').exists(): + targets.append(f'dependencies {ROOT / "node_modules"}') + + if not targets: + info('Nothing installed to remove.') + return cfg + + print(' This will permanently remove:') + for t in targets: + print(f' {t}') + print() + info(f'The repository itself ({ROOT}) is left alone — delete it manually if you want it gone.') + print() + + if input(" Type 'delete' to confirm: ").strip() != 'delete': + info('Cancelled — no changes made') + return cfg + + # ── Service ─────────────────────────────────────────────────────────────── + if service_installed(): + print() + print(' Removing systemd service...') + sudo_run(['systemctl', 'stop', 'dataflow']) + sudo_run(['systemctl', 'disable', 'dataflow']) + r = sudo_run(['rm', '-f', str(SERVICE_FILE)]) + if r.returncode != 0: + err(f'Could not remove {SERVICE_FILE} — check sudo permissions') + else: + sudo_run(['systemctl', 'daemon-reload']) + ok(f'Service stopped, disabled, and {SERVICE_FILE} removed') + + # ── nginx ───────────────────────────────────────────────────────────────── + if conf_path: + print() + print(' Removing nginx site...') + r = sudo_run(['rm', '-f', str(conf_path)]) + if r.returncode != 0: + err(f'Could not remove {conf_path} — check sudo permissions') + elif sudo_run(['nginx', '-t'], capture_output=True).returncode != 0: + err('nginx config test failed after removal — not reloading; check nginx manually') + else: + sudo_run(['systemctl', 'reload', 'nginx']) + ok(f'{conf_path} removed and nginx reloaded') + + # ── Database ────────────────────────────────────────────────────────────── + if cfg and can_connect(cfg): + print() + print(f' Dropping the database requires PostgreSQL admin credentials.') + admin = { + 'user': prompt('PostgreSQL admin username', 'postgres'), + 'password': prompt('PostgreSQL admin password', secret=True), + 'host': cfg['DB_HOST'], + 'port': cfg['DB_PORT'], + } + r = psql_admin(admin, 'SELECT 1') + if r.returncode != 0: + err(f'Cannot connect as admin — database and user left in place\n{r.stderr.strip()}') + else: + r = psql_admin(admin, f'DROP DATABASE IF EXISTS {db_name}') + if r.returncode != 0: + err(f'Could not drop database "{db_name}"\n{r.stderr.strip()}') + else: + ok(f'Database "{db_name}" dropped') + r = psql_admin(admin, f'DROP USER IF EXISTS {db_user}') + if r.returncode != 0: + err(f'Could not drop user {db_user}\n{r.stderr.strip()}') + else: + ok(f'User {db_user} dropped') + + # ── Generated files ─────────────────────────────────────────────────────── + print() + for path, label in [(ENV_FILE, 'config'), + (ROOT / 'public', 'built UI'), + (ROOT / 'node_modules', 'dependencies')]: + if not path.exists(): + continue + if path.is_dir(): + shutil.rmtree(path, ignore_errors=True) + else: + path.unlink() + ok(f'Removed {label} ({path})') + + print() + ok('Uninstall complete') + return None + + def action_set_login_credentials(cfg): header('Set login credentials (LOGIN_USER / LOGIN_PASSWORD_HASH in .env)') @@ -786,13 +907,14 @@ def action_set_login_credentials(cfg): MENU = [ ('Database configuration and deployment dialog (.env)', action_configure), ('Redeploy "dataflow" schema only (database/schema.sql)', action_deploy_schema), - ('Redeploy SQL functions only (database/queries/)', action_deploy_functions), + ('Redeploy SQL functions only (database/*.sql)', action_deploy_functions), ('Build UI (ui/ → public/)', action_build_ui), ('Set up nginx reverse proxy', action_setup_nginx), ('Install dataflow systemd service unit', action_install_service), ('Start / restart dataflow.service', action_restart_service), ('Stop dataflow.service', action_stop_service), ('Set login credentials', action_set_login_credentials), + ('Uninstall (service, nginx, database, .env, build)', action_uninstall), ] def main(): @@ -805,14 +927,11 @@ def main(): show_status(cfg) db_target = f'into "{cfg["DB_NAME"]}" on {cfg["DB_HOST"]}' if cfg else '(not configured)' - DB_ACTIONS = { - 'Deploy "dataflow" schema (database/schema.sql)', - 'Deploy SQL functions (database/functions.sql)', - } + DB_ACTIONS = {action_deploy_schema, action_deploy_functions} print(bold('Actions')) - for i, (label, _) in enumerate(MENU, 1): - suffix = f' {dim(db_target)}' if label in DB_ACTIONS else '' + for i, (label, fn) in enumerate(MENU, 1): + suffix = f' {dim(db_target)}' if fn in DB_ACTIONS else '' print(f' {cyan(str(i))}. {label}{suffix}') print(f' {cyan("q")}. Quit') print() @@ -828,13 +947,12 @@ def main(): if 0 <= idx < len(MENU): label, fn = MENU[idx] import inspect - sig = inspect.signature(fn) - if len(sig.parameters) == 0: - result = fn() - elif len(sig.parameters) == 1: - result = fn(cfg) - if label.startswith('Configure') and result is not None: - cfg = result + # cfg is reloaded from .env at the top of every loop, so a return + # value is only ever informational + if len(inspect.signature(fn).parameters) == 0: + fn() + else: + fn(cfg) pause() else: warn('Invalid choice — enter a number from the list above') diff --git a/uninstall.sh b/uninstall.sh deleted file mode 100755 index 3c38e39..0000000 --- a/uninstall.sh +++ /dev/null @@ -1,85 +0,0 @@ -#!/bin/bash -# -# Dataflow Uninstall Script -# Removes database user, database, and optionally .env -# - -echo "⚠️ Dataflow Uninstall" -echo "=====================" -echo "" - -# Load .env if it exists -if [ -f .env ]; then - export $(cat .env | grep -v '^#' | xargs) -fi - -DB_NAME=${DB_NAME:-dataflow} -DB_USER=${DB_USER:-dataflow} - -echo "⚠️ This will permanently delete:" -echo " - Database: $DB_NAME" -echo " - User: $DB_USER" -echo "" -read -p "Type 'delete' to confirm: " CONFIRM -if [ "$CONFIRM" != "delete" ]; then - echo "Cancelled." - exit 0 -fi - -# Prompt for admin credentials -echo "" -echo "📋 PostgreSQL Admin Credentials" -echo "" -read -p "Admin username [postgres]: " ADMIN_USER -ADMIN_USER=${ADMIN_USER:-postgres} -read -s -p "Admin password: " ADMIN_PASS -echo "" - -DB_HOST=${DB_HOST:-localhost} -DB_PORT=${DB_PORT:-5432} -DB_NAME=${DB_NAME:-dataflow} -DB_USER=${DB_USER:-dataflow} - -# Test admin connection -echo "" -echo "🔍 Testing PostgreSQL admin connection..." -export PGPASSWORD="$ADMIN_PASS" -if ! psql -U "$ADMIN_USER" -h "$DB_HOST" -p "$DB_PORT" -d postgres -c '\q' 2>/dev/null; then - echo "✗ Cannot connect to PostgreSQL" - exit 1 -fi -echo "✓ Connected" - -# Drop database -echo "" -echo "🗄️ Dropping database..." -psql -U "$ADMIN_USER" -h "$DB_HOST" -p "$DB_PORT" -d postgres -c "DROP DATABASE IF EXISTS $DB_NAME;" 2>/dev/null || true -echo "✓ Database dropped" - -# Drop user -echo "" -echo "👤 Dropping user..." -psql -U "$ADMIN_USER" -h "$DB_HOST" -p "$DB_PORT" -d postgres -c "DROP USER IF EXISTS $DB_USER;" 2>/dev/null || true -echo "✓ User dropped" - -unset PGPASSWORD - -# Optionally remove .env -echo "" -read -p "Remove .env file? [y/N]: " REMOVE_ENV -if [ "$REMOVE_ENV" == "y" ] || [ "$REMOVE_ENV" == "Y" ]; then - rm -f .env - echo "✓ .env removed" -fi - -# Optionally remove node_modules -echo "" -read -p "Remove node_modules? [y/N]: " REMOVE_MODULES -if [ "$REMOVE_MODULES" == "y" ] || [ "$REMOVE_MODULES" == "Y" ]; then - rm -rf node_modules - echo "✓ node_modules removed" -fi - -echo "" -echo "✅ Uninstall complete!" -echo ""