From 715bbcba0d6c8a7074f8baaae7318ac7aebf1b15 Mon Sep 17 00:00:00 2001 From: Sreekanth Date: Sun, 27 Sep 2026 18:24:04 +0530 Subject: [PATCH 1/3] Remove sync run() method Signed-off-by: Sreekanth --- .../manifests/sink/sink_log.py | 17 ++- .../pynumaflow_lite/_batchmap_server.py | 7 -- .../pynumaflow_lite/_map_server.py | 7 -- .../pynumaflow_lite/_mapstream_server.py | 7 -- .../pynumaflow_lite/_sink_server.py | 103 +++++++++++++----- .../pynumaflow_lite/batchmapper.pyi | 1 - .../pynumaflow_lite/mapper.pyi | 1 - .../pynumaflow_lite/mapstreamer.pyi | 1 - .../pynumaflow_lite/sinker.pyi | 3 +- .../pynumaflow-lite/tests/examples/map_cat.py | 12 +- .../tests/examples/map_cat_class.py | 12 +- .../tests/examples/sink_log.py | 11 +- .../tests/examples/sink_log_class.py | 11 +- 13 files changed, 123 insertions(+), 70 deletions(-) diff --git a/packages/pynumaflow-lite/manifests/sink/sink_log.py b/packages/pynumaflow-lite/manifests/sink/sink_log.py index 1b1c656c..614c9956 100644 --- a/packages/pynumaflow-lite/manifests/sink/sink_log.py +++ b/packages/pynumaflow-lite/manifests/sink/sink_log.py @@ -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) @@ -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()) diff --git a/packages/pynumaflow-lite/pynumaflow_lite/_batchmap_server.py b/packages/pynumaflow-lite/pynumaflow_lite/_batchmap_server.py index 32b15304..69f33237 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/_batchmap_server.py +++ b/packages/pynumaflow-lite/pynumaflow_lite/_batchmap_server.py @@ -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() diff --git a/packages/pynumaflow-lite/pynumaflow_lite/_map_server.py b/packages/pynumaflow-lite/pynumaflow_lite/_map_server.py index ff477680..04b5d97e 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/_map_server.py +++ b/packages/pynumaflow-lite/pynumaflow_lite/_map_server.py @@ -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() diff --git a/packages/pynumaflow-lite/pynumaflow_lite/_mapstream_server.py b/packages/pynumaflow-lite/pynumaflow_lite/_mapstream_server.py index c4d25ace..31ac047a 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/_mapstream_server.py +++ b/packages/pynumaflow-lite/pynumaflow_lite/_mapstream_server.py @@ -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() diff --git a/packages/pynumaflow-lite/pynumaflow_lite/_sink_server.py b/packages/pynumaflow-lite/pynumaflow_lite/_sink_server.py index 4dd29892..c8ce034a 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/_sink_server.py +++ b/packages/pynumaflow-lite/pynumaflow_lite/_sink_server.py @@ -1,27 +1,55 @@ 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() @@ -29,22 +57,57 @@ def stop(self) -> None: 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 @@ -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() diff --git a/packages/pynumaflow-lite/pynumaflow_lite/batchmapper.pyi b/packages/pynumaflow-lite/pynumaflow_lite/batchmapper.pyi index a60d81b5..e87630d8 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/batchmapper.pyi +++ b/packages/pynumaflow-lite/pynumaflow_lite/batchmapper.pyi @@ -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: ... diff --git a/packages/pynumaflow-lite/pynumaflow_lite/mapper.pyi b/packages/pynumaflow-lite/pynumaflow_lite/mapper.pyi index 9a18861f..2ccbfc27 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/mapper.pyi +++ b/packages/pynumaflow-lite/pynumaflow_lite/mapper.pyi @@ -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: ... diff --git a/packages/pynumaflow-lite/pynumaflow_lite/mapstreamer.pyi b/packages/pynumaflow-lite/pynumaflow_lite/mapstreamer.pyi index 06232864..91bf70bc 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/mapstreamer.pyi +++ b/packages/pynumaflow-lite/pynumaflow_lite/mapstreamer.pyi @@ -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: ... diff --git a/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi b/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi index cb61cde0..4db31b50 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi +++ b/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi @@ -103,11 +103,12 @@ class SinkAsyncServer: *, sock_file: str | None = ..., 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: ... + async def wait_for_termination(self) -> None: ... async def __aenter__(self) -> SinkAsyncServer: ... async def __aexit__( self, diff --git a/packages/pynumaflow-lite/tests/examples/map_cat.py b/packages/pynumaflow-lite/tests/examples/map_cat.py index bb38d770..20f21f1c 100644 --- a/packages/pynumaflow-lite/tests/examples/map_cat.py +++ b/packages/pynumaflow-lite/tests/examples/map_cat.py @@ -1,3 +1,5 @@ +import asyncio + from pynumaflow_lite.mapper import Datum, MapAsyncServer, Message @@ -26,9 +28,13 @@ async def map_handler(datum: Datum) -> list[Message]: return [Message(datum.value, keys=datum.keys, user_metadata=user_metadata)] -if __name__ == "__main__": - MapAsyncServer( +async def main() -> None: + await MapAsyncServer( map_handler, sock_file="/tmp/var/run/numaflow/map.sock", server_info_file="/tmp/var/run/numaflow/mapper-server-info", - ).run() + ).serve() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/packages/pynumaflow-lite/tests/examples/map_cat_class.py b/packages/pynumaflow-lite/tests/examples/map_cat_class.py index fba24163..450e9414 100644 --- a/packages/pynumaflow-lite/tests/examples/map_cat_class.py +++ b/packages/pynumaflow-lite/tests/examples/map_cat_class.py @@ -1,3 +1,5 @@ +import asyncio + from pynumaflow_lite.mapper import Datum, MapAsyncServer, Mapper, Message @@ -27,9 +29,13 @@ async def handler(self, datum: Datum) -> list[Message]: return [Message(datum.value, keys=datum.keys, user_metadata=user_metadata)] -if __name__ == "__main__": - MapAsyncServer( +async def main() -> None: + await MapAsyncServer( SimpleCat(), sock_file="/tmp/var/run/numaflow/map.sock", server_info_file="/tmp/var/run/numaflow/mapper-server-info", - ).run() + ).serve() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/packages/pynumaflow-lite/tests/examples/sink_log.py b/packages/pynumaflow-lite/tests/examples/sink_log.py index c649e28f..346e1334 100644 --- a/packages/pynumaflow-lite/tests/examples/sink_log.py +++ b/packages/pynumaflow-lite/tests/examples/sink_log.py @@ -1,3 +1,4 @@ +import asyncio import collections.abc import logging @@ -37,9 +38,13 @@ async def async_handler( return responses -if __name__ == "__main__": - sinker.SinkAsyncServer( +async def main() -> None: + await sinker.SinkAsyncServer( async_handler, sock_file="/tmp/var/run/numaflow/sink.sock", server_info_file="/tmp/var/run/numaflow/sinker-server-info", - ).run() + ).serve() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/packages/pynumaflow-lite/tests/examples/sink_log_class.py b/packages/pynumaflow-lite/tests/examples/sink_log_class.py index dda93180..436c5d84 100644 --- a/packages/pynumaflow-lite/tests/examples/sink_log_class.py +++ b/packages/pynumaflow-lite/tests/examples/sink_log_class.py @@ -1,3 +1,4 @@ +import asyncio import logging from collections.abc import AsyncIterable @@ -39,9 +40,13 @@ async def handler(self, datums: AsyncIterable[sinker.Datum]) -> list[sinker.Resp return responses -if __name__ == "__main__": - sinker.SinkAsyncServer( +async def main() -> None: + await sinker.SinkAsyncServer( SimpleLogSink(), sock_file="/tmp/var/run/numaflow/sink.sock", server_info_file="/tmp/var/run/numaflow/sinker-server-info", - ).run() + ).serve() + + +if __name__ == "__main__": + asyncio.run(main()) From 43d8f21aeb7214410427b6e25157c969a5d925e0 Mon Sep 17 00:00:00 2001 From: Sreekanth Date: Sun, 27 Sep 2026 19:05:55 +0530 Subject: [PATCH 2/3] Report all errors using exception group Signed-off-by: Sreekanth --- .../pynumaflow_lite/_sink_dtypes.py | 4 +- .../pynumaflow_lite/sinker.pyi | 22 ++-- packages/pynumaflow-lite/src/sink/mod.rs | 116 ++++++++---------- packages/pynumaflow-lite/src/sink/server.rs | 83 +++++-------- packages/pynumaflow-lite/tests/test_sink.py | 22 +++- .../pynumaflow-lite/tests/test_sink_unit.py | 20 ++- 6 files changed, 129 insertions(+), 138 deletions(-) diff --git a/packages/pynumaflow-lite/pynumaflow_lite/_sink_dtypes.py b/packages/pynumaflow-lite/pynumaflow_lite/_sink_dtypes.py index 840dd65f..04e8e724 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/_sink_dtypes.py +++ b/packages/pynumaflow-lite/pynumaflow_lite/_sink_dtypes.py @@ -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 @@ -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. diff --git a/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi b/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi index 4db31b50..1e80398c 100644 --- a/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi +++ b/packages/pynumaflow-lite/pynumaflow_lite/sinker.pyi @@ -4,6 +4,8 @@ import datetime as _dt from collections.abc import AsyncIterator, Awaitable, Callable from types import TracebackType +from ._sink_dtypes import Sinker as Sinker + class NackOptions: """Per-message redelivery options for a nack.""" @@ -38,6 +40,8 @@ class Message: class Response: id: str error: str | None + serve_response: bytes | None + on_success_message: Message | None nack_options: NackOptions | None @staticmethod @@ -56,11 +60,11 @@ class Response: def __eq__(self, other: object) -> bool: ... class Datum: + id: str keys: list[str] value: bytes watermark: _dt.datetime event_time: _dt.datetime - id: str headers: dict[str, str] user_metadata: dict[str, dict[str, bytes]] system_metadata: dict[str, dict[str, bytes]] @@ -68,14 +72,14 @@ class Datum: def __init__( self, *, - keys: list[str] = ..., - value: bytes = ..., - id: str = ..., + id: str, + keys: list[str] | None = ..., + value: bytes | None = ..., event_time: _dt.datetime | None = ..., watermark: _dt.datetime | None = ..., - headers: dict[str, str] = ..., - user_metadata: dict[str, dict[str, bytes]] = ..., - system_metadata: dict[str, dict[str, bytes]] = ..., + headers: dict[str, str] | None = ..., + user_metadata: dict[str, dict[str, bytes]] | None = ..., + system_metadata: dict[str, dict[str, bytes]] | None = ..., ) -> None: ... def __repr__(self) -> str: ... def __str__(self) -> str: ... @@ -92,10 +96,6 @@ class _SinkAsyncServer: def wait_ready(self, timeout: float = ...) -> Awaitable[None]: ... def stop(self) -> None: ... -class Sinker: - def __call__(self, datums: AsyncIterator[Datum]) -> Awaitable[list[Response]]: ... - async def handler(self, datums: AsyncIterator[Datum]) -> list[Response]: ... - class SinkAsyncServer: def __init__( self, diff --git a/packages/pynumaflow-lite/src/sink/mod.rs b/packages/pynumaflow-lite/src/sink/mod.rs index 50d32c98..bb927eac 100644 --- a/packages/pynumaflow-lite/src/sink/mod.rs +++ b/packages/pynumaflow-lite/src/sink/mod.rs @@ -21,24 +21,7 @@ use pyo3::prelude::*; use std::sync::Mutex; use crate::nack::NackOptions; - -fn bytes_literal(value: &[u8]) -> String { - format!("b\"{}\"", String::from_utf8_lossy(value).escape_debug()) -} - -fn metadata_literal(metadata: &HashMap>>) -> String { - let groups: Vec = metadata - .iter() - .map(|(group, kv)| { - let entries: Vec = kv - .iter() - .map(|(key, value)| format!("{:?}: {}", key, bytes_literal(value))) - .collect(); - format!("{:?}: {{{}}}", group, entries.join(", ")) - }) - .collect(); - format!("{{{}}}", groups.join(", ")) -} +use crate::pyrs::{bytes_literal, py_repr}; fn system_metadata_to_hash_map( value: sink::SystemMetadata, @@ -99,17 +82,13 @@ impl Message { } } - fn __repr__(&self) -> String { - format!( + fn __repr__(&self, py: Python<'_>) -> PyResult { + Ok(format!( "Message(value={}, keys={}, user_metadata={})", bytes_literal(&self.value), - self.keys - .as_ref() - .map_or_else(|| "None".to_string(), |keys| format!("{keys:?}")), - self.user_metadata - .as_ref() - .map_or_else(|| "None".to_string(), metadata_literal) - ) + py_repr(py, &self.keys)?, + py_repr(py, &self.user_metadata)?, + )) } } @@ -136,8 +115,12 @@ pub struct Response { pub response_type: ResponseType, #[pyo3(get)] pub error: Option, + /// Payload for the serving store. It is set only for a serve response. + #[pyo3(get)] pub serve_response: Option>, - pub on_success_msg: Option, + /// Message for the OnSuccess sink. It is set only for an on_success response. + #[pyo3(get)] + pub on_success_message: Option, /// Options sent back to the source when nacking the message. #[pyo3(get)] pub nack_options: Option, @@ -154,7 +137,7 @@ impl Response { response_type: ResponseType::Success, error: None, serve_response: None, - on_success_msg: None, + on_success_message: None, nack_options: None, } } @@ -168,7 +151,7 @@ impl Response { response_type: ResponseType::Failure, error: Some(error), serve_response: None, - on_success_msg: None, + on_success_message: None, nack_options: None, } } @@ -182,7 +165,7 @@ impl Response { response_type: ResponseType::Fallback, error: None, serve_response: None, - on_success_msg: None, + on_success_message: None, nack_options: None, } } @@ -196,7 +179,7 @@ impl Response { response_type: ResponseType::Serve, error: None, serve_response: Some(payload), - on_success_msg: None, + on_success_message: None, nack_options: None, } } @@ -211,7 +194,7 @@ impl Response { response_type: ResponseType::OnSuccess, error: None, serve_response: None, - on_success_msg: message, + on_success_message: message, nack_options: None, } } @@ -225,37 +208,38 @@ impl Response { response_type: ResponseType::Nack, error: None, serve_response: None, - on_success_msg: None, + on_success_message: None, nack_options, } } - fn __repr__(&self) -> String { - match self.response_type { - ResponseType::Success => format!("Response.success(id={:?})", self.id), + fn __repr__(&self, py: Python<'_>) -> PyResult { + let id = py_repr(py, &self.id)?; + Ok(match self.response_type { + ResponseType::Success => format!("Response.success(id={id})"), ResponseType::Failure => format!( - "Response.failure(id={:?}, error={:?})", - self.id, - self.error.as_deref().unwrap_or_default() + "Response.failure(id={id}, error={})", + py_repr(py, self.error.as_deref().unwrap_or_default())? ), - ResponseType::Fallback => format!("Response.fallback(id={:?})", self.id), + ResponseType::Fallback => format!("Response.fallback(id={id})"), ResponseType::Serve => format!( - "Response.serve(id={:?}, payload={})", - self.id, + "Response.serve(id={id}, payload={})", bytes_literal(self.serve_response.as_deref().unwrap_or_default()) ), ResponseType::OnSuccess => format!( - "Response.on_success(id={:?}, message={})", - self.id, - self.on_success_msg - .as_ref() - .map_or_else(|| "None".to_string(), |m| m.__repr__()) + "Response.on_success(id={id}, message={})", + match &self.on_success_message { + Some(message) => message.__repr__(py)?, + None => "None".to_string(), + } ), ResponseType::Nack => format!( - "Response.nack(id={:?}, nack_options={:?})", - self.id, self.nack_options + "Response.nack(id={id}, nack_options={})", + self.nack_options + .as_ref() + .map_or_else(|| "None".to_string(), NackOptions::__repr__) ), - } + }) } } @@ -286,7 +270,7 @@ impl From for sink::Response { response_type, err: value.error, serve_response: value.serve_response, - on_success_msg: value.on_success_msg.map(|m| m.into()), + on_success_msg: value.on_success_message.map(|m| m.into()), nack_options: value.nack_options.map(Into::into), } } @@ -326,9 +310,9 @@ impl Datum { #[new] #[pyo3(signature = ( *, + id: "str", keys: "list[str] | None"=None, value: "bytes | None"=None, - id: "str | None"=None, event_time: "datetime.datetime | None"=None, watermark: "datetime.datetime | None"=None, headers: "dict[str, str] | None"=None, @@ -337,9 +321,9 @@ impl Datum { ) -> "Datum")] #[allow(clippy::too_many_arguments)] fn new( + id: String, keys: Option>, value: Option>, - id: Option, event_time: Option>, watermark: Option>, headers: Option>, @@ -351,29 +335,29 @@ impl Datum { value: value.unwrap_or_default(), watermark: watermark.unwrap_or(DateTime::::UNIX_EPOCH), event_time: event_time.unwrap_or(DateTime::::UNIX_EPOCH), - id: id.unwrap_or_default(), + id, headers: headers.unwrap_or_default(), user_metadata: user_metadata.unwrap_or_default(), system_metadata: system_metadata.unwrap_or_default(), } } - fn __repr__(&self) -> String { - format!( - "Datum(keys={:?}, value={}, watermark={}, event_time={}, id={:?}, headers={:?}, user_metadata={}, system_metadata={})", - self.keys, + fn __repr__(&self, py: Python<'_>) -> PyResult { + Ok(format!( + "Datum(id={}, keys={}, value={}, watermark={}, event_time={}, headers={}, user_metadata={}, system_metadata={})", + py_repr(py, &self.id)?, + py_repr(py, &self.keys)?, bytes_literal(&self.value), self.watermark, self.event_time, - self.id, - self.headers, - metadata_literal(&self.user_metadata), - metadata_literal(&self.system_metadata) - ) + py_repr(py, &self.headers)?, + py_repr(py, &self.user_metadata)?, + py_repr(py, &self.system_metadata)?, + )) } - fn __str__(&self) -> String { - self.__repr__() + fn __str__(&self, py: Python<'_>) -> PyResult { + self.__repr__(py) } } diff --git a/packages/pynumaflow-lite/src/sink/server.rs b/packages/pynumaflow-lite/src/sink/server.rs index f588bdf7..69d80077 100644 --- a/packages/pynumaflow-lite/src/sink/server.rs +++ b/packages/pynumaflow-lite/src/sink/server.rs @@ -1,38 +1,29 @@ +use std::sync::{Arc, Mutex}; + use numaflow::shared::ServerExtras; use numaflow::sink; - use pyo3::exceptions::PyTypeError; use pyo3::prelude::*; -use std::sync::{Arc, Mutex}; -use tokio::task::JoinHandle; + +use crate::pyrs::{combine_errors, format_error}; pub(crate) struct PySinkRunner { pub(crate) event_loop: Arc>, pub(crate) py_func: Arc>, - pub(crate) error_slot: Arc>>, - pub(crate) shutdown_tx: Arc>>>, + pub(crate) errors: Arc>>, } impl PySinkRunner { - fn fail(&self, error: PyErr) { - Python::attach(|py| error.print(py)); - - let mut error_slot = self.error_slot.lock().unwrap(); - if error_slot.is_none() { - *error_slot = Some(error); - } - drop(error_slot); - - if let Some(tx) = self.shutdown_tx.lock().unwrap().take() { - let _ = tx.send(()); - } - } - - async fn fail_sink(&self, error: PyErr, forwarder: JoinHandle<()>) -> Vec { - self.fail(error); - forwarder.abort(); - let _ = forwarder.await; - Vec::new() + fn fail(&self, error: PyErr) -> ! { + // Keep every error. start() raises them together, which lets Python + // format every traceback instead of Rust printing them by hand. + let message = Python::attach(|py| format_error(py, &error)); + self.errors.lock().unwrap().push(error); + + // numaflow catches this panic, sends a gRPC error for this batch, and + // starts the server shutdown. An empty result would instead look like a + // batch that the handler wrote with no responses. + panic!("{message}"); } } @@ -55,7 +46,7 @@ impl sink::Sinker for PySinkRunner { // When input ends, dropping tx closes the channel }); - // Call the Python coroutine: py_func(datums: AsyncIterable[Datum]) -> list[Response] + // Call the Python coroutine: py_func(datums: AsyncIterator[Datum]) -> list[Response] let fut = match Python::attach(|py| -> PyResult<_> { let locals = pyo3_async_runtimes::TaskLocals::new(self.event_loop.bind(py).clone()); let py_func = self.py_func.clone(); @@ -76,12 +67,12 @@ impl sink::Sinker for PySinkRunner { }) }) { Ok(fut) => fut, - Err(error) => return self.fail_sink(error, forwarder).await, + Err(error) => self.fail(error), }; let result = match fut.await { Ok(result) => result, - Err(error) => return self.fail_sink(error, forwarder).await, + Err(error) => self.fail(error), }; // Ensure forwarder completes @@ -101,10 +92,7 @@ impl sink::Sinker for PySinkRunner { }) }) { Ok(responses) => responses, - Err(error) => { - self.fail(error); - return Vec::new(); - } + Err(error) => self.fail(error), }; responses @@ -131,25 +119,15 @@ pub(super) async fn start( }); let event_loop = rx.await.unwrap(); - let error_slot = Arc::new(Mutex::new(None)); - let (internal_shutdown_tx, internal_shutdown_rx) = tokio::sync::oneshot::channel(); - let (server_shutdown_tx, server_shutdown_rx) = tokio::sync::oneshot::channel(); - - tokio::spawn(async move { - tokio::select! { - _ = shutdown_rx => {}, - _ = internal_shutdown_rx => {}, - } - let _ = server_shutdown_tx.send(()); - }); - - let (sig_handle, combined_rx) = crate::pyrs::setup_sig_handler(server_shutdown_rx); + let errors = Arc::new(Mutex::new(Vec::new())); + // Shutdown has two sources, and neither one needs a channel here. The Python + // side signals stop() through shutdown_rx. An uncaught Python error panics in + // fail(), and numaflow then shuts the server down on its own. let py_runner = PySinkRunner { py_func: Arc::new(py_func), event_loop: event_loop.clone(), - error_slot: error_slot.clone(), - shutdown_tx: Arc::new(Mutex::new(Some(internal_shutdown_tx))), + errors: Arc::clone(&errors), }; let server = numaflow::sink::Server::new(py_runner) @@ -157,7 +135,7 @@ pub(super) async fn start( .with_server_info_file(info_file); let result = server - .start_with_shutdown(combined_rx) + .start_with_shutdown(shutdown_rx) .await .map_err(|e| pyo3::PyErr::new::(e.to_string())); @@ -173,14 +151,9 @@ pub(super) async fn start( // Wait for the blocking asyncio thread to finish. let _ = py_asyncio_loop_handle.await; - // if not finished, abort it - if !sig_handle.is_finished() { - println!("Aborting signal handler"); - sig_handle.abort(); - } - - if let Some(error) = error_slot.lock().unwrap().take() { - return Err(error); + let errors = std::mem::take(&mut *errors.lock().unwrap()); + if !errors.is_empty() { + return Err(Python::attach(|py| combine_errors(py, errors))); } result diff --git a/packages/pynumaflow-lite/tests/test_sink.py b/packages/pynumaflow-lite/tests/test_sink.py index bef59904..ebd6f806 100644 --- a/packages/pynumaflow-lite/tests/test_sink.py +++ b/packages/pynumaflow-lite/tests/test_sink.py @@ -32,7 +32,10 @@ def test_sink_data_types_match_pythonic_api(): assert datum.event_time assert datum.user_metadata == {} assert isinstance(datum.system_metadata, dict) - assert 'value=b"value"' in repr(datum) + assert "value=b'value'" in repr(datum) + assert repr(sinker.Datum(id="id-2", value=b"\xff")).startswith("Datum(id='id-2', keys=[], value=b'\\xff'") + with pytest.raises(TypeError): + sinker.Datum(value=b"value") # id is required message = sinker.Message( b"value", @@ -54,6 +57,23 @@ def test_sink_data_types_match_pythonic_api(): assert response.error is None assert response == sinker.Response.success("id-1") assert sinker.Response.failure("id-1", "boom").error == "boom" + assert response.serve_response is None + assert response.on_success_message is None + + serve = sinker.Response.serve("id-1", b"payload") + assert serve.serve_response == b"payload" + assert repr(serve) == "Response.serve(id='id-1', payload=b'payload')" + + on_success = sinker.Response.on_success("id-1", message) + assert on_success.on_success_message == message + assert repr(on_success) == ( + "Response.on_success(id='id-1', message=Message(value=b'value', keys=['key'], " + "user_metadata={'custom_info': {'version': b'1.0.0'}}))" + ) + + nack = sinker.Response.nack("id-1", sinker.NackOptions(delay=5)) + assert repr(nack) == f"Response.nack(id='id-1', nack_options={nack.nack_options!r})" + assert "Some(" not in repr(nack) assert not hasattr(sinker, "KeyValueGroup") assert not hasattr(sinker, "Responses") assert not hasattr(sinker, "PyAsyncDatumStream") diff --git a/packages/pynumaflow-lite/tests/test_sink_unit.py b/packages/pynumaflow-lite/tests/test_sink_unit.py index c6f281fd..7805d9c5 100644 --- a/packages/pynumaflow-lite/tests/test_sink_unit.py +++ b/packages/pynumaflow-lite/tests/test_sink_unit.py @@ -49,7 +49,11 @@ async def _run_sink_client(sock_path: Path) -> subprocess.CompletedProcess[str]: ) -async def _exercise_server(tmp_path: Path, handler) -> None: +async def _exercise_server( + tmp_path: Path, + handler, + clients: list[subprocess.CompletedProcess[str]] | None = None, +) -> None: path_id = f"{tmp_path.name[:12]}-{uuid.uuid4().hex[:12]}" sock_path = Path(f"/tmp/pnl-{path_id}.sock") server_info_path = Path(f"/tmp/pnl-{path_id}.info") @@ -61,7 +65,9 @@ async def _exercise_server(tmp_path: Path, handler) -> None: try: async with server: - await _run_sink_client(sock_path) + client = await _run_sink_client(sock_path) + if clients is not None: + clients.append(client) finally: for path in (sock_path, server_info_path): with suppress(FileNotFoundError): @@ -74,8 +80,16 @@ async def handler(datums: AsyncIterator[sinker.Datum]) -> list[sinker.Response]: raise RuntimeError("sink exploded") return [] + clients: list[subprocess.CompletedProcess[str]] = [] with pytest.raises(RuntimeError, match="sink exploded"): - asyncio.run(_exercise_server(tmp_path, handler)) + asyncio.run(_exercise_server(tmp_path, handler, clients)) + + # The client must get a gRPC error that carries the handler error, not an + # empty batch result. + [client] = clients + assert client.returncode != 0 + assert "sink exploded" in client.stderr + assert "results.len()" not in client.stderr def test_sink_server_rejects_non_list_response(tmp_path: Path): From 015f9e93aad7a2983ff507526120c5a4ddcecf13 Mon Sep 17 00:00:00 2001 From: Sreekanth Date: Sun, 27 Sep 2026 19:15:08 +0530 Subject: [PATCH 3/3] Update sinker example Signed-off-by: Sreekanth --- .../manifests/sink/.dockerignore | 1 + .../pynumaflow-lite/manifests/sink/Dockerfile | 77 ++++++++++--------- .../manifests/sink/pyproject.toml | 8 +- .../pynumaflow-lite/manifests/sink/uv.lock | 8 ++ 4 files changed, 51 insertions(+), 43 deletions(-) create mode 100644 packages/pynumaflow-lite/manifests/sink/.dockerignore create mode 100644 packages/pynumaflow-lite/manifests/sink/uv.lock diff --git a/packages/pynumaflow-lite/manifests/sink/.dockerignore b/packages/pynumaflow-lite/manifests/sink/.dockerignore new file mode 100644 index 00000000..21d0b898 --- /dev/null +++ b/packages/pynumaflow-lite/manifests/sink/.dockerignore @@ -0,0 +1 @@ +.venv/ diff --git a/packages/pynumaflow-lite/manifests/sink/Dockerfile b/packages/pynumaflow-lite/manifests/sink/Dockerfile index 90ba07c7..9e7b9f6a 100644 --- a/packages/pynumaflow-lite/manifests/sink/Dockerfile +++ b/packages/pynumaflow-lite/manifests/sink/Dockerfile @@ -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"] diff --git a/packages/pynumaflow-lite/manifests/sink/pyproject.toml b/packages/pynumaflow-lite/manifests/sink/pyproject.toml index 8e475beb..226b8e38 100644 --- a/packages/pynumaflow-lite/manifests/sink/pyproject.toml +++ b/packages/pynumaflow-lite/manifests/sink/pyproject.toml @@ -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" - diff --git a/packages/pynumaflow-lite/manifests/sink/uv.lock b/packages/pynumaflow-lite/manifests/sink/uv.lock new file mode 100644 index 00000000..d0e09b73 --- /dev/null +++ b/packages/pynumaflow-lite/manifests/sink/uv.lock @@ -0,0 +1,8 @@ +version = 1 +revision = 3 +requires-python = "==3.11.*" + +[[package]] +name = "sink-log" +version = "0.1.0" +source = { virtual = "." }