diff options
| author | Jenkins <jenkins@review.openstack.org> | 2015-07-08 23:36:35 +0000 |
|---|---|---|
| committer | Gerrit Code Review <review@openstack.org> | 2015-07-08 23:36:35 +0000 |
| commit | b77eac140fb7684ae283a40099cfac3324b2bddd (patch) | |
| tree | e9a2c040c1088a8de48ebe15ca6b5dd28cca4125 /taskflow/examples | |
| parent | 17b819bae21fbd27ffdb9575930f8ad37a8fb437 (diff) | |
| parent | 40d19c7696f1e0b7d75eacbd271974ee9155c019 (diff) | |
| download | taskflow-b77eac140fb7684ae283a40099cfac3324b2bddd.tar.gz | |
Merge "Handle conductor ctrl-c more appropriately"
Diffstat (limited to 'taskflow/examples')
| -rw-r--r-- | taskflow/examples/99_bottles.py | 54 |
1 files changed, 35 insertions, 19 deletions
diff --git a/taskflow/examples/99_bottles.py b/taskflow/examples/99_bottles.py index 9959255..90894e9 100644 --- a/taskflow/examples/99_bottles.py +++ b/taskflow/examples/99_bottles.py @@ -54,33 +54,47 @@ JB_CONF = { 'board': 'zookeeper', 'path': '/taskflow/99-bottles-demo', } -DB_URI = r"sqlite:////tmp/bottles.db" -PART_DELAY = 1.0 +PERSISTENCE_URI = r"sqlite:////tmp/bottles.db" +TAKE_DOWN_DELAY = 1.0 +PASS_AROUND_DELAY = 3.0 HOW_MANY_BOTTLES = 99 -class TakeABottleDownPassItAround(task.Task): - def execute(self, bottles_left): +class TakeABottleDown(task.Task): + def execute(self): sys.stdout.write('Take one down, ') - time.sleep(PART_DELAY) + sys.stdout.flush() + time.sleep(TAKE_DOWN_DELAY) + + +class PassItAround(task.Task): + def execute(self): sys.stdout.write('pass it around, ') - time.sleep(PART_DELAY) + sys.stdout.flush() + time.sleep(PASS_AROUND_DELAY) + + +class Conclusion(task.Task): + def execute(self, bottles_left): sys.stdout.write('%s bottles of beer on the wall...\n' % bottles_left) + sys.stdout.flush() def make_bottles(count): s = lf.Flow("bottle-song") for bottle in reversed(list(range(1, count + 1))): - t = TakeABottleDownPassItAround("take-bottle-%s" % bottle, - inject={"bottles_left": bottle - 1}) - s.add(t) + take_bottle = TakeABottleDown("take-bottle-%s" % bottle) + pass_it = PassItAround("pass-%s-around" % bottle) + next_bottles = Conclusion("next-bottles-%s" % (bottle - 1), + inject={"bottles_left": bottle - 1}) + s.add(take_bottle, pass_it, next_bottles) return s def run_conductor(): print("Starting conductor with pid: %s" % ME) my_name = "conductor-%s" % ME - persist_backend = persistence_backends.fetch(DB_URI) + persist_backend = persistence_backends.fetch(PERSISTENCE_URI) with contextlib.closing(persist_backend): with contextlib.closing(persist_backend.get_connection()) as conn: conn.upgrade() @@ -90,17 +104,18 @@ def run_conductor(): with contextlib.closing(job_backend): cond = conductor_backends.fetch('blocking', my_name, job_backend, persistence=persist_backend) - # Run forever, and kill -9 me... - # - # TODO(harlowja): it would be nicer if we could handle - # ctrl-c better... - cond.run() + # Run forever, and kill -9 or ctrl-c me... + try: + cond.run() + finally: + cond.stop() + cond.wait() def run_poster(): print("Starting poster with pid: %s" % ME) my_name = "poster-%s" % ME - persist_backend = persistence_backends.fetch(DB_URI) + persist_backend = persistence_backends.fetch(PERSISTENCE_URI) with contextlib.closing(persist_backend): with contextlib.closing(persist_backend.get_connection()) as conn: conn.upgrade() @@ -128,11 +143,12 @@ def run_poster(): def main(): if len(sys.argv) == 1: sys.stderr.write("%s p|c\n" % os.path.basename(sys.argv[0])) - return - if sys.argv[1] == 'p': + elif sys.argv[1] == 'p': run_poster() - if sys.argv[1] == 'c': + elif sys.argv[1] == 'c': run_conductor() + else: + sys.stderr.write("%s p|c\n" % os.path.basename(sys.argv[0])) if __name__ == '__main__': |
