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
1 change: 1 addition & 0 deletions packages/pynumaflow-lite/manifests/sink/.dockerignore
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
.venv/
77 changes: 41 additions & 36 deletions packages/pynumaflow-lite/manifests/sink/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", "sink_log.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", "sink_log.py"]
8 changes: 1 addition & 7 deletions packages/pynumaflow-lite/manifests/sink/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +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"

17 changes: 12 additions & 5 deletions packages/pynumaflow-lite/manifests/sink/sink_log.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
import asyncio
import logging
from collections.abc import AsyncIterable

from pynumaflow_lite import sinker
from pynumaflow_lite.sinker import Sinker
from pynumaflow_lite.sinker import Datum, Response, SinkAsyncServer, Sinker

# Configure logging
logging.basicConfig(level=logging.INFO)
Expand All @@ -14,15 +14,22 @@ class SimpleLogSink(Sinker):
Simple log sink that logs each message and returns success responses.
"""

async def handler(self, datums: AsyncIterable[sinker.Datum]) -> list[sinker.Response]:
async def handler(self, datums: AsyncIterable[Datum]) -> list[Response]:
responses = []
async for msg in datums:
_LOGGER.info("User Defined Sink: %s", msg.value.decode("utf-8"))
responses.append(sinker.Response.success(msg.id))
responses.append(Response.success(msg.id))
# if we are not able to write to sink and if we have a fallback sink configured
# we can use Response.fallback(msg.id) to write the message to fallback sink
return responses


async def main() -> None:
print("Starting sink server")
# `serve` returns when SIGINT or SIGTERM arrives.
await SinkAsyncServer(SimpleLogSink()).serve()
print("Sink server stopped")


if __name__ == "__main__":
sinker.SinkAsyncServer(SimpleLogSink()).run()
asyncio.run(main())
8 changes: 8 additions & 0 deletions packages/pynumaflow-lite/manifests/sink/uv.lock

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

Original file line number Diff line number Diff line change
Expand Up @@ -123,10 +123,3 @@ async def __aexit__(
await self._task
finally:
self._task = None

def run(self) -> None:
"""Run the batchmap server in a new event loop until it stops."""
try:
asyncio.run(self.serve())
except KeyboardInterrupt:
self.stop()
7 changes: 0 additions & 7 deletions packages/pynumaflow-lite/pynumaflow_lite/_map_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,10 +123,3 @@ async def __aexit__(
await self._task
finally:
self._task = None

def run(self) -> None:
"""Run the map server in a new event loop until it stops."""
try:
asyncio.run(self.serve())
except KeyboardInterrupt:
self.stop()
Original file line number Diff line number Diff line change
Expand Up @@ -123,10 +123,3 @@ async def __aexit__(
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()
4 changes: 2 additions & 2 deletions packages/pynumaflow-lite/pynumaflow_lite/_sink_dtypes.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
from abc import ABCMeta, abstractmethod
from collections.abc import AsyncIterable
from collections.abc import AsyncIterator

from pynumaflow_lite.sinker import Datum, Response

Expand All @@ -13,7 +13,7 @@ def __call__(self, *args, **kwargs):
return self.handler(*args, **kwargs)

@abstractmethod
async def handler(self, datums: AsyncIterable[Datum]) -> list[Response]:
async def handler(self, datums: AsyncIterator[Datum]) -> list[Response]:
"""
Implement this handler function for sink.
Process the stream of datums and return responses.
Expand Down
103 changes: 75 additions & 28 deletions packages/pynumaflow-lite/pynumaflow_lite/_sink_server.py
Original file line number Diff line number Diff line change
@@ -1,50 +1,113 @@
from __future__ import annotations

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

from .pynumaflow_lite import sinker as _sinker

Datum: TypeAlias = _sinker.Datum
Response: TypeAlias = _sinker.Response

_SHUTDOWN_SIGNALS = (signal.SIGINT, signal.SIGTERM)


class SinkAsyncServer:
def __init__(
self,
handler: Any,
handler: Callable[[AsyncIterator[Datum]], Awaitable[list[Response]]],
*,
sock_file: str | None = None,
server_info_file: str | None = None,
install_signal_handlers: bool = True,
) -> None:
self._core = _sinker._SinkAsyncServer(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:
await self._core.start(self._handler)
"""Run the sink 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("sink 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("sink 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) -> SinkAsyncServer:
"""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("sink server is already serving")

self._task = asyncio.create_task(self.serve())
self._task = asyncio.create_task(self._serve(install_signal_handlers=False))
try:
await self.wait_ready()
except asyncio.CancelledError:
self.stop()
if self._task is not None:
await self._task
raise
except Exception:
except BaseException:
self.stop()
if self._task is not None:
await self._task
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

Expand All @@ -60,19 +123,3 @@ async def __aexit__(
await self._task
finally:
self._task = None

async def _main(self) -> None:
loop = asyncio.get_running_loop()
try:
loop.add_signal_handler(signal.SIGINT, self.stop)
loop.add_signal_handler(signal.SIGTERM, self.stop)
except (NotImplementedError, RuntimeError):
pass

await self.serve()

def run(self) -> None:
try:
asyncio.run(self._main())
except KeyboardInterrupt:
self.stop()
1 change: 0 additions & 1 deletion packages/pynumaflow-lite/pynumaflow_lite/batchmapper.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,6 @@ class BatchMapAsyncServer:
server_info_file: str | None = ...,
install_signal_handlers: bool = ...,
) -> None: ...
def run(self) -> None: ...
async def serve(self) -> None: ...
def stop(self) -> None: ...
async def wait_ready(self, timeout: float = ...) -> None: ...
Expand Down
1 change: 0 additions & 1 deletion packages/pynumaflow-lite/pynumaflow_lite/mapper.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,6 @@ class MapAsyncServer:
server_info_file: str | None = ...,
install_signal_handlers: bool = ...,
) -> None: ...
def run(self) -> None: ...
async def serve(self) -> None: ...
def stop(self) -> None: ...
async def wait_ready(self, timeout: float = ...) -> None: ...
Expand Down
1 change: 0 additions & 1 deletion packages/pynumaflow-lite/pynumaflow_lite/mapstreamer.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,6 @@ class MapStreamAsyncServer:
server_info_file: str | None = ...,
install_signal_handlers: bool = ...,
) -> None: ...
def run(self) -> None: ...
async def serve(self) -> None: ...
def stop(self) -> None: ...
async def wait_ready(self, timeout: float = ...) -> None: ...
Expand Down
Loading
Loading