Coverage for backend/django/core/auxiliary/services/result_summary/cache.py: 84%

156 statements  

« prev     ^ index     » next       coverage.py v7.10.7, created at 2026-07-22 05:22 +0000

1import logging 

2from enum import StrEnum 

3 

4from common.models.notifications.payloads import ( 

5 NotificationServiceMessageType, 

6 ResultSummaryAvailablePayload, 

7 ResultSummaryCalculatingPayload, 

8 ResultSummaryCalculationError, 

9 ResultSummaryCalculationStatus, 

10 ResultSummaryFailedPayload, 

11 ResultSummaryStatusPayload, 

12) 

13from common.services import messaging 

14from core.auxiliary.models.Scenario import Scenario, ScenarioTabTypeEnum 

15from core.auxiliary.models.ScenarioResultSummaryCache import ( 

16 ScenarioResultSummaryCache, 

17 ScenarioResultSummaryStatus, 

18) 

19from core.auxiliary.models.Solution import Solution 

20from core.auxiliary.models.Task import Task, TaskType 

21from core.auxiliary.services.result_summary.cache_payloads import ( 

22 ResultSummaryCachePayload, 

23 cached_summary_to_result_summary, 

24) 

25from core.auxiliary.services.result_summary.contracts import ( 

26 ResultSummary, 

27 ResultSummaryRequest, 

28) 

29from core.auxiliary.services.result_summary.summaries import ( 

30 build_cached_result_summaries_with_errors, 

31 get_result_summaries, 

32) 

33from django.db import transaction 

34from django.db.models import Count, Max 

35from django.utils import timezone 

36from pydantic import ValidationError 

37 

38logger = logging.getLogger(__name__) 

39 

40 

41class ResultSummaryCacheInvalidationReason(StrEnum): 

42 mss_solve_started = "mss_solve_started" 

43 mss_input_replaced = "mss_input_replaced" 

44 mss_input_cleared = "mss_input_cleared" 

45 mss_input_cell_changed = "mss_input_cell_changed" 

46 mss_input_column_changed = "mss_input_column_changed" 

47 mss_time_series_settings_changed = "mss_time_series_settings_changed" 

48 result_summary_cache_schema_changed = "result_summary_cache_schema_changed" 

49 

50 

51class ResultSummaryCacheUnavailable(Exception): 

52 """Raised when MSS aggregate summaries are not ready for a 200 response.""" 

53 

54 def __init__( 

55 self, 

56 *, 

57 payload: ResultSummaryStatusPayload, 

58 response_status: int, 

59 ) -> None: 

60 super().__init__(payload.status.value) 

61 self.payload = payload 

62 self.response_status = response_status 

63 

64 

65def invalidate_result_summary_cache_for_scenario( 

66 *, scenario_id: int, reason: ResultSummaryCacheInvalidationReason 

67) -> None: 

68 """Delete cached aggregate summaries when MSS inputs/results change.""" 

69 

70 deleted_count, _ = ScenarioResultSummaryCache.objects.filter( 

71 scenario_id=scenario_id 

72 ).delete() 

73 logger.info( 

74 "Invalidated %s result summary cache rows for scenario_id=%s reason=%s.", 

75 deleted_count, 

76 scenario_id, 

77 reason, 

78 ) 

79 

80 

81def build_result_summary_cache_for_successful_mss_task( 

82 *, parent_task_id: int, scenario_id: int 

83) -> None: 

84 """Build MSS result-summary cache rows after all child solves succeeded.""" 

85 

86 try: 

87 task = Task.objects.select_related("metadata").get(id=parent_task_id) 

88 scenario = Scenario.objects.get(id=scenario_id) 

89 except (Task.DoesNotExist, Scenario.DoesNotExist): 

90 logger.info( 

91 "Skipping result summary cache build for missing parent_task_id=%s scenario_id=%s.", 

92 parent_task_id, 

93 scenario_id, 

94 ) 

95 return 

96 

97 if not _should_build_cache(task=task, scenario=scenario): 

98 return 

99 

100 simulation_object_ids = _solved_simulation_object_ids(scenario_id=scenario.id) 

101 if not simulation_object_ids: 101 ↛ 102line 101 didn't jump to line 102 because the condition on line 101 was never true

102 return 

103 

104 calculation_started_at = timezone.now() 

105 with transaction.atomic(): 

106 for simulation_object_id in simulation_object_ids: 

107 cache, created = ScenarioResultSummaryCache.objects.select_for_update().get_or_create( 

108 scenario_id=scenario.id, 

109 simulationObject_id=simulation_object_id, 

110 defaults={ 

111 "flowsheet_state": scenario.flowsheet_state, 

112 "task_id": task.id, 

113 }, 

114 ) 

115 if not created and _cache_belongs_to_newer_task(cache, task): 115 ↛ 116line 115 didn't jump to line 116 because the condition on line 115 was never true

116 continue 

117 cache.flowsheet_state = scenario.flowsheet_state 

118 cache.task_id = task.id 

119 cache.status = ScenarioResultSummaryStatus.CALCULATING 

120 cache.payload = {} 

121 cache.errors = [] 

122 cache.calculation_started_at = calculation_started_at 

123 cache.calculation_completed_at = None 

124 cache.calculation_failed_at = None 

125 cache.save( 

126 update_fields=[ 

127 "flowsheet_state", 

128 "task", 

129 "status", 

130 "payload", 

131 "errors", 

132 "calculation_started_at", 

133 "calculation_completed_at", 

134 "calculation_failed_at", 

135 ] 

136 ) 

137 

138 _send_calculating_notification( 

139 flowsheet_id=scenario.flowsheet_state.flowsheet_id, 

140 scenario_id=scenario.id, 

141 simulation_object_ids=simulation_object_ids, 

142 source_task_id=task.id, 

143 ) 

144 

145 for simulation_object_id in simulation_object_ids: 

146 _build_cache_for_simulation_object( 

147 task=task, 

148 scenario=scenario, 

149 simulation_object_id=simulation_object_id, 

150 ) 

151 

152 

153def get_cached_result_summaries( 

154 request: ResultSummaryRequest, 

155) -> list[ResultSummary] | None: 

156 """Return cached summaries for ready MSS cache rows, if present.""" 

157 

158 cache = _get_cache_row(request) 

159 if cache is None: 159 ↛ 160line 159 didn't jump to line 160 because the condition on line 159 was never true

160 return None 

161 if cache.status != ScenarioResultSummaryStatus.AVAILABLE: 

162 return None 

163 

164 try: 

165 payload = ResultSummaryCachePayload.model_validate(cache.payload) 

166 except ValidationError: 

167 invalidate_result_summary_cache_for_scenario( 

168 scenario_id=request.scenario_id, 

169 reason=ResultSummaryCacheInvalidationReason.result_summary_cache_schema_changed, 

170 ) 

171 return None 

172 return [ 

173 cached_summary_to_result_summary( 

174 summary, 

175 request.histogram_options_for(summary.key), 

176 ) 

177 for summary in payload.summaries 

178 ] 

179 

180 

181def get_result_summaries_with_cache( 

182 request: ResultSummaryRequest, 

183) -> list[ResultSummary]: 

184 """Use cached MSS aggregate summaries, preserving live calculation elsewhere.""" 

185 

186 scenario = Scenario.objects.get(id=request.scenario_id) 

187 if scenario.state_name != ScenarioTabTypeEnum.MultiSteadyState: 

188 return get_result_summaries(request) 

189 

190 cache = _get_cache_row(request) 

191 if cache is None: 

192 if _has_source_solutions(request): 

193 # Older MSS results may predate the cache row. Keep those visible 

194 # rather than returning a permanent pending state for historical solves. 

195 return get_result_summaries(request) 

196 # Some selected objects legitimately produce no solution rows. Treat 

197 # that as empty result data instead of an endless post-processing state. 

198 return [] 

199 

200 cached = get_cached_result_summaries(request) 

201 if cached is not None: 

202 return cached 

203 if ( 

204 cache.status == ScenarioResultSummaryStatus.AVAILABLE 

205 and _get_cache_row(request) is None 

206 and _has_source_solutions(request) 

207 ): 

208 return get_result_summaries(request) 

209 

210 status_payload = get_result_summary_status(request) 

211 response_status = ( 

212 409 

213 if status_payload.status 

214 in { 

215 ResultSummaryCalculationStatus.failed, 

216 ResultSummaryCalculationStatus.partial, 

217 } 

218 else 202 

219 ) 

220 raise ResultSummaryCacheUnavailable( 

221 payload=status_payload, 

222 response_status=response_status, 

223 ) 

224 

225 

226def get_result_summary_status( 

227 request: ResultSummaryRequest, 

228) -> ResultSummaryStatusPayload: 

229 """Return status metadata for one MSS scenario/object summary cache row.""" 

230 

231 cache = _get_cache_row(request) 

232 if cache is None: 232 ↛ 233line 232 didn't jump to line 233 because the condition on line 232 was never true

233 return ResultSummaryStatusPayload( 

234 scenario_id=request.scenario_id, 

235 simulation_object_id=request.simulation_object_id, 

236 source_task_id=None, 

237 status=ResultSummaryCalculationStatus.pending, 

238 errors=[], 

239 ) 

240 

241 return ResultSummaryStatusPayload( 

242 scenario_id=request.scenario_id, 

243 simulation_object_id=request.simulation_object_id, 

244 source_task_id=cache.task_id, 

245 status=ResultSummaryCalculationStatus(cache.status), 

246 errors=[ 

247 ResultSummaryCalculationError.model_validate(error) 

248 for error in cache.errors 

249 ], 

250 ) 

251 

252 

253def _should_build_cache(*, task: Task, scenario: Scenario) -> bool: 

254 metadata = task.metadata 

255 return ( 

256 task.task_type == TaskType.IDAES_SOLVE 

257 and scenario.state_name == ScenarioTabTypeEnum.MultiSteadyState 

258 and metadata is not None 

259 and metadata.scheduled_tasks > 0 

260 and metadata.successful_tasks == metadata.scheduled_tasks 

261 and metadata.failed_tasks == 0 

262 and metadata.cancelled_tasks == 0 

263 ) 

264 

265 

266def _solved_simulation_object_ids(*, scenario_id: int) -> list[int]: 

267 return list( 

268 Solution.objects.filter(scenario_id=scenario_id) 

269 .exclude(property__property__set__simulationObject_id__isnull=True) 

270 .order_by() 

271 .values_list("property__property__set__simulationObject_id", flat=True) 

272 .distinct() 

273 ) 

274 

275 

276def _cache_belongs_to_newer_task( 

277 cache: ScenarioResultSummaryCache, task: Task 

278) -> bool: 

279 # Result-summary builds can overlap when a user reruns a scenario quickly. 

280 # A late callback from an older task must not roll the cache back to stale 

281 # post-processing state. 

282 return cache.task_id is not None and cache.task_id > task.id 

283 

284 

285def _build_cache_for_simulation_object( 

286 *, task: Task, scenario: Scenario, simulation_object_id: int 

287) -> None: 

288 try: 

289 build = build_cached_result_summaries_with_errors( 

290 scenario_id=scenario.id, 

291 simulation_object_id=simulation_object_id, 

292 source_task_id=task.id, 

293 flowsheet_id=scenario.flowsheet_state.flowsheet_id, 

294 ) 

295 status = ( 

296 ScenarioResultSummaryStatus.PARTIAL 

297 if build.errors 

298 else ScenarioResultSummaryStatus.AVAILABLE 

299 ) 

300 if not build.payload.summaries: 300 ↛ 301line 300 didn't jump to line 301 because the condition on line 300 was never true

301 status = ScenarioResultSummaryStatus.FAILED 

302 errors = build.errors or [ 

303 ResultSummaryCalculationError( 

304 message="No aggregate result summaries could be calculated.", 

305 cause="no_usable_summaries", 

306 simulation_object_id=simulation_object_id, 

307 ) 

308 ] 

309 else: 

310 errors = build.errors 

311 except Exception as error: 

312 build = None 

313 status = ScenarioResultSummaryStatus.FAILED 

314 errors = [ 

315 ResultSummaryCalculationError( 

316 message="Aggregate result summaries could not be calculated.", 

317 cause=error.__class__.__name__, 

318 simulation_object_id=simulation_object_id, 

319 ) 

320 ] 

321 

322 diagnostics = _source_solution_diagnostics( 

323 scenario_id=scenario.id, 

324 simulation_object_id=simulation_object_id, 

325 ) 

326 completed_at = timezone.now() 

327 payload = ( 

328 build.payload.model_dump(mode="json", by_alias=True) 

329 if build is not None 

330 else {} 

331 ) 

332 error_payload = [error.model_dump(mode="json") for error in errors] 

333 

334 with transaction.atomic(): 

335 try: 

336 cache = ScenarioResultSummaryCache.objects.select_for_update().get( 

337 scenario_id=scenario.id, 

338 simulationObject_id=simulation_object_id, 

339 ) 

340 except ScenarioResultSummaryCache.DoesNotExist: 

341 return 

342 

343 if cache.task_id != task.id: 343 ↛ 344line 343 didn't jump to line 344 because the condition on line 343 was never true

344 return 

345 

346 cache.status = status 

347 cache.payload = payload 

348 cache.errors = error_payload 

349 cache.solution_count = diagnostics["solution_count"] or 0 

350 cache.row_count = diagnostics["row_count"] or 0 

351 cache.source_solution_created_at_max = diagnostics[ 

352 "source_solution_created_at_max" 

353 ] 

354 if status == ScenarioResultSummaryStatus.FAILED: 354 ↛ 355line 354 didn't jump to line 355 because the condition on line 354 was never true

355 cache.calculation_failed_at = completed_at 

356 cache.calculation_completed_at = None 

357 else: 

358 cache.calculation_completed_at = completed_at 

359 cache.calculation_failed_at = None 

360 cache.save() 

361 

362 if status == ScenarioResultSummaryStatus.FAILED: 362 ↛ 363line 362 didn't jump to line 363 because the condition on line 362 was never true

363 _send_failed_notification( 

364 flowsheet_id=scenario.flowsheet_state.flowsheet_id, 

365 scenario_id=scenario.id, 

366 simulation_object_id=simulation_object_id, 

367 source_task_id=task.id, 

368 errors=errors, 

369 ) 

370 else: 

371 _send_available_notification( 

372 flowsheet_id=scenario.flowsheet_state.flowsheet_id, 

373 scenario_id=scenario.id, 

374 simulation_object_id=simulation_object_id, 

375 source_task_id=task.id, 

376 errors=errors, 

377 ) 

378 

379 

380def _source_solution_diagnostics( 

381 *, scenario_id: int, simulation_object_id: int 

382) -> dict[str, int | object]: 

383 return Solution.objects.filter( 

384 scenario_id=scenario_id, 

385 property__property__set__simulationObject_id=simulation_object_id, 

386 ).aggregate( 

387 solution_count=Count("id"), 

388 row_count=Count("solve_index", distinct=True), 

389 source_solution_created_at_max=Max("created_at"), 

390 ) 

391 

392 

393def _has_source_solutions(request: ResultSummaryRequest) -> bool: 

394 return Solution.objects.filter( 

395 scenario_id=request.scenario_id, 

396 property__property__set__simulationObject_id=request.simulation_object_id, 

397 ).exists() 

398 

399 

400def _get_cache_row( 

401 request: ResultSummaryRequest, 

402) -> ScenarioResultSummaryCache | None: 

403 return ( 

404 ScenarioResultSummaryCache.objects.filter( 

405 scenario_id=request.scenario_id, 

406 simulationObject_id=request.simulation_object_id, 

407 ) 

408 .select_related("task") 

409 .first() 

410 ) 

411 

412 

413def _send_calculating_notification( 

414 *, 

415 flowsheet_id: int, 

416 scenario_id: int, 

417 simulation_object_ids: list[int], 

418 source_task_id: int, 

419) -> None: 

420 payload = ResultSummaryCalculatingPayload( 

421 scenario_id=scenario_id, 

422 simulation_object_ids=simulation_object_ids, 

423 source_task_id=source_task_id, 

424 ) 

425 messaging.send_flowsheet_notification_message( 

426 flowsheet_id, 

427 payload.model_dump(mode="json"), 

428 NotificationServiceMessageType.RESULT_SUMMARY_CALCULATING, 

429 ) 

430 

431 

432def _send_available_notification( 

433 *, 

434 flowsheet_id: int, 

435 scenario_id: int, 

436 simulation_object_id: int, 

437 source_task_id: int, 

438 errors: list[ResultSummaryCalculationError], 

439) -> None: 

440 payload = ResultSummaryAvailablePayload( 

441 scenario_id=scenario_id, 

442 simulation_object_id=simulation_object_id, 

443 source_task_id=source_task_id, 

444 errors=errors, 

445 ) 

446 messaging.send_flowsheet_notification_message( 

447 flowsheet_id, 

448 payload.model_dump(mode="json"), 

449 NotificationServiceMessageType.RESULT_SUMMARY_AVAILABLE, 

450 ) 

451 

452 

453def _send_failed_notification( 

454 *, 

455 flowsheet_id: int, 

456 scenario_id: int, 

457 simulation_object_id: int, 

458 source_task_id: int, 

459 errors: list[ResultSummaryCalculationError], 

460) -> None: 

461 payload = ResultSummaryFailedPayload( 

462 scenario_id=scenario_id, 

463 simulation_object_id=simulation_object_id, 

464 source_task_id=source_task_id, 

465 errors=errors, 

466 ) 

467 messaging.send_flowsheet_notification_message( 

468 flowsheet_id, 

469 payload.model_dump(mode="json"), 

470 NotificationServiceMessageType.RESULT_SUMMARY_FAILED, 

471 )