"""Pipeline simulator 5,000,000.""" from video.logic.activity_task_worker_pool import merge_if_multiple_inputs def simulated_pipeline(functions, inputs={}): """Pipeline simulator 5,000,000. Chain activity task handlers together without using a state machine or running any dev server. Args: functions: Activity task handler functions you want to call and chain each function's outputs to the next function's inputs. inputs: Initial inputs. """ for f in functions: outputs = f({ **merge_if_multiple_inputs(inputs), '{}_job_id'.format(f.__name__): 123, }) inputs = {**inputs, **outputs} return merge_if_multiple_inputs(inputs)