futures = {executor.submit(run_trajectory, t): t for t in trajectories} for fut in as_completed(futures): try: collect(fut.result(timeout=120)) # per-trajectory watchdog except TimeoutError: mark(futures[fut], "timeout") # record it, keep sweeping executor.shutdown(wait=False, cancel_futures=True) # never wait on a stuck thread