Spaces:
Paused
Paused
| """Validated DAG executor using the existing specialist/GIS architecture.""" | |
| import asyncio | |
| import json | |
| import time | |
| from datetime import datetime, UTC | |
| from pathlib import Path | |
| import numpy as np | |
| from scipy import ndimage | |
| from satquery_engine.schemas import AnalysisResponse, ArtifactRef, EvidenceItem, GeoVerdict, TraceStep, VerdictStatus, TaskType, TaskPlan, PlanNode | |
| from satquery_engine.services.intents import plan_query, validate_dag | |
| from satquery_engine.services.input_configuration import classify_configuration, policy_blockers, node_policy_error, acquisition_time | |
| from satquery_engine.services.raster import validate_inputs, render_preview | |
| from satquery_engine.services.alignment import align_pair | |
| from satquery_engine.services.buildings import detect_buildings, match_buildings | |
| from satquery_engine.services.measurements import measure_cover, measure_cover_change | |
| from satquery_engine.services.model_registry import registry | |
| from satquery_engine.services.confidence import confidence_breakdown | |
| from satquery_engine.services.report import write_pdf_report, write_manifest | |
| from satquery_engine.services.geolocation import parse_user_coordinates | |
| async def analyze(*,result_id,query,pair_type,image_paths,output_dir,progress=None,input_filenames=None,input_asset_ids=None): | |
| started=time.perf_counter(); output_dir.mkdir(parents=True,exist_ok=True) | |
| trace=[]; artifacts=[]; evidence=[]; findings=[]; statistics={}; models=[]; limitations=[]; results={} | |
| def emit(message): | |
| if progress: progress(message) | |
| def add_path(path): | |
| rel=path.relative_to(output_dir).as_posix() | |
| mime={".png":"image/png",".tif":"image/tiff",".geojson":"application/geo+json",".pdf":"application/pdf",".json":"application/json"}.get(path.suffix,"application/octet-stream") | |
| if not any(a.url.endswith('/'+rel) for a in artifacts): | |
| name = path.name if not any(a.name == path.name for a in artifacts) else rel.replace('/', '_') | |
| import uuid | |
| artifacts.append(ArtifactRef(artifact_id=str(uuid.uuid5(uuid.NAMESPACE_URL,f"{result_id}/{rel}")),name=name,url=f"/artifacts/{result_id}/{rel}",mime_type=mime)) | |
| emit("Reading and checking image metadata") | |
| metadata,quality=await asyncio.to_thread(validate_inputs,image_paths) | |
| if input_filenames: | |
| metadata = [asset.model_copy(update={"filename": name}) for asset, name in zip(metadata, input_filenames)] | |
| if input_asset_ids: | |
| metadata=[asset.model_copy(update={"asset_id":ident}) for asset,ident in zip(metadata,input_asset_ids,strict=True)] | |
| configuration,clarification=classify_configuration(metadata,pair_type) | |
| emit("Understanding your request and planning analysis") | |
| user_coordinates=parse_user_coordinates(query) | |
| plan=(TaskPlan(task=TaskType.GROUNDING,application="user_location",specific_task="user_location", | |
| intents=["USER_LOCATION"],tools=["validate","synthesize","report"], | |
| reason="Validated explicit user coordinates without model inference.", | |
| nodes=[PlanNode(node_id="validate",tool="validate"), | |
| PlanNode(node_id="synthesize",tool="synthesize",depends_on=["validate"]), | |
| PlanNode(node_id="report",tool="report",depends_on=["synthesize"])]) | |
| if user_coordinates is not None else plan_query(query,len(image_paths),pair_type,configuration)) | |
| validate_dag(plan) | |
| clarification=plan.clarification or (clarification if len(image_paths)==2 else None) | |
| if plan.requires_temporal_relationship and configuration.startswith("PAIR_BITEMPORAL"): | |
| dates=[a.acquisition_date for a in metadata] | |
| if all(dates): | |
| try: | |
| order=sorted(range(2),key=lambda i: acquisition_time(dates[i])) | |
| image_paths=[image_paths[i] for i in order]; metadata=[metadata[i] for i in order] | |
| except ValueError: pass | |
| quality.blockers.extend(policy_blockers(plan,metadata,configuration)) | |
| quality.compatible=not quality.blockers and not clarification | |
| trace.append(TraceStep(step=1,component="validation_and_planner",action="Checked metadata, pair semantics and allow-listed task DAG", | |
| status="ok" if quality.compatible else "blocked",duration_ms=round((time.perf_counter()-started)*1000), | |
| details={"configuration":configuration,"intents":plan.intents,"blockers":quality.blockers})) | |
| for i,path in enumerate(image_paths): | |
| preview=output_dir/f"preview_{i+1}.png" | |
| await asyncio.to_thread(render_preview,path,preview) | |
| add_path(preview) | |
| if user_coordinates is not None and quality.compatible: | |
| latitude,longitude=user_coordinates | |
| location={"type":"user_location","latitude":latitude,"longitude":longitude,"source":"USER_COORDINATES"} | |
| statistics["user_location"]=location | |
| summary=f"Marked the user-provided location at latitude {latitude}, longitude {longitude}." | |
| point_file=output_dir/"user_location.geojson" | |
| point_file.write_text(json.dumps({"type":"FeatureCollection", | |
| "properties":{"crs":"EPSG:4326","coordinate_space":"geographic","source":"USER_COORDINATES"}, | |
| "features":[{"type":"Feature","geometry":{"type":"Point","coordinates":[longitude,latitude]}, | |
| "properties":location}]}),encoding="utf-8") | |
| add_path(point_file) | |
| evidence.append(EvidenceItem(evidence_id="evidence_1",kind="user_location",producer="user", | |
| summary=summary,confidence=1.0,confidence_source="authoritative_user_input", | |
| artifact_url=artifacts[-1].url,metrics=location)) | |
| findings.append({"text":summary,"evidence_ids":["evidence_1"],"status":"supported"}) | |
| trace.append(TraceStep(step=len(trace)+1,component="user_location",action="Validated and mapped explicit coordinates", | |
| status="ok",duration_ms=0,details={"source":"USER_COORDINATES"})) | |
| working=list(image_paths); completed={"validate"}; failed=set() | |
| from satquery_engine.services.surface_context import SurfaceContext | |
| surfaces = SurfaceContext(output_dir / "surface_context") | |
| preflight={n.node_id:node_policy_error(n,metadata) for n in plan.nodes} | |
| if quality.compatible and plan.task not in {TaskType.UNSUPPORTED,TaskType.UNCLEAR}: | |
| for node in plan.nodes: | |
| if node.tool in {"validate","synthesize","report"}: continue | |
| if any(d in failed for d in node.depends_on): | |
| failed.add(node.node_id) | |
| trace.append(TraceStep(step=len(trace)+1,component=node.tool,action="Skipped because required evidence is unavailable",status="skipped",duration_ms=0)) | |
| continue | |
| tick=time.perf_counter(); emit({"register":"Aligning images","buildings":"Running building detection","building_match":"Comparing building footprints", "spectral":"Measuring image features","spectral_change":"Measuring change","land_cover":"Mapping land cover from image evidence","change":"Measuring image differences","fusion":"Checking optical and radar evidence","vlm":"Describing image evidence"}[node.tool]) | |
| sub=output_dir/node.node_id; sub.mkdir(exist_ok=True) | |
| params=node.parameters | |
| try: | |
| if preflight.get(node.node_id): raise ValueError(preflight[node.node_id]) | |
| result=None; summary="" | |
| if node.tool=="register": | |
| a,b,registration=await asyncio.to_thread(align_pair,*working,output_dir,not plan.requires_temporal_relationship) | |
| working=[a,b]; add_path(output_dir/"registration_report.json") | |
| quality.checks["alignment_score"]=registration.alignment_score | |
| quality.checks["registration_method"]=registration.method | |
| if plan.requires_temporal_relationship and registration.alignment_score<.35: | |
| raise ValueError("The images could not be aligned reliably enough for change analysis.") | |
| limitations.append(registration.message) | |
| elif node.tool=="buildings": | |
| emit("Checking water boundaries before counting buildings") | |
| water_context = await asyncio.to_thread(surfaces.water, working[params["asset"]]) | |
| for path in water_context.get("paths", []): add_path(path) | |
| models.extend(water_context.get("models_used", [])) | |
| result=await asyncio.to_thread(detect_buildings,working[params["asset"]],sub,emit,water_result=water_context) | |
| summary=f'I found approximately {result["count"]} visible building footprints in image {params["asset"]+1}.' | |
| if result.get("count_reliability") == "COUNT_UNRELIABLE": | |
| summary += ' This is a low-confidence detection count; an accurate building inventory cannot be established from this result.' | |
| if result["count"] == 0: | |
| summary += ' No footprints were detected; this does not establish that no buildings are present.' | |
| if any(i in plan.intents for i in ["BUILT_UP_ANALYSIS", "BUILT_UP_CHANGE", "BUILDING_FOOTPRINT"]): | |
| summary += f' Their measured footprint area is {result["area_m2"]:.1f} square metres.' if result["area_m2"] is not None else f' They cover {result["selected_pixels"]} image pixels; geographic area is unknown.' | |
| models.append(result["model_id"]) | |
| elif node.tool=="building_match": | |
| result=await asyncio.to_thread(match_buildings,results["buildings_a"],results["buildings_b"],sub) | |
| summary=f'The detected building count changed from {result["before_count"]} to {result["after_count"]}. There are {result["possible_new_count"]} possible new and {result["possible_removed_count"]} possibly removed footprints.' | |
| if result['net_footprint_area_m2'] is not None: | |
| summary += f' The net detected footprint area change is {result["net_footprint_area_m2"]:+.1f} square metres; this measures buildings, not all built-up surfaces.' | |
| elif node.tool=="spectral": | |
| if params["target"] == "water": | |
| result=await asyncio.to_thread(surfaces.water,working[0],largest=params.get("largest",False),strict=params.get("strict",False)) | |
| else: | |
| result=await asyncio.to_thread(measure_cover,working[0],sub,**params) | |
| unit=f'{result["area_m2"]:.1f} square metres' if result["area_m2"] is not None else f'{result["selected_pixels"]} sampled image pixels' | |
| loc=f' in the {result["location_description"]}' if result.get("location_description") else '' | |
| if params["target"] == "water" and result.get("evidence_state") == "INSUFFICIENT_EVIDENCE": | |
| summary="Water boundaries could not be verified from this RGB image because the compatible aerial models failed." | |
| else: | |
| summary=f'{result["method"]} identifies {params["target"]}-like regions{loc} covering {result["coverage_percent"]:.2f}% of valid pixels ({unit}).' | |
| elif node.tool=="land_cover": | |
| from satquery_engine.services.land_cover import classify_land_cover_composite | |
| result=await asyncio.to_thread(classify_land_cover_composite,working[0],sub,query, | |
| water_result=await asyncio.to_thread(surfaces.water,working[0]), building_result=results.get("buildings_a"), | |
| vegetation_result=results.get("vegetation_measure")) | |
| summary=result["summary"] | |
| elif node.tool=="spectral_change": | |
| result=await asyncio.to_thread(measure_cover_change,*working,sub,params["target"]) | |
| summary=f'The area meeting the {params["target"]} index threshold changed by {result["net_percentage_points"]:+.2f} percentage points across shared valid pixels.' | |
| elif node.tool=="change": | |
| from satquery_engine.services.raster import deterministic_change_detection | |
| result=await asyncio.to_thread(deterministic_change_detection,*working,sub) | |
| result.pop("mask_path"); result.pop("geojson_path") | |
| result["method"]="Radiometric difference threshold; not a learned change model" | |
| result["evidence_strength"]=0.0 | |
| result["limitations"]=["This measures image differences, which can reflect illumination or season. The learned change specialist is unavailable; the type of change is unverified."] | |
| summary=f'I measured image differences over {result["changed_percent"]:.2f}% of the sampled pixels. I cannot establish what caused them.' | |
| elif node.tool=="fusion": | |
| external=await registry.invoke("croma",{"optical_path":str(working[0]),"sar_path":str(working[1])}) | |
| if not external.get("available"): raise ValueError("The optical-SAR fusion specialist is unavailable. No fused result was generated.") | |
| # A representation alone cannot ground semantic spatial claims. | |
| raise ValueError("The fusion service returned a representation, but a validated task-specific decoder is not configured. No geographic fusion claim was generated.") | |
| elif node.tool=="vlm": | |
| if params.get("counting"): | |
| raise ValueError("A dedicated detector for this object type is unavailable. A language model will not supply an exact count.") | |
| facts = {name:{k:v for k,v in value.items() if k in {"count","coverage_percent","area_m2","breakdown","method","limitations"}} for name,value in results.items()} | |
| explanation_query = query + "\nExplain only the following backend evidence. Do not invent or revise counts, areas, coordinates, confidence or percentages. Do not introduce numerical claims in the scene description. Backend facts: " + json.dumps(facts,allow_nan=False) | |
| external=await registry.invoke_vlm(query=explanation_query,image_paths=[str(p) for p in working],tile_paths=[],metadata=[a.model_dump(mode="json") for a in metadata]) | |
| if not external.get("available"): raise ValueError("The image description model is unavailable. No description was generated.") | |
| if params.get("grounding"): | |
| raise ValueError("A validated spatial detector for this object is unavailable; text alone cannot establish its location.") | |
| summary=external.get("answer","") | |
| if not summary: raise ValueError("The image description model returned no answer.") | |
| # Numerical measurements never come from language output. | |
| import re | |
| if re.search(r"\d",summary): raise ValueError("The description included unsupported numerical claims and was withheld.") | |
| models.append(external["model"]) | |
| result={"method":external["model"],"evidence_strength":0.0,"paths":[], | |
| "limitations":["This description is a model interpretation; no exact geographic measurements were made."]} | |
| if result is not None: | |
| from satquery_engine.services.quality_gate import validate_result_evidence | |
| validate_result_evidence(result,output_dir / "surface_context" if node.tool == "spectral" and params["target"] == "water" else sub) | |
| results[node.node_id]=result | |
| models.extend(result.get("models_used", [])) | |
| for path in result.get("paths",[]): add_path(path) | |
| metrics={k:v for k,v in result.items() if k not in {"paths","features","limitations","mask_path","geojson_path"} and not isinstance(v,np.ndarray)} | |
| statistics[node.node_id]=metrics | |
| eid=f"evidence_{len(evidence)+1}" | |
| evidence.append(EvidenceItem(evidence_id=eid,kind=node.tool,producer=result.get("model_id",result.get("method",node.tool)),summary=summary, | |
| confidence=float(result.get("evidence_strength",0)),confidence_source="uncalibrated_model_score" if node.tool=="buildings" else "documented_measurement_strength", | |
| artifact_url=next((a.url for a in artifacts if any(a.url.endswith('/'+p.relative_to(output_dir).as_posix()) for p in result.get("paths",[]) if 'overlay' in p.name)),None),metrics=metrics, | |
| timestamp=datetime.now(UTC).isoformat(),asset_indices=[params["asset"]] if "asset" in params else list(range(len(working))))) | |
| findings.append({"text":summary,"evidence_ids":[eid],"status":"supported_with_limitations"}) | |
| limitations.extend(result.get("limitations",[])) | |
| completed.add(node.node_id) | |
| trace.append(TraceStep(step=len(trace)+1,component=node.tool,action=node.node_id,status="ok",duration_ms=round((time.perf_counter()-tick)*1000),details={"preprocessing":result.get("preprocessing") if result else None,"checkpoint_sha256":result.get("checkpoint_sha256") if result else None})) | |
| except (ValueError,ImportError,RuntimeError,OSError) as exc: | |
| failed.add(node.node_id); limitations.append(str(exc)) | |
| findings.append({"text":str(exc),"evidence_ids":[],"status":"insufficient_evidence"}) | |
| trace.append(TraceStep(step=len(trace)+1,component=node.tool,action=node.node_id,status="unavailable",duration_ms=round((time.perf_counter()-tick)*1000),details={"message":str(exc)})) | |
| emit("Collecting evidence and preparing your answer") | |
| limitations=list(dict.fromkeys(limitations+quality.warnings)) | |
| if clarification: | |
| answer=clarification; status=VerdictStatus.INSUFFICIENT_EVIDENCE | |
| elif quality.blockers: | |
| answer=" ".join(quality.blockers); status=VerdictStatus.INVALID_INPUT | |
| elif plan.task==TaskType.UNSUPPORTED: | |
| answer="This request is outside the supported satellite-image analysis tasks."; status=VerdictStatus.UNSUPPORTED_TASK | |
| else: | |
| answer=" ".join(f["text"] for f in findings) or "No evidence was produced for this request." | |
| status=(VerdictStatus.SUPPORTED if user_coordinates is not None and evidence else | |
| VerdictStatus.SUPPORTED_WITH_LIMITATIONS if evidence else VerdictStatus.INSUFFICIENT_EVIDENCE) | |
| if any(r.get("count_reliability") == "COUNT_UNRELIABLE" or | |
| (r.get("benchmark_f1") is not None and r["benchmark_f1"]<.6) for r in results.values()): | |
| status=VerdictStatus.LOW_CONFIDENCE | |
| if any((r.get("aerial_water_agreement") or {}).get("withheld_fraction", 0) > .5 | |
| for r in results.values()): | |
| status=VerdictStatus.LOW_CONFIDENCE | |
| if any(r.get("evidence_state") == "INSUFFICIENT_EVIDENCE" for r in results.values()): | |
| status=(VerdictStatus.INSUFFICIENT_EVIDENCE if plan.task == TaskType.WATER_ANALYSIS | |
| else VerdictStatus.LOW_CONFIDENCE) | |
| breakdown=confidence_breakdown(quality,evidence) | |
| verdict=GeoVerdict(status=status,answer=answer,confidence=breakdown.final_score, | |
| confidence_kind="authoritative_user_input" if user_coordinates is not None and evidence else "uncalibrated_evidence_strength", | |
| limitations=limitations,confidence_breakdown=breakdown) | |
| # List canonical downloadable outputs before serialization, so JSON/PDF/API agree. | |
| add_path(output_dir/"GeoProof_Report.pdf"); add_path(output_dir/"analysis.json") | |
| response=AnalysisResponse(result_id=result_id,query=query,generated_at=datetime.now(UTC).isoformat(),task_plan=plan, | |
| mode="automatic_evidence_analysis",assets=metadata,quality=quality,evidence=evidence,verdict=verdict,trace=trace, | |
| artifacts=artifacts,input_configuration=configuration,findings=findings,statistics=statistics,models_used=list(dict.fromkeys(models)), | |
| clarification=clarification,report_id=result_id,timings={"analysis_ms":round((time.perf_counter()-started)*1000)}) | |
| emit("Creating your report") | |
| await asyncio.to_thread(write_pdf_report,output_dir,response.model_dump(mode="json")) | |
| await asyncio.to_thread(write_manifest,output_dir,response.model_dump(mode="json")) | |
| return response | |