def get_message( # type: ignore[override]
self,
channel: str,
body: Any,
sub: "_StreamHandlerMixin",
) -> Any:
message: BatchStreamMessage | DefaultStreamMessage
if sub.stream_sub.batch:
message = BatchStreamMessage(
type="bstream",
channel=channel,
data=[{bDATA_KEY: body}],
message_ids=[],
)
else:
message = DefaultStreamMessage(
type="stream",
channel=channel,
data={bDATA_KEY: body},
message_ids=[],
)
if sub.stream_sub.claim_min_idle_time is not None:
# The in-memory broker has no claiming: every delivery is a new
# message, so expose the new-message metadata values
message["idle_times"] = [0]
message["delivery_counts"] = [0]
return message