(self, func)
| 64 | |
| 65 | @_common_def(inputs=["inp:[]"], outputs=["out"]) |
| 66 | def reduce_def(self, func): |
| 67 | ret = [] |
| 68 | for inp in self.inp: |
| 69 | ret.append(inp.recv()) |
| 70 | |
| 71 | all_empty = True |
| 72 | for envelope in ret: |
| 73 | all_empty = all_empty and envelope is None |
| 74 | if all_empty: |
| 75 | return |
| 76 | |
| 77 | for envelope in ret: |
| 78 | assert envelope is not None |
| 79 | |
| 80 | msgs = [ envelope.msg for envelope in ret ] |
| 81 | |
| 82 | self.out.send(ret[0].repack(func(msgs))) |
| 83 | |
| 84 | |
| 85 | @_common_def(inputs=["inp"]) |