from collections.abc import Callable from ..core.database import API_DELETE, API_UPDATE, ChangeSet, g_delete_prefix, try_read from ..model.elements import get_node_links, is_pipe, is_pump, is_valve from ..model.reservoirs import unset_reservoir_by_pattern from ..model.tanks import unset_tank_by_curve from ..model.pumps import unset_pump_by_curve, unset_pump_by_pattern from ..model.tags import delete_tag_by_node, delete_tag_by_link from ..model.demands import delete_demand_by_junction from ..model.status import delete_status_by_link from ..model.energy import delete_pump_energy_by_pump, unset_pump_energy_by_pattern, unset_pump_energy_by_curve from ..model.emitters import delete_emitter_by_junction from ..model.quality import delete_quality_by_node from ..model.sources import delete_source_by_node, unset_source_by_pattern from ..model.reactions import delete_pipe_reaction_by_pipe, delete_tank_reaction_by_tank from ..model.mixing import delete_mixing_by_tank from ..gis.vertices import delete_vertex_by_link from ..gis.labels import unset_label_by_node from ..model.options import generate_v2, generate_v3 def expand_junction_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.junctions where node_id = %s", (id,)) if row == None: return result links = get_node_links(name, id) for link in links: if is_pipe(name, link): result.merge(expand_pipe_delete(name, ChangeSet(g_delete_prefix | {'type': 'pipe', 'id': link}))) if is_pump(name, link): result.merge(expand_pump_delete(name, ChangeSet(g_delete_prefix | {'type': 'pump', 'id': link}))) if is_valve(name, link): result.merge(expand_valve_delete(name, ChangeSet(g_delete_prefix | {'type': 'valve', 'id': link}))) result.merge(delete_tag_by_node(name, id)) result.merge(delete_demand_by_junction(name, id)) result.merge(delete_emitter_by_junction(name, id)) result.merge(delete_quality_by_node(name, id)) result.merge(delete_source_by_node(name, id)) result.merge(unset_label_by_node(name, id)) result.merge(cs) return result def expand_reservoir_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.reservoirs where node_id = %s", (id,)) if row == None: return result links = get_node_links(name, id) for link in links: if is_pipe(name, link): result.merge(expand_pipe_delete(name, ChangeSet(g_delete_prefix | {'type': 'pipe', 'id': link}))) if is_pump(name, link): result.merge(expand_pump_delete(name, ChangeSet(g_delete_prefix | {'type': 'pump', 'id': link}))) if is_valve(name, link): result.merge(expand_valve_delete(name, ChangeSet(g_delete_prefix | {'type': 'valve', 'id': link}))) result.merge(delete_tag_by_node(name, id)) result.merge(delete_quality_by_node(name, id)) result.merge(delete_source_by_node(name, id)) result.merge(unset_label_by_node(name, id)) result.merge(cs) return result def expand_tank_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.tanks where node_id = %s", (id,)) if row == None: return result links = get_node_links(name, id) for link in links: if is_pipe(name, link): result.merge(expand_pipe_delete(name, ChangeSet(g_delete_prefix | {'type': 'pipe', 'id': link}))) if is_pump(name, link): result.merge(expand_pump_delete(name, ChangeSet(g_delete_prefix | {'type': 'pump', 'id': link}))) if is_valve(name, link): result.merge(expand_valve_delete(name, ChangeSet(g_delete_prefix | {'type': 'valve', 'id': link}))) result.merge(delete_tag_by_node(name, id)) result.merge(delete_quality_by_node(name, id)) result.merge(delete_source_by_node(name, id)) result.merge(delete_tank_reaction_by_tank(name, id)) result.merge(delete_mixing_by_tank(name, id)) result.merge(unset_label_by_node(name, id)) result.merge(cs) return result def expand_pipe_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.pipes where link_id = %s", (id,)) if row == None: return result result.merge(delete_tag_by_link(name, id)) result.merge(delete_status_by_link(name, id)) result.merge(delete_pipe_reaction_by_pipe(name, id)) result.merge(delete_vertex_by_link(name, id)) result.merge(cs) return result def expand_pump_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.pumps where link_id = %s", (id,)) if row == None: return result result.merge(delete_tag_by_link(name, id)) result.merge(delete_status_by_link(name, id)) result.merge(delete_pump_energy_by_pump(name, id)) result.merge(delete_vertex_by_link(name, id)) result.merge(cs) return result def expand_valve_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.valves where link_id = %s", (id,)) if row == None: return result result.merge(delete_tag_by_link(name, id)) result.merge(delete_status_by_link(name, id)) result.merge(delete_vertex_by_link(name, id)) result.merge(cs) return result def expand_pattern_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.patterns where id = %s", (id,)) if row == None: return result result.merge(unset_reservoir_by_pattern(name, id)) result.merge(unset_pump_by_pattern(name, id)) result.merge(unset_pump_energy_by_pattern(name, id)) result.merge(unset_source_by_pattern(name, id)) result.merge(cs) return result def expand_curve_delete(name: str, cs: ChangeSet) -> ChangeSet: result = ChangeSet() id = cs.operations[0]['id'] row = try_read(name, "select 1 from network.curves where id = %s", (id,)) if row == None: return result result.merge(unset_tank_by_curve(name, id)) result.merge(unset_pump_by_curve(name, id)) result.merge(unset_pump_energy_by_curve(name, id)) result.merge(cs) return result def expand_v2_options_update(cs: ChangeSet) -> ChangeSet: cs.operations[0]['operation'] = API_UPDATE cs.operations[0]['type'] = 'option' new_cs = cs new_cs.merge(generate_v3(cs)) return new_cs def expand_v3_options_update(cs: ChangeSet) -> ChangeSet: cs.operations[0]['operation'] = API_UPDATE cs.operations[0]['type'] = 'option_v3' new_cs = cs new_cs.merge(generate_v2(cs)) return new_cs def expand_command(name: str, cs: ChangeSet) -> ChangeSet: op = cs.operations[0] operation = op['operation'] element_type = op['type'] if operation == API_DELETE: handler = _DELETE_REWRITERS.get(element_type) if handler: return handler(name, cs) elif operation == API_UPDATE: handler = _UPDATE_REWRITERS.get(element_type) if handler: return handler(cs) return cs DeleteRewriter = Callable[[str, ChangeSet], ChangeSet] UpdateRewriter = Callable[[ChangeSet], ChangeSet] _DELETE_REWRITERS: dict[str, DeleteRewriter] = { "junction": expand_junction_delete, "reservoir": expand_reservoir_delete, "tank": expand_tank_delete, "pipe": expand_pipe_delete, "pump": expand_pump_delete, "valve": expand_valve_delete, "pattern": expand_pattern_delete, "curve": expand_curve_delete, } _UPDATE_REWRITERS: dict[str, UpdateRewriter] = { "option": expand_v2_options_update, "option_v3": expand_v3_options_update, }