Skip to content
Open
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
155 changes: 86 additions & 69 deletions sorts/odd_even_transposition_parallel.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,93 +12,105 @@
"""

import multiprocessing as mp
from typing import Any, Protocol

# lock used to ensure that two processes do not access a pipe at the same time
# NOTE This breaks testing on build runner. May work better locally
# process_lock = mp.Lock()

"""
The function run by the processes that sorts the list

position = the position in the list the process represents, used to know which
neighbor we pass our value to
value = the initial value at list[position]
LSend, RSend = the pipes we use to send to our left and right neighbors
LRcv, RRcv = the pipes we use to receive from our left and right neighbors
resultPipe = the pipe used to send results back to main
"""

class Comparable(Protocol):
def __lt__(self, other: Any, /) -> bool: ...


def oe_process(
position,
value,
l_send,
r_send,
lr_cv,
rr_cv,
result_pipe,
multiprocessing_context,
def oe_process[T: Comparable](
position: int,
value: T,
l_send: tuple[mp.connection.Connection, mp.connection.Connection] | None,

Check warning on line 29 in sorts/odd_even_transposition_parallel.py

View workflow job for this annotation

GitHub Actions / ty

ty (possibly-missing-submodule)

sorts/odd_even_transposition_parallel.py:29:45: possibly-missing-submodule: Submodule `connection` might not have been imported help: Consider explicitly importing `multiprocessing.connection`

Check warning on line 29 in sorts/odd_even_transposition_parallel.py

View workflow job for this annotation

GitHub Actions / ty

ty (possibly-missing-submodule)

sorts/odd_even_transposition_parallel.py:29:19: possibly-missing-submodule: Submodule `connection` might not have been imported help: Consider explicitly importing `multiprocessing.connection`
r_send: tuple[mp.connection.Connection, mp.connection.Connection] | None,

Check warning on line 30 in sorts/odd_even_transposition_parallel.py

View workflow job for this annotation

GitHub Actions / ty

ty (possibly-missing-submodule)

sorts/odd_even_transposition_parallel.py:30:45: possibly-missing-submodule: Submodule `connection` might not have been imported help: Consider explicitly importing `multiprocessing.connection`

Check warning on line 30 in sorts/odd_even_transposition_parallel.py

View workflow job for this annotation

GitHub Actions / ty

ty (possibly-missing-submodule)

sorts/odd_even_transposition_parallel.py:30:19: possibly-missing-submodule: Submodule `connection` might not have been imported help: Consider explicitly importing `multiprocessing.connection`
lr_cv: tuple[mp.connection.Connection, mp.connection.Connection] | None,

Check warning on line 31 in sorts/odd_even_transposition_parallel.py

View workflow job for this annotation

GitHub Actions / ty

ty (possibly-missing-submodule)

sorts/odd_even_transposition_parallel.py:31:44: possibly-missing-submodule: Submodule `connection` might not have been imported help: Consider explicitly importing `multiprocessing.connection`

Check warning on line 31 in sorts/odd_even_transposition_parallel.py

View workflow job for this annotation

GitHub Actions / ty

ty (possibly-missing-submodule)

sorts/odd_even_transposition_parallel.py:31:18: possibly-missing-submodule: Submodule `connection` might not have been imported help: Consider explicitly importing `multiprocessing.connection`
rr_cv: tuple[mp.connection.Connection, mp.connection.Connection] | None,
result_pipe: tuple[mp.connection.Connection, mp.connection.Connection],
multiprocessing_context: Any,
) -> None:
"""
The function run by the processes that sorts the list.

position = the position in the list the process represents, used to know
which neighbor we pass our value to
value = the initial value at list[position]
l_send, r_send = the pipes we use to send to our left and right neighbors
lr_cv, rr_cv = the pipes we use to receive from our left and right
neighbors
result_pipe = the pipe used to send results back to main
"""
process_lock = multiprocessing_context.Lock()

# we perform n swaps since after n swaps we know we are sorted
# we *could* stop early if we are sorted already, but it takes as long to
# find out we are sorted as it does to sort the list with this algorithm
for i in range(10):
if (i + position) % 2 == 0 and r_send is not None:
# send your value to your right neighbor
with process_lock:
r_send[1].send(value)

# receive your right neighbor's value
with process_lock:
temp = rr_cv[0].recv()

# take the lower value since you are on the left
value = min(value, temp)
elif (i + position) % 2 != 0 and l_send is not None:
# send your value to your left neighbor
with process_lock:
l_send[1].send(value)

# receive your left neighbor's value
with process_lock:
temp = lr_cv[0].recv()

# take the higher value since you are on the right
value = max(value, temp)
# after all swaps are performed, send the values back to main
result_pipe[1].send(value)


"""
the function which creates the processes that perform the parallel swaps

arr = the list to be sorted
"""


def odd_even_transposition(arr):
try:
for i in range(10):
if (i + position) % 2 == 0 and r_send is not None:
# send your value to your right neighbor
with process_lock:
r_send[1].send(value)

# receive your right neighbor's value
with process_lock:
temp = rr_cv[0].recv()

# take the lower value since you are on the left
value = min(value, temp)
elif (i + position) % 2 != 0 and l_send is not None:
# send your value to your left neighbor
with process_lock:
l_send[1].send(value)

# receive your left neighbor's value
with process_lock:
temp = lr_cv[0].recv()

# take the higher value since you are on the right
value = max(value, temp)
# after all swaps are performed, send the value back to main
result_pipe[1].send((False, value))
except Exception as e: # noqa: BLE001
result_pipe[1].send((True, e))


def odd_even_transposition[T: Comparable](arr: list[T]) -> list[T]:
"""
>>> odd_even_transposition(list(range(10)[::-1])) == sorted(list(range(10)[::-1]))
True
>>> odd_even_transposition(["a", "x", "c"]) == sorted(["x", "a", "c"])
True
>>> odd_even_transposition([1.9, 42.0, 2.8]) == sorted([1.9, 42.0, 2.8])
True
>>> odd_even_transposition([False, True, False]) == sorted([False, False, True])
True
>>> odd_even_transposition([1, 32.0, 9]) == sorted([False, False, True])
False
>>> odd_even_transposition([1, 32.0, 9]) == sorted([1.0, 32, 9.0])
Sort a list of comparable items using the parallel odd-even transposition
sort algorithm.

Each element is represented by a separate process; neighboring processes
exchange values in alternating odd/even rounds via message-passing pipes.
Items must be mutually comparable (support ``<``); passing a list whose
elements cannot be compared with each other (e.g. mixing ``int`` and
``str``) raises ``TypeError`` inside the worker processes.

:param arr: a list of mutually comparable items
:return: the same list sorted in ascending order

Examples:
>>> odd_even_transposition([5, 4, 3, 2, 1])
[1, 2, 3, 4, 5]
>>> odd_even_transposition([3, 3, 1, 2, 2, 1])
[1, 1, 2, 2, 3, 3]
>>> odd_even_transposition(['c', 'a', 'b'])
['a', 'b', 'c']
>>> odd_even_transposition([3.3, 1.1, 2.2])
[1.1, 2.2, 3.3]
>>> odd_even_transposition(list(range(10)[::-1])) == sorted(range(10))
True
>>> unsorted_list = [-442, -98, -554, 266, -491, 985, -53, -529, 82, -429]
>>> odd_even_transposition(unsorted_list) == sorted(unsorted_list)
True
>>> unsorted_list = [-442, -98, -554, 266, -491, 985, -53, -529, 82, -429]
>>> odd_even_transposition(unsorted_list) == sorted(unsorted_list + [1])
False
"""
if not arr:
return arr

# spawn method is considered safer than fork
multiprocessing_context = mp.get_context("spawn")

Expand Down Expand Up @@ -173,7 +185,12 @@

# wait for the processes to end and write their values to the list
for p in range(len(result_pipe)):
arr[p] = result_pipe[p][0].recv()
is_error, result = result_pipe[p][0].recv()
if is_error:
for proc in process_array_:
proc.terminate()
raise result
arr[p] = result
process_array_[p].join()
return arr

Expand Down
5 changes: 5 additions & 0 deletions tests/test_sorts.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,9 @@
from sorts.merge_insertion_sort import merge_insertion_sort
from sorts.merge_sort import merge_sort
from sorts.odd_even_sort import odd_even_sort
from sorts.odd_even_transposition_parallel import (
odd_even_transposition as odd_even_transposition_parallel,
)
from sorts.odd_even_transposition_single_threaded import odd_even_transposition
from sorts.pancake_sort import pancake_sort
from sorts.patience_sort import patience_sort
Expand Down Expand Up @@ -116,6 +119,7 @@ def test_intro_sort_heap_fallback_preserves_surrounding_items(max_depth: int) ->
merge_sort,
odd_even_sort,
odd_even_transposition,
odd_even_transposition_parallel,
pancake_sort,
patience_sort,
power_sort,
Expand Down Expand Up @@ -198,6 +202,7 @@ def test_rec_insertion_sort(case) -> None:
merge_sort,
odd_even_sort,
odd_even_transposition,
odd_even_transposition_parallel,
pancake_sort,
patience_sort,
power_sort,
Expand Down
Loading