/opt/imh-python/lib/python3.9/site-packages/celery/bin
NameSizeModeActions
__pycache__/-0755rm
amqp.py100230644editdlrm
base.py91730644editdlrm
beat.py25920644editdlrm
call.py23700644editdlrm
celery.py74400644editdlrm
control.py86450644editdlrm
events.py27940644editdlrm
graph.py57960644editdlrm
list.py10580644editdlrm
logtool.py42670644editdlrm
migrate.py21080644editdlrm
multi.py153740644editdlrm
purge.py26080644editdlrm
result.py9760644editdlrm
shell.py48390644editdlrm
upgrade.py30640644editdlrm
worker.py128860644editdlrm
__init__.py00644editdlrm
Edit: /opt/imh-python/lib/python3.9/site-packages/celery/bin/migrate.py (2108B)
"""The ``celery migrate`` command, used to filter and move messages.""" import click from kombu import Connection from celery.bin.base import CeleryCommand, CeleryOption, handle_preload_options from celery.contrib.migrate import migrate_tasks @click.command(cls=CeleryCommand) @click.argument('source') @click.argument('destination') @click.option('-n', '--limit', cls=CeleryOption, type=int, help_group='Migration Options', help='Number of tasks to consume.') @click.option('-t', '--timeout', cls=CeleryOption, type=float, help_group='Migration Options', help='Timeout in seconds waiting for tasks.') @click.option('-a', '--ack-messages', cls=CeleryOption, is_flag=True, help_group='Migration Options', help='Ack messages from source broker.') @click.option('-T', '--tasks', cls=CeleryOption, help_group='Migration Options', help='List of task names to filter on.') @click.option('-Q', '--queues', cls=CeleryOption, help_group='Migration Options', help='List of queues to migrate.') @click.option('-F', '--forever', cls=CeleryOption, is_flag=True, help_group='Migration Options', help='Continually migrate tasks until killed.') @click.pass_context @handle_preload_options def migrate(ctx, source, destination, **kwargs): """Migrate tasks from one broker to another. Warning: This command is experimental, make sure you have a backup of the tasks before you continue. """ # TODO: Use a progress bar def on_migrate_task(state, body, message): ctx.obj.echo(f"Migrating task {state.count}/{state.strtotal}: {body}") migrate_tasks(Connection(source), Connection(destination), callback=on_migrate_task, **kwargs)