35 lines
1004 B
Python
Raw Normal View History

2025-10-15 20:53:52 +02:00
import logging
from typing import AsyncIterable, AsyncIterator, List, Union
2025-10-20 15:42:12 +02:00
from core.models import NbeModel, NbeSchema
2025-10-15 20:53:52 +02:00
2025-10-20 15:42:12 +02:00
T = Union[NbeModel, NbeSchema]
Data = Union[T, List[T]]
2025-10-15 20:53:52 +02:00
Stream = AsyncIterator[Data]
logger = logging.getLogger(__name__)
def _into_ndjson_data(data: Data) -> bytes:
if isinstance(data, list):
return b"".join(item.model_dump_ndjson() for item in data)
else:
return data.model_dump_ndjson()
2025-10-30 11:48:34 +01:00
async def into_ndjson_stream(stream: Stream, *, bootstrap_data: Data = None) -> AsyncIterable[bytes]:
2025-10-15 20:53:52 +02:00
if bootstrap_data is not None:
ndjson_data = _into_ndjson_data(bootstrap_data)
if ndjson_data:
yield ndjson_data
else:
logger.debug("Ignoring streaming bootstrap data because it is empty.")
async for data in stream:
ndjson_data = _into_ndjson_data(data)
if ndjson_data:
yield ndjson_data
else:
logger.debug("Ignoring streaming data because it is empty.")