Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 40 additions & 9 deletions awx/main/models/workflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,9 @@
import json
import logging
import multiprocessing
import os
import signal
import time
from uuid import uuid4
from copy import copy
from urllib.parse import urljoin
Expand Down Expand Up @@ -1206,6 +1209,25 @@ def _render_context_template(conn, template_str, variables):
conn.send(('error', '{}: {}'.format(type(e).__name__, e)[:1024]))


def _reap_render_child(pid, grace=1.0):
"""Wait briefly for the forked render child to exit, then kill it."""
deadline = time.monotonic() + grace
while True:
try:
Comment thread
cigamit marked this conversation as resolved.
if os.waitpid(pid, os.WNOHANG)[0] == pid:
return
except ChildProcessError:
return
if time.monotonic() >= deadline:
break
time.sleep(0.02)
try:
os.kill(pid, signal.SIGKILL)
os.waitpid(pid, 0)
except (ChildProcessError, ProcessLookupError):
pass


class WorkflowApproval(UnifiedJob, JobNotificationMixin):
class Meta:
app_label = 'main'
Expand Down Expand Up @@ -1283,11 +1305,22 @@ def render_context_message(self, ancestor_artifacts):
if not template_str:
return
rendered = None
ctx = multiprocessing.get_context('fork')
parent_conn, child_conn = ctx.Pipe(duplex=False)
worker = ctx.Process(target=_render_context_template, args=(child_conn, template_str, dict(ancestor_artifacts or {})))
# The task manager runs inside a dispatcher pool worker, which is a daemonic
# multiprocessing.Process. multiprocessing refuses to start children from one
# ("daemonic processes are not allowed to have children"), so fork directly;
# the Pipe is only used for its picklable Connection objects.
parent_conn, child_conn = multiprocessing.get_context('fork').Pipe(duplex=False)
pid = None
try:
worker.start()
pid = os.fork()
if pid == 0:
# Child: render, report back, and _exit so the parent's atexit
# handlers and inherited database connection are left untouched.
try:
parent_conn.close()
_render_context_template(child_conn, template_str, dict(ancestor_artifacts or {}))
finally:
os._exit(0)
child_conn.close()
if parent_conn.poll(CONTEXT_TEMPLATE_TIMEOUT):
status, payload = parent_conn.recv()
Expand All @@ -1303,11 +1336,9 @@ def render_context_message(self, ancestor_artifacts):
logger.exception('Unexpected error rendering context_template for approval %s', self.pk)
finally:
parent_conn.close()
if worker.pid is not None:
worker.join(1)
if worker.is_alive():
worker.kill()
worker.join()
child_conn.close() # normally closed right after the fork; also covers fork() itself failing
if pid:
_reap_render_child(pid)
Comment thread
cigamit marked this conversation as resolved.
if rendered and rendered.strip():
self.context_message = rendered
self.save(update_fields=['context_message'])
Expand Down
8 changes: 8 additions & 0 deletions awx/main/tests/functional/models/test_workflow.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,6 @@
# Python
import multiprocessing

import pytest
from unittest import mock
import json
Expand Down Expand Up @@ -973,6 +975,12 @@ def test_non_identifier_artifact_keys(self, approval):
# set_stats keys are arbitrary strings, they must not break rendering
assert self._render(approval, '{{ ok }}', {'not-an-identifier!': 1, 'ok': 'yes'}) == 'yes'

def test_renders_inside_daemonic_process(self, approval, monkeypatch):
# dispatcher pool workers are daemonic multiprocessing children, and
# multiprocessing.Process.start() raises inside one
monkeypatch.setattr(multiprocessing.current_process(), 'daemon', True)
assert self._render(approval, 'context: {{ env }}', {'env': 'prod'}) == 'context: prod'

def test_template_error_does_not_raise(self, approval):
assert self._render(approval, '{% if %}', {'x': 1}) == ''

Expand Down
Loading