(queries, f)
| 38 | return False |
| 39 | |
| 40 | def run_concurrent(queries, f): |
| 41 | pool = Pool(nodes=CLIENT_COUNT) |
| 42 | manager = pathos_multiprocess.Manager() |
| 43 | |
| 44 | barrier = manager.Barrier(CLIENT_COUNT) |
| 45 | barriers = [barrier] * CLIENT_COUNT |
| 46 | |
| 47 | # invoke queries |
| 48 | results = pool.map(f, queries, barriers) |
| 49 | |
| 50 | pool.clear() |
| 51 | |
| 52 | return results |
| 53 | |
| 54 | class testConcurrentQueryFlow(FlowTestsBase): |
| 55 | def __init__(self): |
no test coverage detected