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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
.venv/
77 changes: 41 additions & 36 deletions packages/pynumaflow-lite/manifests/mapstream/Dockerfile
Original file line number Diff line number Diff line change
@@ -1,38 +1,43 @@
FROM python:3.11-slim-bullseye AS builder

ENV PYTHONFAULTHANDLER=1 \
PYTHONUNBUFFERED=1 \
PYTHONHASHSEED=random \
PIP_NO_CACHE_DIR=on \
PIP_DISABLE_PIP_VERSION_CHECK=on \
PIP_DEFAULT_TIMEOUT=100 \
POETRY_HOME="/opt/poetry" \
POETRY_VIRTUALENVS_IN_PROJECT=true \
POETRY_NO_INTERACTION=1 \
PYSETUP_PATH="/opt/pysetup"

ENV PATH="$POETRY_HOME/bin:$PATH"

RUN apt-get update \
&& apt-get install --no-install-recommends -y \
curl \
wget \
# deps for building python deps
build-essential \
&& apt-get install -y git \
&& apt-get clean && rm -rf /var/lib/apt/lists/* \
&& curl -sSL https://install.python-poetry.org | python3 -

FROM builder AS udf

WORKDIR $PYSETUP_PATH
COPY ./ ./

RUN pip install $PYSETUP_PATH/pynumaflow_lite-0.1.0-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl

RUN poetry lock
RUN poetry install --no-cache --no-root && \
rm -rf ~/.cache/pypoetry/
FROM python:3.11-slim-trixie AS builder

CMD ["python", "mapstream_cat.py"]
COPY --from=ghcr.io/astral-sh/uv:0.12.13 /uv /uvx /bin/

ENV UV_COMPILE_BYTECODE=1 \
UV_LINK_MODE=copy \
UV_PYTHON_DOWNLOADS=never

WORKDIR /app

RUN --mount=type=cache,target=/root/.cache/uv \
--mount=type=bind,source=uv.lock,target=uv.lock \
--mount=type=bind,source=pyproject.toml,target=pyproject.toml \
uv sync --locked --no-install-project --no-dev

COPY . .

RUN --mount=type=cache,target=/root/.cache/uv \
uv sync --locked --no-dev

# Install the local pynumaflow-lite wheel that matches the build platform.
RUN uv pip install --no-index --find-links . pynumaflow-lite

FROM python:3.11-slim-trixie

# Setup a non-root user
RUN groupadd --system --gid 999 nonroot \
&& useradd --system --gid 999 --uid 999 --create-home nonroot

COPY --from=builder --chown=nonroot:nonroot /app /app

ENV PATH="/app/.venv/bin:$PATH"

# Keeps Python from buffering stdout and stderr to avoid situations where
# the application crashes without emitting any logs due to buffering.
ENV PYTHONUNBUFFERED=1

# Use the non-root user to run our application
USER nonroot

WORKDIR /app

CMD ["python", "mapstream_cat.py"]
43 changes: 13 additions & 30 deletions packages/pynumaflow-lite/manifests/mapstream/mapstream_cat.py
Original file line number Diff line number Diff line change
@@ -1,41 +1,24 @@
import asyncio
import signal
from collections.abc import AsyncIterator, Callable
from collections.abc import AsyncIterator

from pynumaflow_lite import mapstreamer
from pynumaflow_lite.mapstreamer import Message
from pynumaflow_lite.mapstreamer import Datum, MapStreamAsyncServer, MapStreamer, Message


class SimpleStreamCat(mapstreamer.MapStreamer):
async def handler(self, keys: list[str], datum: mapstreamer.Datum) -> AsyncIterator[Message]:
parts = datum.value.decode("utf-8").split(",")
if not parts:
class SimpleStreamCat(MapStreamer):
async def handler(self, datum: Datum) -> AsyncIterator[Message]:
if not datum.value:
yield Message.to_drop()
return
for s in parts:
yield Message(s.encode(), keys)
for s in datum.value.decode("utf-8").split(","):
yield Message(s.encode(), keys=datum.keys)


async def start(f: Callable[[list[str], mapstreamer.Datum], AsyncIterator[Message]]):
# Use default socket/info file locations; no explicit sock file passed
server = mapstreamer.MapStreamAsyncServer()

# Register loop-level signal handlers so we control shutdown and avoid asyncio.run noise.
loop = asyncio.get_running_loop()
try:
loop.add_signal_handler(signal.SIGINT, lambda: server.stop())
loop.add_signal_handler(signal.SIGTERM, lambda: server.stop())
except (NotImplementedError, RuntimeError):
pass

try:
await server.start(f)
print("Shutting down gracefully...")
except asyncio.CancelledError:
server.stop()
return
async def main() -> None:
print("Starting MapStream server")
# `serve` returns when SIGINT or SIGTERM arrives.
await MapStreamAsyncServer(SimpleStreamCat()).serve()
print("MapStream server stopped")


if __name__ == "__main__":
async_handler = SimpleStreamCat()
asyncio.run(start(async_handler))
asyncio.run(main())
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,6 @@ authors = [
{ name = "Vigith Maurice", email = "vigith@gmail.com" }
]
readme = "README.md"
requires-python = ">=3.11"
requires-python = "==3.11.*"
dependencies = [
]


[build-system]
requires = ["poetry-core>=2.0.0,<3.0.0"]
build-backend = "poetry.core.masonry.api"
8 changes: 8 additions & 0 deletions packages/pynumaflow-lite/manifests/mapstream/uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions packages/pynumaflow-lite/pynumaflow_lite/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from ._map_dtypes import Mapper
from ._map_server import MapAsyncServer
from ._mapstream_dtypes import MapStreamer
from ._mapstream_server import MapStreamAsyncServer
from ._reduce_dtypes import Reducer
from ._reducestreamer_dtypes import ReduceStreamer
from ._session_reduce_dtypes import SessionReducer
Expand Down Expand Up @@ -41,6 +42,7 @@
batchmapper.BatchMapper = BatchMapper
batchmapper.BatchMapAsyncServer = BatchMapAsyncServer
mapstreamer.MapStreamer = MapStreamer
mapstreamer.MapStreamAsyncServer = MapStreamAsyncServer
reducer.Reducer = Reducer
session_reducer.SessionReducer = SessionReducer
reducestreamer.ReduceStreamer = ReduceStreamer
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,10 @@ def __call__(self, *args, **kwargs):
return self.handler(*args, **kwargs)

@abstractmethod
async def handler(self, keys: list[str], datum: Datum) -> AsyncIterator[Message]:
async def handler(self, datum: Datum) -> AsyncIterator[Message]:
"""
Implement this handler function for streaming mapping.
It should be an async generator yielding Message objects.
"""
pass
raise NotImplementedError
yield # makes this an async generator, so overrides that yield type-check
132 changes: 132 additions & 0 deletions packages/pynumaflow-lite/pynumaflow_lite/_mapstream_server.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
from __future__ import annotations

import asyncio
import contextlib
import signal
from collections.abc import AsyncIterator, Callable
from types import TracebackType
from typing import TypeAlias

from .pynumaflow_lite import mapstreamer as _mapstreamer

Datum: TypeAlias = _mapstreamer.Datum
Message: TypeAlias = _mapstreamer.Message

_SHUTDOWN_SIGNALS = (signal.SIGINT, signal.SIGTERM)


class MapStreamAsyncServer:
def __init__(
self,
handler: Callable[[Datum], AsyncIterator[Message]],
*,
sock_file: str | None = None,
server_info_file: str | None = None,
install_signal_handlers: bool = True,
) -> None:
self._core = _mapstreamer._MapStreamAsyncServer(sock_file, server_info_file)
self._handler = handler
self._install_signal_handlers = install_signal_handlers
self._task: asyncio.Task[None] | None = None
self._serving = False
self._installed_signals: list[signal.Signals] = []

async def serve(self) -> None:
"""Run the mapstream server until it stops.

This is the entrypoint for an application that already runs an event
loop. It returns when a shutdown signal arrives or when `stop()` runs.
"""
await self._serve(install_signal_handlers=self._install_signal_handlers)

async def _serve(self, *, install_signal_handlers: bool) -> None:
if self._serving:
raise RuntimeError("mapstream server is already serving")
self._serving = True
try:
if install_signal_handlers:
self._add_signal_handlers()
await self._core.start(self._handler)
finally:
self._remove_signal_handlers()
self._serving = False

def stop(self) -> None:
self._core.stop()

async def wait_ready(self, timeout: float = 30.0) -> None:
await self._core.wait_ready(timeout)

async def wait_for_termination(self) -> None:
"""Wait until the background server task ends.

Use this inside an `async with` block. It raises the handler error if
the server task failed.
"""
if self._task is None:
raise RuntimeError("mapstream server is not serving")
await asyncio.shield(self._task)

def _add_signal_handlers(self) -> None:
loop = asyncio.get_running_loop()
for sig in _SHUTDOWN_SIGNALS:
try:
loop.add_signal_handler(sig, self.stop)
except (NotImplementedError, RuntimeError, OSError):
continue
self._installed_signals.append(sig)

def _remove_signal_handlers(self) -> None:
if not self._installed_signals:
return
try:
loop = asyncio.get_running_loop()
except RuntimeError:
self._installed_signals.clear()
return
for sig in self._installed_signals:
with contextlib.suppress(NotImplementedError, OSError):
loop.remove_signal_handler(sig)
self._installed_signals.clear()

async def __aenter__(self) -> MapStreamAsyncServer:
"""Start the server in a background task and wait until it is ready.

This form is for tests and for code that must run other work next to
the server. It never installs signal handlers.
"""
if self._task is not None and not self._task.done():
raise RuntimeError("mapstream server is already serving")

self._task = asyncio.create_task(self._serve(install_signal_handlers=False))
try:
await self.wait_ready()
except BaseException:
self.stop()
task, self._task = self._task, None
if task is not None:
# Surface the server error, if there is one. It explains the
# failure better than the `wait_ready` error does.
await task
raise
return self

async def __aexit__(
self,
exc_type: type[BaseException] | None,
exc: BaseException | None,
tb: TracebackType | None,
) -> None:
self.stop()
if self._task is not None:
try:
await self._task
finally:
self._task = None

def run(self) -> None:
"""Run the mapstream server in a new event loop until it stops."""
try:
asyncio.run(self.serve())
except KeyboardInterrupt:
self.stop()
Loading
Loading