Update execution state after a task finishes Mutates. This should run atomically (with a lock).
(
dsk, key, state, results, sortkey, delete=True, release_data=release_data
)
| 288 | |
| 289 | |
| 290 | def finish_task( |
| 291 | dsk, key, state, results, sortkey, delete=True, release_data=release_data |
| 292 | ): |
| 293 | """ |
| 294 | Update execution state after a task finishes |
| 295 | |
| 296 | Mutates. This should run atomically (with a lock). |
| 297 | """ |
| 298 | for dep in sorted(state["dependents"][key], key=sortkey, reverse=True): |
| 299 | s = state["waiting"][dep] |
| 300 | s.remove(key) |
| 301 | if not s: |
| 302 | del state["waiting"][dep] |
| 303 | state["ready"].append(dep) |
| 304 | |
| 305 | for dep in state["dependencies"][key]: |
| 306 | if dep in state["waiting_data"]: |
| 307 | s = state["waiting_data"][dep] |
| 308 | s.remove(key) |
| 309 | if not s and dep not in results: |
| 310 | release_data(dep, state, delete=delete) |
| 311 | elif delete and dep not in results: |
| 312 | release_data(dep, state, delete=delete) |
| 313 | |
| 314 | state["finished"].add(key) |
| 315 | state["running"].remove(key) |
| 316 | |
| 317 | return state |
| 318 | |
| 319 | |
| 320 | def nested_get(ind, coll): |