I have a celery tasks setup to check if internal messaging in my app is working or not. I have two tasks send_ping and look_for_pong for the same. Every time send_ping is executed it triggers look_for_pong after message is sent to keep looking for response. If response is found within 5 minutes, it schedules another send_ping task with a delay of 30 minutes. If not, it fires an alert and immediately schedules another send_ping task.
I am trying to build endpoints to start and stop the messaging bots. I am able to successfully start the bot but facing issues in stopping it. It stops the bot eventually after 3-4 ping-pong exchange but does not stop immediately. Pasting the relevant code below.
ping pong celery tasks
@celery.task(priority=HIGH_PRIORITY)
def send_ping():
sent_time = datetime.utcnow()
send_conversation_section(
send_to=subscriber, first_step=conversation_step
)
look_for_pong.delay(ping_time=sent_time)
@celery.task(priority=HIGH_PRIORITY)
def look_for_pong(ping_time):
if pong_response_message:
send_ping.apply_async(countdown=1800)
elif (datetime.utcnow() - ping_time).seconds > 300:
sentry_sdk.capture_message(f"ALERT: Pong not received. Messaging may be down!!")
send_ping.delay()
else:
look_for_pong.apply_async(args=[ping_time],countdown=20)
endpoint to start and kill the bot
@health.route("/api/ping_bot/start", methods=["POST"])
def start_ping_bot():
send_ping.delay()
return make_response(jsonify({"message": "Ping bot started"}), 200)
@health.route("/api/ping_bot/stop", methods=["POST"])
def stop_ping_bot():
query = celery.events.State().tasks_by_type('send_ping')
for uuid, task, status in query:
celery.control.revoke(uuid, terminate=True)
return make_response(jsonify({"message": "Ping bot stopped"}),200)