From 8104514cc5cd400409e01c557c638fc54d1fed86 Mon Sep 17 00:00:00 2001 From: Kapil Dadheech Date: Tue, 7 Jul 2026 11:44:53 +0530 Subject: [PATCH 1/4] Waterbody bug fix --- computing/zoi_layers/zoi1.py | 43 ++-- computing/zoi_layers/zoi2.py | 26 ++- computing/zoi_layers/zoi3.py | 34 +-- .../zoi_layers/zoi_ndvi_from_timeseries.py | 205 ++++++++++++++++++ waterrejuvenation/tasks.py | 115 ++++++++-- waterrejuvenation/utils.py | 131 ++++++++--- 6 files changed, 458 insertions(+), 96 deletions(-) create mode 100644 computing/zoi_layers/zoi_ndvi_from_timeseries.py diff --git a/computing/zoi_layers/zoi1.py b/computing/zoi_layers/zoi1.py index 80670092..5aa1c195 100644 --- a/computing/zoi_layers/zoi1.py +++ b/computing/zoi_layers/zoi1.py @@ -25,7 +25,6 @@ ee_initialize, valid_gee_text, get_gee_dir_path, - is_gee_asset_exists, export_vector_asset_to_gee, make_asset_public, check_task_status, @@ -36,6 +35,7 @@ calculate_zoi_area, wait_for_task_completion, delete_asset_on_GEE, + _waterbody_area_ha, ) from computing.surface_water_bodies.swb import sync_asset_to_db_and_geoserver @@ -50,6 +50,8 @@ def generate_zoi1( app_type="MWS", gee_account_id=None, proj_id=None, + start_date="2017-07-01", + end_date="2025-06-30", ): print("insdie zoi") ee_initialize(gee_account_id) @@ -74,15 +76,12 @@ def generate_zoi1( + description_zoi ) delete_asset_on_GEE(asset_id_zoi) - start_date = "2017-07-01" - end_date = "2025-06-30" zoi_fc = roi.map(compute_zoi) zoi_fc = ee.FeatureCollection(zoi_fc) zoi_rings = zoi_fc.filter(ee.Filter.gt("zoi_wb", 0)).map(create_ring) - if not is_gee_asset_exists(asset_id_zoi): - zoi_task = export_vector_asset_to_gee(zoi_rings, description_zoi, asset_id_zoi) - check_task_status([zoi_task]) - make_asset_public(asset_id_zoi) + zoi_task = export_vector_asset_to_gee(zoi_rings, description_zoi, asset_id_zoi) + check_task_status([zoi_task]) + make_asset_public(asset_id_zoi) if state and district and block: layer_name = f"waterbodies_zoi_{asset_suffix}" print(layer_name) @@ -103,13 +102,6 @@ def generate_zoi1( sync_project_fc_to_geoserver(zoi_rings, proj_obj.name, layer_name, "zoi_layers") -def _waterbody_area_ha(feature): - """Area in hectares; use geometry when area_ored is missing (e.g. water rej layers).""" - geom_area_ha = ee.Number(feature.geometry().area(maxError=1)).divide(10000) - stored_area = feature.get("area_ored") - return ee.Number(ee.Algorithms.If(stored_area, stored_area, geom_area_ha)) - - def compute_zoi(feature): area_of_wb = _waterbody_area_ha(feature) @@ -136,13 +128,32 @@ def y_large_bodies(area): .add(s.multiply(y_large_bodies(area_of_wb)).round()) ) - return feature.set("zoi_wb", zoi) + return feature.set("zoi_wb", zoi).set( + "UID", + ee.Algorithms.If( + feature.get("UID"), + feature.get("UID"), + ee.Algorithms.If( + feature.get("uid"), + feature.get("uid"), + feature.get("MWS_UID"), + ), + ), + ) def create_ring(feature): geom = feature.geometry() # can be point or polygon zoi = ee.Number(feature.get("zoi_wb")) - uid = feature.get("UID") + uid = ee.Algorithms.If( + feature.get("UID"), + feature.get("UID"), + ee.Algorithms.If( + feature.get("uid"), + feature.get("uid"), + feature.get("MWS_UID"), + ), + ) # Make circle buffer from centroid centroid = geom.centroid() diff --git a/computing/zoi_layers/zoi2.py b/computing/zoi_layers/zoi2.py index 03d44196..ad267320 100644 --- a/computing/zoi_layers/zoi2.py +++ b/computing/zoi_layers/zoi2.py @@ -19,6 +19,10 @@ def generate_zoi_ci( gee_account_id=None, proj_id=None, roi=None, + start_date="2017-07-01", + end_date="2025-06-30", + start_year=2017, + end_year=2024, ): from computing.cropping_intensity.cropping_intensity import ( generate_cropping_intensity, @@ -46,6 +50,14 @@ def generate_zoi_ci( + description_ci ) delete_asset_on_GEE(asset_id_ci) + description_zoi_ci = f"cropping_intensity_zoi_{asset_suffix}" + asset_id_zoi_ci = ( + get_gee_dir_path( + asset_folder_list, asset_path=GEE_PATHS[app_type]["GEE_ASSET_PATH"] + ) + + description_zoi_ci + ) + delete_asset_on_GEE(asset_id_zoi_ci) if roi: roi = ee.FeatureCollection(roi) else: @@ -56,20 +68,10 @@ def generate_zoi_ci( asset_folder_list=asset_folder_list, asset_suffix=asset_suffix, app_type=app_type, - start_year=2017, - end_year=2024, + start_year=start_year, + end_year=end_year, gee_account_id=gee_account_id, ) - start_date = "2017-07-01" - end_date = "2025-06-30" - description_zoi_ci = f"cropping_intensity_zoi_{asset_suffix}" - - asset_id_zoi_ci = ( - get_gee_dir_path( - asset_folder_list, asset_path=GEE_PATHS[app_type]["GEE_ASSET_PATH"] - ) - + description_zoi_ci - ) if state and district and block: layer_name = f"waterbodies_zoi_{asset_suffix}" layer_at_geoserver = sync_asset_to_db_and_geoserver( diff --git a/computing/zoi_layers/zoi3.py b/computing/zoi_layers/zoi3.py index fea3b962..2cab5e16 100644 --- a/computing/zoi_layers/zoi3.py +++ b/computing/zoi_layers/zoi3.py @@ -9,7 +9,7 @@ get_gee_dir_path, is_gee_asset_exists, ) -from waterrejuvenation.utils import wait_for_task_completion +from waterrejuvenation.utils import wait_for_task_completion, resolve_zoi_ndvi_input import ee @@ -20,7 +20,10 @@ def get_ndvi_for_zoi( zoi_roi=None, asset_suffix=None, asset_folder_list=None, - start_year="2017-23", + start_date="2017-07-01", + end_date="2025-06-30", + start_year=2017, + end_year=2024, app_type="MWS", gee_account_id=None, proj_id=None, @@ -29,23 +32,9 @@ def get_ndvi_for_zoi( ee_initialize(gee_account_id) from waterrejuvenation.utils import get_ndvi_data - if not proj_id: - description_zoi = "cropping_intensity_zoi_" + asset_suffix - asset_id_zoi = ( - get_gee_dir_path( - asset_folder_list, asset_path=GEE_PATHS[app_type]["GEE_ASSET_PATH"] - ) - + description_zoi - ) - else: - - description_zoi = "cropping_intensity_zoi_" + asset_suffix - asset_id_zoi = ( - get_gee_dir_path( - asset_folder_list, asset_path=GEE_PATHS[app_type]["GEE_ASSET_PATH"] - ) - + description_zoi - ) + zoi_collections = resolve_zoi_ndvi_input( + asset_folder_list, app_type, asset_suffix, zoi_roi=zoi_roi + ) description_ndvi = asset_suffix ndvi_asset_path = ( @@ -55,15 +44,14 @@ def get_ndvi_for_zoi( + description_ndvi ) - zoi_collections = ee.FeatureCollection(asset_id_zoi) - fc = get_ndvi_data(zoi_collections, 2017, 2024, description_ndvi, ndvi_asset_path) + fc = get_ndvi_data( + zoi_collections, start_year, end_year, description_ndvi, ndvi_asset_path + ) task = ee.batch.Export.table.toAsset( collection=fc, description=description_ndvi, assetId=ndvi_asset_path ) task.start() wait_for_task_completion(task) - start_date = "30-06-2017" - end_date = "01-07-2024" if state and district and block: layer_name = f"waterbodies_zoi_{asset_suffix}" layer_at_geoserver = sync_asset_to_db_and_geoserver( diff --git a/computing/zoi_layers/zoi_ndvi_from_timeseries.py b/computing/zoi_layers/zoi_ndvi_from_timeseries.py new file mode 100644 index 00000000..a8a4cc62 --- /dev/null +++ b/computing/zoi_layers/zoi_ndvi_from_timeseries.py @@ -0,0 +1,205 @@ +""" +ZOI NDVI using the ndvi_timeseries compute (HLS-interpolated NDVI). + +Uses the fast band-stack + single reduceRegions path (same idea as ndvi_time_series._generate_ndvi), +while keeping the original NDVI_ JSON property output shape. +""" + +import ee + +from computing.misc.hls_interpolated_ndvi import get_padded_ndvi_ts_image +from computing.surface_water_bodies.swb import sync_asset_to_db_and_geoserver +from computing.utils import sync_project_fc_to_geoserver +from projects.models import Project +from utilities.constants import GEE_PATHS +from utilities.gee_utils import ( + ee_initialize, + get_gee_dir_path, + check_task_status, + is_gee_asset_exists, + export_vector_asset_to_gee, +) +from waterrejuvenation.utils import wait_for_task_completion, resolve_zoi_ndvi_input + + +def _merge_ndvi_year_assets(chunk_assets): + """Same merge logic as waterrejuvenation.utils.merge_assets_chunked_on_year.""" + + def merge_features(feature): + uid = feature.get("UID") + matched_features = [] + for i in range(1, len(chunk_assets)): + matched_feature = ee.Feature( + ee.FeatureCollection(chunk_assets[i]) + .filter(ee.Filter.eq("UID", uid)) + .first() + ) + matched_features.append(matched_feature) + + merged_properties = feature.toDictionary() + for f in matched_features: + merged_properties = merged_properties.combine( + f.toDictionary(), overwrite=False + ) + + return ee.Feature(feature.geometry(), merged_properties) + + return ee.FeatureCollection(chunk_assets[0]).map(merge_features) + + +def _ndvi_bands_to_json_property(feature, year): + """Convert toBands reduceRegions output into NDVI_ JSON string.""" + + props = feature.toDictionary() + keys = props.keys().filter(ee.Filter.stringContains("item", "20")) + + def build_dict(k, acc): + k = ee.String(k) + # toBands prefixes band names with "_" — drop that index + new_key = k.split("_").slice(1).join("_") + return ee.Dictionary(acc).set(new_key, props.get(k)) + + ndvi_dict = ee.Dictionary(keys.iterate(build_dict, ee.Dictionary({}))) + ndvi_json = ee.String.encodeJSON(ndvi_dict) + return feature.set(f"NDVI_{year}", ndvi_json) + + +def build_ndvi_timeseries_from_timeseries_compute( + suitability_vector, + start_year, + end_year, + description, + asset_id, + ndvi_interval_days=16, + reducer_scale=30, + tile_scale=4, +): + """ + Build NDVI time series for ZOI polygons (NDVI_ JSON per feature). + + Fast path: one reduceRegions per hydrological year (not per timestep). + """ + feature_count = suitability_vector.size().getInfo() + print("total") + print(feature_count) + if feature_count == 0: + raise ValueError("Cannot generate NDVI: suitability vector has 0 features") + + task_ids = [] + asset_ids = [] + year = start_year + + while year <= end_year: + start_date = f"{year}-07-01" + end_date = f"{year + 1}-06-30" + ndvi_description = f"ndvi_{year}_{description}" + ndvi_asset_id = f"{asset_id}_ndvi_{year}" + + if is_gee_asset_exists(ndvi_asset_id): + ee.data.deleteAsset(ndvi_asset_id) + + ndvi = get_padded_ndvi_ts_image( + start_date, end_date, suitability_vector.bounds(), ndvi_interval_days + ) + + # Stack all dates into one image, then a single reduceRegions (fast) + ndvi_stack = ndvi.map( + lambda img: img.select("gapfilled_NDVI_lsc").rename( + img.date().format("YYYY-MM-dd") + ) + ).toBands() + + reduced = ndvi_stack.reduceRegions( + collection=suitability_vector, + reducer=ee.Reducer.mean(), + scale=reducer_scale, + tileScale=tile_scale, + ) + + merged_fc = reduced.map(lambda f: _ndvi_bands_to_json_property(f, year)) + + try: + task_id = export_vector_asset_to_gee( + merged_fc, ndvi_description, ndvi_asset_id + ) + if task_id: + print(f"Started export for {year}") + asset_ids.append(ndvi_asset_id) + task_ids.append(task_id) + except Exception as e: + print("Export error:", e) + + year += 1 + + check_task_status(task_ids) + return _merge_ndvi_year_assets(asset_ids) + + +def get_ndvi_for_zoi_from_timeseries_compute( + state=None, + district=None, + block=None, + zoi_roi=None, + asset_suffix=None, + asset_folder_list=None, + start_date="2017-07-01", + end_date="2025-06-30", + start_year=2017, + end_year=2024, + app_type="MWS", + gee_account_id=None, + proj_id=None, +): + """ + Same orchestration and result as get_ndvi_for_zoi (zoi3), but NDVI values come from + the ndvi_timeseries compute (get_padded_ndvi_ts_image) instead of get_ndvi_data. + """ + print("started generating ndvi for zoi (timeseries compute, fast reduce)") + ee_initialize(gee_account_id) + + zoi_collections = resolve_zoi_ndvi_input( + asset_folder_list, app_type, asset_suffix, zoi_roi=zoi_roi + ) + + description_ndvi = asset_suffix + ndvi_asset_path = ( + get_gee_dir_path( + asset_folder_list, asset_path=GEE_PATHS[app_type]["GEE_ASSET_PATH"] + ) + + description_ndvi + ) + + fc = build_ndvi_timeseries_from_timeseries_compute( + zoi_collections, + start_year, + end_year, + description_ndvi, + ndvi_asset_path, + ) + + task = ee.batch.Export.table.toAsset( + collection=fc, description=description_ndvi, assetId=ndvi_asset_path + ) + task.start() + wait_for_task_completion(task) + + if state and district and block: + layer_name = f"waterbodies_zoi_{asset_suffix}" + sync_asset_to_db_and_geoserver( + ndvi_asset_path, + layer_name, + asset_suffix, + start_date, + end_date, + state, + district, + block, + ) + elif proj_id: + proj_obj = Project.objects.get(pk=proj_id) + layer_name = f"waterbodies_zoi_{asset_suffix}" + sync_project_fc_to_geoserver( + fc, proj_obj.name, layer_name, "zoi_layers" + ) + + return fc diff --git a/waterrejuvenation/tasks.py b/waterrejuvenation/tasks.py index 290669fd..36859c52 100644 --- a/waterrejuvenation/tasks.py +++ b/waterrejuvenation/tasks.py @@ -75,10 +75,18 @@ def _gee_safe_property(value): return float(value) if isinstance(value, bool): return value + if isinstance(value, str): + stripped = value.strip() + if stripped and stripped.upper() not in ("N/A", "NAN", "NONE"): + try: + return float(stripped) + except ValueError: + pass return str(value) def _gdf_to_ee_feature_collection(gdf): + gdf = _normalize_uid_in_gdf(gdf) features = [] for _, row in gdf.iterrows(): geom = row.geometry @@ -89,10 +97,40 @@ def _gdf_to_ee_feature_collection(gdf): for col in gdf.columns if col != "geometry" } + # Never let a stray geometry property override the row geometry column. + props.pop("geometry", None) features.append(ee.Feature(ee.Geometry(geom.__geo_interface__), props)) return ee.FeatureCollection(features) +def _extract_uid(row): + """Read waterbody UID from a GeoDataFrame row (handles casing / aliases).""" + if row is None: + return None + for key in ("UID", "uid", "Uid", "MWS_UID", "mws_uid"): + if key not in row.index: + continue + val = row[key] + if val is None or (isinstance(val, float) and pd.isna(val)): + continue + text = str(val).strip() + if text and text.upper() not in ("N/A", "NAN", "NONE"): + return text + return None + + +def _normalize_uid_in_gdf(gdf): + """Ensure a canonical UID column exists for downstream ZOI / API merges.""" + if gdf.empty or "UID" in gdf.columns: + return gdf + for key in ("uid", "Uid", "MWS_UID", "mws_uid"): + if key in gdf.columns: + gdf = gdf.copy() + gdf["UID"] = gdf[key] + return gdf + return gdf + + def _fc_to_gdf(feature_collection): info = feature_collection.getInfo() features = info.get("features") or [] @@ -100,7 +138,51 @@ def _fc_to_gdf(feature_collection): return gpd.GeoDataFrame( columns=["geometry"], geometry="geometry", crs="EPSG:4326" ) - return gpd.GeoDataFrame.from_features(features, crs="EPSG:4326") + gdf = gpd.GeoDataFrame.from_features(features, crs="EPSG:4326") + return _normalize_uid_in_gdf(gdf) + + +def _swb_polygon_geometry(wb_row): + """Return SWB polygon geometry for matched water-rejuvenation exports.""" + geom = wb_row.geometry + if geom is None or geom.is_empty: + return None + if geom.geom_type in ("Polygon", "MultiPolygon"): + return geom + logger.warning( + "SWB feature %s has geometry type %s; expected Polygon/MultiPolygon.", + wb_row.name, + geom.geom_type, + ) + return geom + + +def _merge_matched_desilt_with_swb(desilt_row, wb_row): + """Matched export: SWB polygon geometry + SWB and desilting properties.""" + geometry = _swb_polygon_geometry(wb_row) + if geometry is None: + return None + + wb_props = { + col: wb_row[col] + for col in wb_row.index + if col != "geometry" + } + desilt_props = { + col: desilt_row[col] + for col in desilt_row.index + if col != "geometry" and col.lower() not in ("uid", "mws_uid") + } + props = {**wb_props, **desilt_props, "matched": True} + uid = _extract_uid(wb_row) or _extract_uid(desilt_row) + if uid: + props["UID"] = uid + elif "UID" not in props: + logger.warning( + "Matched waterbody has no UID (wb_index=%s); ZOI merge may fail.", + wb_row.name, + ) + return {**props, "geometry": geometry} def _match_desilting_points_to_waterbodies(desilt_gdf, wb_gdf, max_distance_m=100): @@ -122,28 +204,29 @@ def _match_desilting_points_to_waterbodies(desilt_gdf, wb_gdf, max_distance_m=10 for idx, point_row in desilt_gdf.iterrows(): point_metric = desilt_metric.loc[idx].geometry - props = {col: point_row[col] for col in desilt_gdf.columns if col != "geometry"} + desilt_props = { + col: point_row[col] for col in desilt_gdf.columns if col != "geometry" + } intersect_hits = wb_metric[wb_metric.intersects(point_metric)] if not intersect_hits.empty: - wb_idx = intersect_hits.index[0] - out_geom = wb_gdf.loc[wb_idx].geometry - props["matched"] = True - props["match_type"] = "intersect" - matched_rows.append({**props, "geometry": out_geom}) + wb_row = wb_gdf.loc[intersect_hits.index[0]] + row = _merge_matched_desilt_with_swb(point_row, wb_row) + if row: + row["match_type"] = "intersect" + matched_rows.append(row) continue near_hits = wb_metric[wb_metric.intersects(point_metric.buffer(max_distance_m))] if not near_hits.empty: - wb_idx = near_hits.index[0] - out_geom = wb_gdf.loc[wb_idx].geometry - props["matched"] = True - props["match_type"] = "near" - matched_rows.append({**props, "geometry": out_geom}) + wb_row = wb_gdf.loc[near_hits.index[0]] + row = _merge_matched_desilt_with_swb(point_row, wb_row) + if row: + row["match_type"] = "near" + matched_rows.append(row) continue - props["matched"] = False - props["match_type"] = "none" + props = {**desilt_props, "matched": False, "match_type": "none"} unmatched_rows.append({**props, "geometry": point_row.geometry}) matched_gdf = ( @@ -983,6 +1066,10 @@ def BuildWaterBodyLayer( matched_gdf, unmatched_gdf = _match_desilting_points_to_waterbodies( desilt_gdf, wb_gdf, max_distance_m=100 ) + matched_gdf = _normalize_uid_in_gdf(matched_gdf) + if len(matched_gdf) > 0: + geom_types = matched_gdf.geometry.geom_type.value_counts().to_dict() + logger.info("BuildWaterBodyLayer matched geometry types: %s", geom_types) matched_fc = _gdf_to_ee_feature_collection(matched_gdf) unmatched_fc = _gdf_to_ee_feature_collection(unmatched_gdf) diff --git a/waterrejuvenation/utils.py b/waterrejuvenation/utils.py index b61fb7b7..60c2b249 100644 --- a/waterrejuvenation/utils.py +++ b/waterrejuvenation/utils.py @@ -340,7 +340,15 @@ def create_ring(feature): zoi = ee.Number(feature.get("zoi_wb")) waterbody_name = feature.get("waterbody_name") impactfull = feature.get("impactful") - uid = feature.get("UID") + uid = ee.Algorithms.If( + feature.get("UID"), + feature.get("UID"), + ee.Algorithms.If( + feature.get("uid"), + feature.get("uid"), + feature.get("MWS_UID"), + ), + ) # Make circle buffer from centroid centroid = geom.centroid() @@ -403,6 +411,72 @@ def calculate_zoi_area(zoi): return ee.Number.parse(area_hectares.format("%.2f")) +def ensure_uid_on_fc(fc): + """Ensure every feature has a canonical UID for NDVI / merge steps.""" + + def ensure_uid(feature): + uid = ee.Algorithms.If( + feature.get("UID"), + feature.get("UID"), + ee.Algorithms.If( + feature.get("uid"), + feature.get("uid"), + ee.Algorithms.If( + feature.get("MWS_UID"), + feature.get("MWS_UID"), + feature.get("system:index"), + ), + ), + ) + return feature.set("UID", uid) + + return fc.map(ensure_uid) + + +def resolve_zoi_ndvi_input( + asset_folder_list, app_type, asset_suffix, zoi_roi=None +): + """ + Pick polygons for ZOI NDVI: prefer cropping_intensity_zoi, else zoi_ rings. + """ + if zoi_roi: + fc = ee.FeatureCollection(zoi_roi) + count = fc.size().getInfo() + logger.info("NDVI input: explicit roi (%s features)", count) + if count == 0: + raise ValueError(f"NDVI input roi is empty: {zoi_roi}") + return ensure_uid_on_fc(fc) + + base = get_gee_dir_path( + asset_folder_list, asset_path=GEE_PATHS[app_type]["GEE_ASSET_PATH"] + ) + ci_asset = base + f"cropping_intensity_zoi_{asset_suffix}" + zoi_asset = base + f"zoi_{asset_suffix}" + + ci_fc = ee.FeatureCollection(ci_asset) + ci_size = ci_fc.size().getInfo() + if ci_size > 0: + logger.info( + "NDVI input: cropping_intensity_zoi (%s features) from %s", + ci_size, + ci_asset, + ) + return ensure_uid_on_fc(ci_fc) + + zoi_fc = ee.FeatureCollection(zoi_asset) + zoi_size = zoi_fc.size().getInfo() + logger.info( + "NDVI input: fallback zoi_ (%s features) from %s", + zoi_size, + zoi_asset, + ) + if zoi_size == 0: + raise ValueError( + f"No ZOI features for NDVI (tried {ci_asset} and {zoi_asset})" + ) + return ensure_uid_on_fc(zoi_fc) + + def get_ndvi_data(suitability_vector, start_year, end_year, description, asset_id): """ Extracts and exports NDVI data for a set of features by aggregating NDVI values @@ -421,6 +495,13 @@ def get_ndvi_data(suitability_vector, start_year, end_year, description, asset_i Returns: ee.FeatureCollection: Merged NDVI time series across years. """ + suitability_vector = ensure_uid_on_fc(suitability_vector) + feature_count = suitability_vector.size().getInfo() + print("total") + print(feature_count) + if feature_count == 0: + raise ValueError("Cannot generate NDVI: suitability vector has 0 features") + task_ids = [] asset_ids = [] # Loop over each year @@ -465,39 +546,18 @@ def annotate(feature): # Map image-wise extraction and flatten to a single FeatureCollection all_ndvi = ndvi.map(map_image).flatten() - # Extract all unique UIDs from the input feature collection - uids = suitability_vector.aggregate_array("UID") - count = uids.size() # Server-side count - print("total") - print(count.getInfo()) - - # For each UID, filter NDVI features and aggregate to dict - def build_feature(uid): - """ - Reconstruct a single feature by merging its NDVI values across all images - into one property NDVI_ as a JSON dictionary {date: value}. - """ - # Get the geometry and properties of the original feature - feature_geom = ee.Feature( - suitability_vector.filter(ee.Filter.eq("UID", uid)).first() - ) - - # Filter all NDVI records related to this UID + # Reconstruct one row per input polygon with NDVI_ JSON property + def build_feature(input_feature): + uid = input_feature.get("UID") filtered = all_ndvi.filter(ee.Filter.eq("UID", uid)) - - # Create dictionary: {date: ndvi} date_ndvi_list = filtered.aggregate_array("ndvi_date").zip( filtered.aggregate_array("ndvi") ) - - # Convert to dictionary and encode as JSON string ndvi_dict = ee.Dictionary(date_ndvi_list.flatten()) ndvi_json = ee.String.encodeJSON(ndvi_dict) + return input_feature.set(f"NDVI_{start_year}", ndvi_json) - return feature_geom.set(f"NDVI_{start_year}", ndvi_json) - - # Apply feature-wise aggregation - merged_fc = ee.FeatureCollection(uids.map(build_feature)) + merged_fc = suitability_vector.map(build_feature) # Export as single-row-per-feature collection try: @@ -562,7 +622,9 @@ def get_ndvi_for_zoi( + asset_suffix_ndvi ) - zoi_collections = ee.FeatureCollection(zoi_asset_path) + zoi_collections = resolve_zoi_ndvi_input( + asset_folder, app_type, f"{proj_obj.name}_{proj_obj.id}".lower(), zoi_roi=zoi_asset_path + ) fc = get_ndvi_data(zoi_collections, 2017, 2024, asset_suffix_ndvi, ndvi_asset_path) task = ee.batch.Export.table.toAsset( @@ -607,10 +669,17 @@ def merge_features(feature): def _waterbody_area_ha(feature): - """Area in hectares; use geometry when area_ored is missing (e.g. water rej layers).""" + """Area in hectares; use geometry when area_ored is missing or non-numeric.""" geom_area_ha = ee.Number(feature.geometry().area(maxError=1)).divide(10000) - stored_area = feature.get("area_ored") - return ee.Number(ee.Algorithms.If(stored_area, stored_area, geom_area_ha)) + raw_area = feature.get("area_ored") + # Water-rej exports may store area_ored as a string; coerce before numeric ops. + stored_str = ee.String(ee.Algorithms.If(raw_area, raw_area, "0")) + is_na = stored_str.compareTo("N/A").eq(0) + safe_str = ee.String(ee.Algorithms.If(is_na, "0", stored_str)) + stored_numeric = ee.Number.parse(safe_str) + return ee.Number( + ee.Algorithms.If(stored_numeric.gt(0), stored_numeric, geom_area_ha) + ) def compute_zoi(feature): From 1696c3210a811a23607672d10f59485e446412e2 Mon Sep 17 00:00:00 2001 From: Kapil Dadheech Date: Wed, 8 Jul 2026 11:14:40 +0530 Subject: [PATCH 2/4] fix uuid issue --- waterrejuvenation/tasks.py | 187 ++++++++++++++++++++++++++++++------- waterrejuvenation/utils.py | 59 ++++++++++++ 2 files changed, 214 insertions(+), 32 deletions(-) diff --git a/waterrejuvenation/tasks.py b/waterrejuvenation/tasks.py index 36859c52..7790329c 100644 --- a/waterrejuvenation/tasks.py +++ b/waterrejuvenation/tasks.py @@ -40,6 +40,8 @@ wait_for_task_completion, delete_asset_on_GEE, find_nearest_water_pixel, + format_waterbody_uid_value, + id_text, ) from computing.surface_water_bodies.swb import generate_swb_layer @@ -64,70 +66,185 @@ def is_nan(value): ) -def _gee_safe_property(value): +def _is_string_id_property(prop_name): + """Identifier / label fields must stay strings in GEE + GeoServer exports.""" + lowered = str(prop_name).lower() + if lowered in { + "uid", + "mws_uid", + "desilt_id", + "pond_id", + "village_id", + "waterbody_name", + "village", + "state", + "district", + "taluka", + "tehsil", + "block", + }: + return True + return lowered.endswith("_uid") or lowered.endswith("_id") + + +def _gee_safe_property(value, prop_name=None): + import json import numpy as np + from datetime import date, datetime + from shapely.geometry.base import BaseGeometry if is_nan(value): - return "N/A" + return None + if prop_name and _is_string_id_property(prop_name): + text = id_text(value) + return text + if isinstance(value, BaseGeometry): + return value.wkt + if isinstance(value, (np.bool_,)): + return bool(value) if isinstance(value, (np.integer,)): return int(value) if isinstance(value, (np.floating,)): return float(value) if isinstance(value, bool): return value + if isinstance(value, int): + return value + if isinstance(value, float): + return value + if isinstance(value, (datetime, date, pd.Timestamp)): + return value.isoformat() + if isinstance(value, (list, tuple, set)): + return json.dumps( + [_gee_safe_property(item) for item in value], + default=str, + ) + if isinstance(value, dict): + # GeoJSON Feature/Geometry dicts cannot be stored as table properties. + if value.get("type") in { + "Feature", + "FeatureCollection", + "Point", + "MultiPoint", + "LineString", + "MultiLineString", + "Polygon", + "MultiPolygon", + "GeometryCollection", + }: + return json.dumps(value) + return json.dumps( + {str(k): _gee_safe_property(v) for k, v in value.items()}, + default=str, + ) if isinstance(value, str): stripped = value.strip() - if stripped and stripped.upper() not in ("N/A", "NAN", "NONE"): + if not stripped or stripped.upper() in ("N/A", "NAN", "NONE"): + return None + if prop_name and prop_name.lower() in {"area_ored", "zoi", "zoi_wb", "zoi_area"}: try: return float(stripped) except ValueError: - pass + return stripped + return stripped return str(value) +def _sync_gdf_to_project_geoserver(gdf, project_name, layer_name, workspace): + """Push a local GeoDataFrame to GeoServer without an EE getInfo round-trip.""" + import os + + import geopandas as gpd + + from computing.utils import fix_invalid_geometry_in_gdf, push_shape_to_geoserver + + if gdf is None or gdf.empty: + logger.warning("No features to sync for layer %s", layer_name) + return None + + state_dir = os.path.join("data/fc_to_shape", project_name) + os.makedirs(state_dir, exist_ok=True) + path = os.path.join(state_dir, layer_name) + out_gdf = gpd.GeoDataFrame(gdf.copy(), geometry="geometry", crs="EPSG:4326") + out_gdf = fix_invalid_geometry_in_gdf(out_gdf) + out_gdf.to_file(path + ".gpkg", driver="GPKG") + return push_shape_to_geoserver( + path, workspace=workspace, layer_name=layer_name, file_type="gpkg" + ) + + def _gdf_to_ee_feature_collection(gdf): gdf = _normalize_uid_in_gdf(gdf) - features = [] + geojson_features = [] for _, row in gdf.iterrows(): geom = row.geometry if geom is None or geom.is_empty: continue - props = { - col: _gee_safe_property(row[col]) - for col in gdf.columns - if col != "geometry" - } - # Never let a stray geometry property override the row geometry column. - props.pop("geometry", None) - features.append(ee.Feature(ee.Geometry(geom.__geo_interface__), props)) - return ee.FeatureCollection(features) + props = {} + for col in gdf.columns: + if col == "geometry": + continue + safe_val = _gee_safe_property(row[col], prop_name=col) + if safe_val is not None: + props[col] = safe_val + geojson_features.append( + { + "type": "Feature", + "geometry": geom.__geo_interface__, + "properties": props, + } + ) + if not geojson_features: + return ee.FeatureCollection([]) + return ee.FeatureCollection(geojson_features) def _extract_uid(row): """Read waterbody UID from a GeoDataFrame row (handles casing / aliases).""" if row is None: return None - for key in ("UID", "uid", "Uid", "MWS_UID", "mws_uid"): - if key not in row.index: - continue - val = row[key] - if val is None or (isinstance(val, float) and pd.isna(val)): - continue - text = str(val).strip() - if text and text.upper() not in ("N/A", "NAN", "NONE"): - return text - return None + mws_val = None + uid_val = None + for key in ("MWS_UID", "mws_uid"): + if key in row.index: + mws_val = row[key] + break + for key in ("UID", "uid", "Uid"): + if key in row.index: + uid_val = row[key] + break + uid, _ = format_waterbody_uid_value(mws_val, uid_val) + return uid def _normalize_uid_in_gdf(gdf): - """Ensure a canonical UID column exists for downstream ZOI / API merges.""" - if gdf.empty or "UID" in gdf.columns: + """Ensure canonical UID / MWS_UID columns exist with underscore format.""" + if gdf.empty: return gdf + gdf = gdf.copy() for key in ("uid", "Uid", "MWS_UID", "mws_uid"): - if key in gdf.columns: - gdf = gdf.copy() + if key in gdf.columns and "UID" not in gdf.columns: gdf["UID"] = gdf[key] - return gdf + break + if "UID" not in gdf.columns: + return gdf + + mws_col = next( + (col for col in ("MWS_UID", "mws_uid") if col in gdf.columns), + None, + ) + + def _normalize_row(row): + mws_val = row[mws_col] if mws_col else None + uid, mws_val = format_waterbody_uid_value(mws_val, row["UID"]) + row = row.copy() + if uid is not None: + row["UID"] = uid + if mws_val is not None and mws_col: + row[mws_col] = mws_val + return row + + gdf = gdf.apply(_normalize_row, axis=1) return gdf @@ -174,7 +291,13 @@ def _merge_matched_desilt_with_swb(desilt_row, wb_row): if col != "geometry" and col.lower() not in ("uid", "mws_uid") } props = {**wb_props, **desilt_props, "matched": True} - uid = _extract_uid(wb_row) or _extract_uid(desilt_row) + mws_val = wb_row["MWS_UID"] if "MWS_UID" in wb_row.index else None + if mws_val is None and "mws_uid" in wb_row.index: + mws_val = wb_row["mws_uid"] + uid_val = wb_row["UID"] if "UID" in wb_row.index else None + uid, mws_val = format_waterbody_uid_value(mws_val, uid_val) + if mws_val is not None: + props["MWS_UID"] = mws_val if uid: props["UID"] = uid elif "UID" not in props: @@ -1135,8 +1258,8 @@ def BuildWaterBodyLayer( # ------------------------------------------------------------------ layer_name = f"waterbodies_{proj_obj.name}_{proj_obj.id}".lower() if len(matched_gdf) > 0: - sync_project_fc_to_geoserver( - matched_fc, + _sync_gdf_to_project_geoserver( + matched_gdf, proj_obj.name, layer_name, "swb", diff --git a/waterrejuvenation/utils.py b/waterrejuvenation/utils.py index 60c2b249..96ba1737 100644 --- a/waterrejuvenation/utils.py +++ b/waterrejuvenation/utils.py @@ -411,6 +411,65 @@ def calculate_zoi_area(zoi): return ee.Number.parse(area_hectares.format("%.2f")) +def is_nan_value(value): + return ( + value is None + or (isinstance(value, float) and math.isnan(value)) + or pd.isna(value) + ) + + +def id_text(value): + """Normalize an identifier value to a clean string.""" + if is_nan_value(value): + return None + if isinstance(value, float) and value.is_integer(): + value = int(value) + if isinstance(value, int): + text = str(value) + else: + text = str(value).strip() + if not text or text.upper() in ("N/A", "NAN", "NONE"): + return None + if text.endswith(".0"): + stem = text[:-2] + if stem.replace("_", "").isdigit(): + text = stem + if "e" in text.lower(): + try: + as_float = float(text) + if as_float.is_integer(): + text = str(int(as_float)) + except ValueError: + pass + return text + + +def format_waterbody_uid_value(mws_uid, uid): + """ + Restore SWB UID format: _, e.g. 25_34833_101. + + GEE/pandas may coerce this to a number (12129065580), dropping the underscore + before the per-waterbody index (12129065_580). + """ + mws_str = id_text(mws_uid) + uid_str = id_text(uid) + if not uid_str: + return None, mws_str + if "_" in uid_str: + return uid_str, mws_str + if not mws_str: + return uid_str, mws_str + + mws_digits = mws_str.replace("_", "") + uid_digits = uid_str.replace("_", "") + if uid_digits.startswith(mws_digits) and len(uid_digits) > len(mws_digits): + suffix = uid_digits[len(mws_digits) :] + if suffix.isdigit(): + return f"{mws_str}_{suffix}", mws_str + return uid_str, mws_str + + def ensure_uid_on_fc(fc): """Ensure every feature has a canonical UID for NDVI / merge steps.""" From 331f6da36903c4b48d173121146eee25316dcc86 Mon Sep 17 00:00:00 2001 From: Kapil Dadheech Date: Mon, 13 Jul 2026 14:34:54 +0530 Subject: [PATCH 3/4] Require explicit ZOI start/end dates instead of hardcoded years. Reject API and water-rej compute flows when the analysis window is missing so date ranges stay caller-provided year to year. Co-authored-by: Cursor --- computing/api.py | 56 +++++++++++++++-- computing/zoi_layers/zoi.py | 60 +++++++++++++++++-- computing/zoi_layers/zoi1.py | 8 ++- computing/zoi_layers/zoi2.py | 13 ++-- computing/zoi_layers/zoi3.py | 12 ++-- .../zoi_layers/zoi_ndvi_from_timeseries.py | 12 ++-- waterrejuvenation/models.py | 10 ++++ waterrejuvenation/tasks.py | 23 +++++++ waterrejuvenation/views.py | 29 +++++++-- 9 files changed, 196 insertions(+), 27 deletions(-) diff --git a/computing/api.py b/computing/api.py index 84155f5a..b913db7d 100644 --- a/computing/api.py +++ b/computing/api.py @@ -1512,25 +1512,71 @@ def generate_ndvi_timeseries(request): def generate_zoi_to_gee(request): print("Inside generate zoi layers") try: - state = request.data.get("state").lower() - district = request.data.get("district").lower() - block = request.data.get("block").lower() + state = request.data.get("state") + district = request.data.get("district") + block = request.data.get("block") gee_account_id = request.data.get("gee_account_id") + start_date = request.data.get("start_date") or request.data.get("startDate") + end_date = request.data.get("end_date") or request.data.get("endDate") + start_year = request.data.get("start_year") or request.data.get("startYear") + end_year = request.data.get("end_year") or request.data.get("endYear") + + if not state or not district or not block: + return Response( + {"error": "state, district, and block are required."}, + status=status.HTTP_400_BAD_REQUEST, + ) + + # Allow start_year/end_year as hydrological-year shorthand. + if not start_date and start_year is not None: + start_date = f"{int(start_year)}-07-01" + if not end_date and end_year is not None: + end_date = f"{int(end_year) + 1}-06-30" + + if not start_date or not end_date: + return Response( + { + "error": ( + "start_date and end_date are required (YYYY-MM-DD), " + "or provide start_year and end_year (hydrological years)." + ) + }, + status=status.HTTP_400_BAD_REQUEST, + ) + + from computing.zoi_layers.zoi import _resolve_zoi_time_window + + try: + start_date, end_date, _, _ = _resolve_zoi_time_window(start_date, end_date) + except ValueError as exc: + return Response({"error": str(exc)}, status=status.HTTP_400_BAD_REQUEST) + + state = state.lower() + district = district.lower() + block = block.lower() + generate_zoi.apply_async( kwargs={ "state": state, "district": district, "block": block, "gee_account_id": gee_account_id, + "start_date": start_date, + "end_date": end_date, }, queue="waterbody", ) return Response( - {"Success": "Successfully initiated"}, status=status.HTTP_200_OK + { + "Success": "Successfully initiated", + "start_date": start_date, + "end_date": end_date, + }, + status=status.HTTP_200_OK, ) except Exception as e: - print("Exception in generate_mining_to_gee api :: ", e) + print("Exception in generate_zoi_to_gee api :: ", e) return Response({"Exception": e}, status=status.HTTP_500_INTERNAL_SERVER_ERROR) diff --git a/computing/zoi_layers/zoi.py b/computing/zoi_layers/zoi.py index b8be98b2..5a1fcc17 100644 --- a/computing/zoi_layers/zoi.py +++ b/computing/zoi_layers/zoi.py @@ -1,10 +1,47 @@ +from datetime import datetime + from computing.zoi_layers.zoi1 import generate_zoi1 from computing.zoi_layers.zoi2 import generate_zoi_ci from computing.zoi_layers.zoi3 import get_ndvi_for_zoi -from projects.models import Project -from utilities.gee_utils import ee_initialize, valid_gee_text, check_task_status -from waterrejuvenation.utils import wait_for_task_completion, delete_asset_on_GEE from nrm_app.celery import app +from projects.models import Project +from utilities.gee_utils import ee_initialize, valid_gee_text + + +def _resolve_zoi_time_window(start_date=None, end_date=None): + """ + Validate and normalize ZOI date window. + + Requires explicit start_date and end_date (YYYY-MM-DD). No year defaults — + callers must provide the analysis window. + """ + if not start_date or not end_date: + raise ValueError( + "start_date and end_date are required (YYYY-MM-DD). " + "Pass both parameters explicitly; no default date window is applied." + ) + + start_date = str(start_date).strip() + end_date = str(end_date).strip() + + try: + start_dt = datetime.strptime(start_date, "%Y-%m-%d") + end_dt = datetime.strptime(end_date, "%Y-%m-%d") + except ValueError as exc: + raise ValueError("start_date and end_date must be in YYYY-MM-DD format.") from exc + + if start_dt > end_dt: + raise ValueError("start_date must be less than or equal to end_date.") + + # ZOI CI/NDVI use hydrological years (July -> June), represented by start-year. + start_year = start_dt.year if start_dt.month >= 7 else start_dt.year - 1 + end_year = end_dt.year if end_dt.month >= 7 else end_dt.year - 1 + if start_year > end_year: + raise ValueError( + "Provided date window does not contain a valid hydrological year." + ) + + return start_date, end_date, start_year, end_year @app.task() @@ -18,6 +55,8 @@ def generate_zoi( app_type="MWS", gee_account_id=None, proj_id=None, + start_date=None, + end_date=None, ): print(f"gee account id {gee_account_id}") ee_initialize(gee_account_id) @@ -31,6 +70,10 @@ def generate_zoi( asset_folder_list = [proj_obj.name.lower()] asset_suffix = f"{proj_obj.name}_{proj_obj.id}".lower() + start_date, end_date, start_year, end_year = _resolve_zoi_time_window( + start_date, end_date + ) + generate_zoi1( state, district, @@ -41,6 +84,8 @@ def generate_zoi( app_type, gee_account_id, proj_id, + start_date=start_date, + end_date=end_date, ) generate_zoi_ci( @@ -52,10 +97,13 @@ def generate_zoi( app_type, gee_account_id, proj_id, + start_date=start_date, + end_date=end_date, + start_year=start_year, + end_year=end_year, ) if proj_id: - get_ndvi_for_zoi( state=state, district=district, @@ -65,4 +113,8 @@ def generate_zoi( app_type=app_type, gee_account_id=gee_account_id, proj_id=proj_id, + start_date=start_date, + end_date=end_date, + start_year=start_year, + end_year=end_year, ) diff --git a/computing/zoi_layers/zoi1.py b/computing/zoi_layers/zoi1.py index 5aa1c195..3f89d367 100644 --- a/computing/zoi_layers/zoi1.py +++ b/computing/zoi_layers/zoi1.py @@ -50,10 +50,14 @@ def generate_zoi1( app_type="MWS", gee_account_id=None, proj_id=None, - start_date="2017-07-01", - end_date="2025-06-30", + start_date=None, + end_date=None, ): print("insdie zoi") + if not start_date or not end_date: + raise ValueError( + "start_date and end_date are required for ZOI generation (YYYY-MM-DD)." + ) ee_initialize(gee_account_id) description = "swb3_" + asset_suffix asset_id = ( diff --git a/computing/zoi_layers/zoi2.py b/computing/zoi_layers/zoi2.py index ad267320..c7e44946 100644 --- a/computing/zoi_layers/zoi2.py +++ b/computing/zoi_layers/zoi2.py @@ -19,15 +19,20 @@ def generate_zoi_ci( gee_account_id=None, proj_id=None, roi=None, - start_date="2017-07-01", - end_date="2025-06-30", - start_year=2017, - end_year=2024, + start_date=None, + end_date=None, + start_year=None, + end_year=None, ): from computing.cropping_intensity.cropping_intensity import ( generate_cropping_intensity, ) + if not start_date or not end_date or start_year is None or end_year is None: + raise ValueError( + "start_date, end_date, start_year, and end_year are required for ZOI CI." + ) + if state and district and block: asset_suffix = ( valid_gee_text(district.lower()) + "_" + valid_gee_text(block.lower()) diff --git a/computing/zoi_layers/zoi3.py b/computing/zoi_layers/zoi3.py index 2cab5e16..afecc374 100644 --- a/computing/zoi_layers/zoi3.py +++ b/computing/zoi_layers/zoi3.py @@ -20,15 +20,19 @@ def get_ndvi_for_zoi( zoi_roi=None, asset_suffix=None, asset_folder_list=None, - start_date="2017-07-01", - end_date="2025-06-30", - start_year=2017, - end_year=2024, + start_date=None, + end_date=None, + start_year=None, + end_year=None, app_type="MWS", gee_account_id=None, proj_id=None, ): print("started generating ndvi") + if not start_date or not end_date or start_year is None or end_year is None: + raise ValueError( + "start_date, end_date, start_year, and end_year are required for ZOI NDVI." + ) ee_initialize(gee_account_id) from waterrejuvenation.utils import get_ndvi_data diff --git a/computing/zoi_layers/zoi_ndvi_from_timeseries.py b/computing/zoi_layers/zoi_ndvi_from_timeseries.py index a8a4cc62..3bf3a9d7 100644 --- a/computing/zoi_layers/zoi_ndvi_from_timeseries.py +++ b/computing/zoi_layers/zoi_ndvi_from_timeseries.py @@ -142,10 +142,10 @@ def get_ndvi_for_zoi_from_timeseries_compute( zoi_roi=None, asset_suffix=None, asset_folder_list=None, - start_date="2017-07-01", - end_date="2025-06-30", - start_year=2017, - end_year=2024, + start_date=None, + end_date=None, + start_year=None, + end_year=None, app_type="MWS", gee_account_id=None, proj_id=None, @@ -155,6 +155,10 @@ def get_ndvi_for_zoi_from_timeseries_compute( the ndvi_timeseries compute (get_padded_ndvi_ts_image) instead of get_ndvi_data. """ print("started generating ndvi for zoi (timeseries compute, fast reduce)") + if not start_date or not end_date or start_year is None or end_year is None: + raise ValueError( + "start_date, end_date, start_year, and end_year are required for ZOI NDVI." + ) ee_initialize(gee_account_id) zoi_collections = resolve_zoi_ndvi_input( diff --git a/waterrejuvenation/models.py b/waterrejuvenation/models.py index 3f4397fb..b6b4afd1 100644 --- a/waterrejuvenation/models.py +++ b/waterrejuvenation/models.py @@ -62,6 +62,9 @@ def __str__(self): def save(self, *args, **kwargs): """Override save to calculate file hash before saving""" + start_date = kwargs.pop("start_date", None) + end_date = kwargs.pop("end_date", None) + if not self.excel_hash and self.file: # Calculate hash for new file self.file.seek(0) @@ -74,6 +77,11 @@ def save(self, *args, **kwargs): print(f"is processing required: {self.is_processing_required}") print(f"is lullc required: {self.is_lulc_required}") if self.is_compute: + if self.is_processing_required and (not start_date or not end_date): + raise ValueError( + "start_date and end_date are required when is_compute and " + "is_processing_required are true (YYYY-MM-DD)." + ) Upload_Desilting_Points.apply_async( kwargs={ "file_obj_id": self.id, @@ -81,6 +89,8 @@ def save(self, *args, **kwargs): "is_lulc_required": self.is_lulc_required, "is_processing_required": self.is_processing_required, "is_closest_wp": self.is_closest_wp, + "start_date": start_date, + "end_date": end_date, }, queue="waterbody1", ) diff --git a/waterrejuvenation/tasks.py b/waterrejuvenation/tasks.py index 7790329c..0c897eed 100644 --- a/waterrejuvenation/tasks.py +++ b/waterrejuvenation/tasks.py @@ -376,6 +376,8 @@ def Upload_Desilting_Points( is_lulc_required=True, gee_account_id=None, is_processing_required=True, + start_date=None, + end_date=None, ): import pandas as pd from .models import WaterbodiesFileUploadLog, WaterbodiesDesiltingLog @@ -387,6 +389,12 @@ def normalize(val): return None return val + if is_processing_required and (not start_date or not end_date): + raise ValueError( + "start_date and end_date are required when processing desilting points " + "through ZOI (YYYY-MM-DD)." + ) + ee_initialize(gee_account_id) wb_obj = WaterbodiesFileUploadLog.objects.get(pk=file_obj_id) @@ -504,6 +512,8 @@ def normalize(val): is_lulc_required=is_lulc_required, gee_account_id=gee_account_id, proj_id=proj_obj.id, + start_date=start_date, + end_date=end_date, ) wb_obj.process = True @@ -515,7 +525,14 @@ def Generate_lulc_mws( is_lulc_required=True, gee_account_id=None, proj_id=None, + start_date=None, + end_date=None, ): + if not start_date or not end_date: + raise ValueError( + "start_date and end_date are required for water-rej ZOI generation " + "(YYYY-MM-DD)." + ) proj_obj = Project.objects.get(pk=proj_id) asset_suffix = f"{proj_obj.name}_{proj_obj.id}".lower() asset_folder = [proj_obj.name.lower()] @@ -579,6 +596,8 @@ def Generate_lulc_mws( asset_suffix=asset_suffix, asset_folder=asset_folder, app_type="WATERBODY", + start_date=start_date, + end_date=end_date, ) @@ -719,6 +738,8 @@ def Genereate_zoi_and_zoi_indicator( asset_suffix=None, asset_folder=None, roi=None, + start_date=None, + end_date=None, ): print(f"roi: {roi}") ee_initialize(gee_project_id) @@ -737,6 +758,8 @@ def Genereate_zoi_and_zoi_indicator( app_type=app_type, gee_account_id=gee_project_id, proj_id=proj_id, + start_date=start_date, + end_date=end_date, ) diff --git a/waterrejuvenation/views.py b/waterrejuvenation/views.py index aed9559d..90dcc232 100644 --- a/waterrejuvenation/views.py +++ b/waterrejuvenation/views.py @@ -113,17 +113,36 @@ def create(self, request, *args, **kwargs): continue # Prepare data for serializer + is_compute = request.data.get("is_compute", False) + is_processing_required = request.data.get("is_processing_required", True) + start_date = request.data.get("start_date") or request.data.get("startDate") + end_date = request.data.get("end_date") or request.data.get("endDate") + + def _as_bool(val, default=False): + if val is None: + return default + if isinstance(val, bool): + return val + return str(val).lower() in ("true", "1", "yes") + + compute_enabled = _as_bool(is_compute, False) + processing_enabled = _as_bool(is_processing_required, True) + if compute_enabled and processing_enabled and (not start_date or not end_date): + errors.append( + "start_date and end_date are required when compute processing " + "is enabled (YYYY-MM-DD)." + ) + continue + data = { "name": request.data.get("name", valid_gee_text(uploaded_file.name)), "file": uploaded_file, "project": project.id, "gee_account_id": request.data.get("gee_account_id"), "is_lulc_required": request.data.get("is_lulc_required", True), - "is_processing_required": request.data.get( - "is_processing_required", True - ), + "is_processing_required": is_processing_required, "is_closest_wp": request.data.get("is_closest_wp", True), - "is_compute": request.data.get("is_compute", False), + "is_compute": is_compute, } print(data) serializer = self.get_serializer(data=data) @@ -150,6 +169,8 @@ def create(self, request, *args, **kwargs): project=project, uploaded_by=request.user, excel_hash=excel_hash, + start_date=start_date, + end_date=end_date, ) # # Convert KML to GeoJSON From 71648ff3689243464a326f13e9916148ba2517a5 Mon Sep 17 00:00:00 2001 From: Kapil Dadheech Date: Mon, 13 Jul 2026 14:47:27 +0530 Subject: [PATCH 4/4] Remove remaining year hardcoding from waterrej precipitation and nearest-water paths. Co-authored-by: Cursor --- computing/utils.py | 9 ++++- waterrejuvenation/tasks.py | 76 ++++++++++++++++++++++++++------------ waterrejuvenation/utils.py | 46 ++++++++++++++--------- 3 files changed, 89 insertions(+), 42 deletions(-) diff --git a/computing/utils.py b/computing/utils.py index 1d7a73b3..2faf8919 100644 --- a/computing/utils.py +++ b/computing/utils.py @@ -353,8 +353,15 @@ def get_agri_year_key(season_key): def calculate_precipitation_season( - geojson_filepath, draught_asset_id, start_year=2017, end_year=2024 + geojson_filepath, draught_asset_id, start_year=None, end_year=None ): + if start_year is None or end_year is None: + raise ValueError( + "start_year and end_year are required for calculate_precipitation_season." + ) + start_year = int(start_year) + end_year = int(end_year) + # Load the GeoJSON file with open(geojson_filepath, "r") as f: feature_collection = json.load(f) diff --git a/waterrejuvenation/tasks.py b/waterrejuvenation/tasks.py index 0c897eed..4f262062 100644 --- a/waterrejuvenation/tasks.py +++ b/waterrejuvenation/tasks.py @@ -22,7 +22,7 @@ from computing.water_rejuvenation.water_rejuventation import ( find_watersheds_for_point_with_buffer, ) -from computing.zoi_layers.zoi import generate_zoi +from computing.zoi_layers.zoi import generate_zoi, _resolve_zoi_time_window from projects.models import Project from utilities.constants import SITE_DATA_PATH, GEE_PATHS @@ -389,12 +389,16 @@ def normalize(val): return None return val - if is_processing_required and (not start_date or not end_date): + if (is_processing_required or is_closest_wp) and (not start_date or not end_date): raise ValueError( "start_date and end_date are required when processing desilting points " - "through ZOI (YYYY-MM-DD)." + "or finding the closest water pixel (YYYY-MM-DD)." ) + start_year = end_year = None + if start_date and end_date: + _, _, start_year, end_year = _resolve_zoi_time_window(start_date, end_date) + ee_initialize(gee_account_id) wb_obj = WaterbodiesFileUploadLog.objects.get(pk=file_obj_id) @@ -444,7 +448,11 @@ def normalize(val): print("inside closest wp") try: result_dict = find_nearest_water_pixel( - dsilting_obj_log.lat, dsilting_obj_log.lon, 1500 + dsilting_obj_log.lat, + dsilting_obj_log.lon, + 1500, + start_year=start_year, + end_year=end_year, ) print(result_dict) except Exception as e: @@ -528,11 +536,9 @@ def Generate_lulc_mws( start_date=None, end_date=None, ): - if not start_date or not end_date: - raise ValueError( - "start_date and end_date are required for water-rej ZOI generation " - "(YYYY-MM-DD)." - ) + start_date, end_date, start_year, end_year = _resolve_zoi_time_window( + start_date, end_date + ) proj_obj = Project.objects.get(pk=proj_id) asset_suffix = f"{proj_obj.name}_{proj_obj.id}".lower() asset_folder = [proj_obj.name.lower()] @@ -558,8 +564,8 @@ def Generate_lulc_mws( make_asset_public(mws_asset_id) if is_lulc_required: clip_lulc_v3( - start_year=2017, - end_year=2024, + start_year=start_year, + end_year=end_year, gee_account_id=gee_account_id, roi_path=mws_asset_id, asset_folder=asset_folder, @@ -570,7 +576,11 @@ def Generate_lulc_mws( except Exception as e: logger.error(f"Error in Generating Lulc and mws layer: {str(e)}") Generate_water_balance_indicator( - mws_asset_id, proj_id=proj_obj.id, gee_account_id=gee_account_id + mws_asset_id, + proj_id=proj_obj.id, + gee_account_id=gee_account_id, + start_date=start_date, + end_date=end_date, ) asset_suffix_swb3 = f"swb3_{proj_obj.name}+{proj_obj.id}" asset_id_swb = ( @@ -580,7 +590,10 @@ def Generate_lulc_mws( + asset_suffix_swb3 ) BuildMWSLayer( - gee_account_id=gee_account_id, proj_id=proj_obj.id, app_type="WATERBODY" + gee_account_id=gee_account_id, + proj_id=proj_obj.id, + app_type="WATERBODY", + export_year_range=(start_year, end_year), ) asset_suffix_wb = f"waterbodies_{asset_suffix}".lower() asset_id_wb = ( @@ -602,9 +615,14 @@ def Generate_lulc_mws( @shared_task() -def Generate_water_balance_indicator(mws_asset_id, proj_id, gee_account_id=None): +def Generate_water_balance_indicator( + mws_asset_id, proj_id, gee_account_id=None, start_date=None, end_date=None +): print(f"project id {gee_account_id}") + start_date, end_date, start_year, end_year = _resolve_zoi_time_window( + start_date, end_date + ) proj_obj = Project.objects.get(pk=proj_id) logger.info("Generating SWB layer for given lat long") asset_folder = [str(proj_obj.name).lower()] @@ -655,8 +673,8 @@ def Generate_water_balance_indicator(mws_asset_id, proj_id, gee_account_id=None) asset_suffix=asset_suffix, asset_folder_list=asset_folder, app_type="WATERBODY", - start_year="2017", - end_year="2024", + start_year=str(start_year), + end_year=str(end_year), is_all_classes=True, gee_account_id=gee_account_id, ) @@ -680,8 +698,8 @@ def Generate_water_balance_indicator(mws_asset_id, proj_id, gee_account_id=None) asset_suffix=asset_suffix, asset_folder_list=asset_folder, app_type="WATERBODY", - start_date="2017-06-30", - end_date="2025-07-1", + start_date=start_date, + end_date=end_date, is_annual=False, ) make_asset_public(asset_id_prec) @@ -691,12 +709,14 @@ def Generate_water_balance_indicator(mws_asset_id, proj_id, gee_account_id=None) asset_suffix=asset_suffix, asset_folder_list=asset_folder, app_type="WATERBODY", - start_year=2017, - end_year=2024, + start_year=start_year, + end_year=end_year, gee_account_id=gee_account_id, state=proj_obj.state_soi.state_name, ) - dst_filename = "drought_" + asset_suffix + "_" + str(2017) + "_" + str(2022) + dst_filename = ( + "drought_" + asset_suffix + "_" + str(start_year) + "_" + str(end_year) + ) draught_asset_id = ( get_gee_dir_path( asset_folder, asset_path=GEE_PATHS["WATERBODY"]["GEE_ASSET_PATH"] @@ -884,7 +904,7 @@ def BuildMWSLayer( block=None, district=None, drought_asset_override=None, # optional: full path to drought asset if you want to override default - export_year_range=(2017, 2022), # for naming drought asset + export_year_range=None, # (start_year, end_year) for naming drought asset ): """ Full BuildMWSLayer: builds final MWS waterbody FC, joins drought properties (flat, prefixed), @@ -902,6 +922,10 @@ def BuildMWSLayer( try: # initialize GEE + if export_year_range is None: + raise ValueError( + "export_year_range=(start_year, end_year) is required for BuildMWSLayer." + ) ee_initialize(gee_account_id) # ------------------------- @@ -953,10 +977,10 @@ def BuildMWSLayer( # ------------------------- # Drought asset id (default naming) # ------------------------- + start_y, end_y = export_year_range if drought_asset_override: draught_asset_id = drought_asset_override else: - start_y, end_y = export_year_range dst_filename = f"drought_{asset_suffix}" draught_asset_id = ( get_gee_dir_path( @@ -984,8 +1008,12 @@ def BuildMWSLayer( # ------------------------- # Create final_fc using your domain function # ------------------------- + start_y, end_y = export_year_range final_fc = calculate_precipitation_season( - mws_geojson_op, draught_asset_id=draught_asset_id + mws_geojson_op, + draught_asset_id=draught_asset_id, + start_year=start_y, + end_year=end_y, ) final_fc = ee.FeatureCollection(final_fc) diff --git a/waterrejuvenation/utils.py b/waterrejuvenation/utils.py index 96ba1737..667d318a 100644 --- a/waterrejuvenation/utils.py +++ b/waterrejuvenation/utils.py @@ -30,7 +30,7 @@ logger = logging.getLogger(__name__) from datetime import datetime -from nrm_app.settings import lulc_years, water_classes +from nrm_app.settings import water_classes import pandas as pd import math @@ -38,16 +38,6 @@ import requests import json -years = [ - "2017_2018", - "2018_2019", - "2019_2020", - "2020_2021", - "2021_2022", - "2022_2023", - "2023_2024", -] - # utils.py EXPECTED_EXCEL_HEADERS = { @@ -160,16 +150,16 @@ def fine_closest_wb_pixel(lon, lat): ee_initialize() -def find_nearest_water_pixel(lat, lon, distance_threshold): +def find_nearest_water_pixel(lat, lon, distance_threshold, start_year=None, end_year=None): """ Finds the nearest water pixel within a given threshold. Parameters: lat (float): Latitude of the input location. lon (float): Longitude of the input location. - lulc_years (list of str): List of LULC year identifiers (e.g., '2017_2018'). - water_classes (list of int): List of LULC class codes considered as water. distance_threshold (float): Maximum distance (in meters) to consider. + start_year (int): First hydrological year start (required). + end_year (int): Last hydrological year start (required). Returns: dict: { @@ -179,6 +169,16 @@ def find_nearest_water_pixel(lat, lon, distance_threshold): 'distance_m': float (if success) } """ + if start_year is None or end_year is None: + raise ValueError( + "start_year and end_year are required for find_nearest_water_pixel." + ) + start_year = int(start_year) + end_year = int(end_year) + if start_year > end_year: + raise ValueError("start_year must be less than or equal to end_year.") + + lulc_year_ids = [f"{y}_{y + 1}" for y in range(start_year, end_year + 1)] print("Given lat long") print(f"-----{lat}----{lon}---") @@ -186,7 +186,7 @@ def find_nearest_water_pixel(lat, lon, distance_threshold): # Create water masks from each LULC year water_masks = [] - for year in lulc_years: + for year in lulc_year_ids: image = ee.Image(f"{PAN_INDIA_LULC_V3_DATASET}{year}") water_mask = ( image.select("predicted_label") @@ -670,8 +670,18 @@ def merge_features(feature): def get_ndvi_for_zoi( - zoi_asset_path, asset_suffix, asset_folder, proj_id=None, app_type="WATER_REJ" + zoi_asset_path, + asset_suffix, + asset_folder, + proj_id=None, + app_type="WATER_REJ", + start_year=None, + end_year=None, ): + if start_year is None or end_year is None: + raise ValueError( + "start_year and end_year are required for get_ndvi_for_zoi." + ) proj_obj = Project.objects.get(pk=proj_id) asset_suffix_ndvi = f"zoi_ndvi_{proj_obj.name}_{proj_obj.id}" ndvi_asset_path = ( @@ -685,7 +695,9 @@ def get_ndvi_for_zoi( asset_folder, app_type, f"{proj_obj.name}_{proj_obj.id}".lower(), zoi_roi=zoi_asset_path ) - fc = get_ndvi_data(zoi_collections, 2017, 2024, asset_suffix_ndvi, ndvi_asset_path) + fc = get_ndvi_data( + zoi_collections, start_year, end_year, asset_suffix_ndvi, ndvi_asset_path + ) task = ee.batch.Export.table.toAsset( collection=fc, description=asset_suffix_ndvi, assetId=ndvi_asset_path )