Skip to content

Commit e5909ba

Browse files
committed
Handle unparsed streaming artifact metadata
1 parent 000ab4e commit e5909ba

5 files changed

Lines changed: 115 additions & 9 deletions

File tree

src/polystore/__init__.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,10 +26,10 @@
2626
get_backend,
2727
)
2828
from .constants import Backend, MemoryType, TransportMode
29-
from .disk import DiskStorageBackend
29+
from .disk import DiskBackend, DiskStorageBackend
3030
from .filemanager import FileManager
3131
from .formats import FileFormat, DEFAULT_IMAGE_EXTENSIONS
32-
from .memory import MemoryStorageBackend
32+
from .memory import MemoryBackend, MemoryStorageBackend
3333
from .metadata_writer import (
3434
AtomicMetadataWriter,
3535
MetadataWriteError,
@@ -76,7 +76,9 @@
7676
"register_cleanup_callback",
7777
"STORAGE_BACKENDS",
7878
"DiskStorageBackend",
79+
"DiskBackend",
7980
"MemoryStorageBackend",
81+
"MemoryBackend",
8082
"FileManager",
8183
"file_lock",
8284
"atomic_write_json",

src/polystore/disk.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -834,3 +834,6 @@ def _save_rois(self, rois: List, output_path: Path, images_dir: str = None, **kw
834834

835835
logger.info(f"Saved {roi_count} ROIs to .roi.zip archive: {output_path}")
836836
return str(output_path)
837+
838+
839+
DiskBackend = DiskStorageBackend

src/polystore/memory.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -657,3 +657,6 @@ def __init__(self, target: str):
657657

658658
def __repr__(self):
659659
return f"<MemorySymlink → {self.target}>"
660+
661+
662+
MemoryBackend = MemoryStorageBackend

src/polystore/streaming/_streaming_backend.py

Lines changed: 27 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -9,8 +9,9 @@
99
import os
1010
import time
1111
import uuid
12+
from dataclasses import dataclass
1213
from pathlib import Path
13-
from typing import Any, Callable, List, Set, Union
14+
from typing import Any, Callable, List, Mapping, Set, Union
1415
import numpy as np
1516

1617
from ..base import DataSink
@@ -24,6 +25,27 @@
2425
logger = logging.getLogger(__name__)
2526

2627

28+
@dataclass(frozen=True)
29+
class StreamingComponentMetadata:
30+
"""Message metadata for one streamed item."""
31+
32+
parsed_filename_metadata: Mapping[str, Any] | None
33+
source: str
34+
35+
def to_payload(self) -> dict[str, Any]:
36+
if self.parsed_filename_metadata is None:
37+
metadata: dict[str, Any] = {}
38+
elif isinstance(self.parsed_filename_metadata, Mapping):
39+
metadata = dict(self.parsed_filename_metadata)
40+
else:
41+
raise TypeError(
42+
"Streaming filename parser must return a mapping or None, "
43+
f"got {type(self.parsed_filename_metadata).__name__}."
44+
)
45+
metadata["source"] = self.source
46+
return metadata
47+
48+
2749
class StreamingBackend(DataSink):
2850
"""
2951
Abstract base class for ZeroMQ-based streaming backends.
@@ -165,12 +187,10 @@ def _parse_component_metadata(self, file_path: Union[str, Path], microscope_hand
165187
Component metadata dict with source added
166188
"""
167189
filename = os.path.basename(str(file_path))
168-
component_metadata = microscope_handler.parser.parse_filename(filename)
169-
170-
# Add pre-built source value directly
171-
component_metadata['source'] = source
172-
173-
return component_metadata
190+
return StreamingComponentMetadata(
191+
microscope_handler.parser.parse_filename(filename),
192+
source,
193+
).to_payload()
174194

175195
def _detect_data_type(self, data: Any):
176196
"""

tests/test_streaming_metadata.py

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
from types import SimpleNamespace
2+
3+
import pytest
4+
5+
from polystore.streaming._streaming_backend import StreamingBackend
6+
7+
8+
class MetadataProbeStreamingBackend(StreamingBackend):
9+
VIEWER_TYPE = "probe"
10+
SHM_PREFIX = "probe_"
11+
12+
def save_batch(self, data_list, file_paths, **kwargs):
13+
raise NotImplementedError
14+
15+
16+
def test_streaming_component_metadata_accepts_unparsed_artifact_filename() -> None:
17+
backend = MetadataProbeStreamingBackend()
18+
microscope_handler = SimpleNamespace(
19+
parser=SimpleNamespace(parse_filename=lambda _filename: None)
20+
)
21+
22+
metadata = backend._parse_component_metadata(
23+
"A01_s001_w1_z001_t001_Nuclei_step3_rois.roi.zip",
24+
microscope_handler,
25+
source="IdentifyPrimaryObjects",
26+
)
27+
28+
assert metadata == {"source": "IdentifyPrimaryObjects"}
29+
30+
31+
def test_streaming_batch_items_accept_unparsed_artifact_filename() -> None:
32+
backend = MetadataProbeStreamingBackend()
33+
microscope_handler = SimpleNamespace(
34+
parser=SimpleNamespace(parse_filename=lambda _filename: None)
35+
)
36+
37+
batch_images, image_ids = backend._prepare_batch_items(
38+
[object()],
39+
["A01_s001_w1_z001_t001_Nuclei_step3_rois.roi.zip"],
40+
microscope_handler,
41+
"IdentifyPrimaryObjects",
42+
lambda _data, _path, _data_type: ({"payload": "ok"}, "image"),
43+
)
44+
45+
assert len(image_ids) == 1
46+
assert batch_images[0]["metadata"] == {"source": "IdentifyPrimaryObjects"}
47+
assert batch_images[0]["payload"] == "ok"
48+
49+
50+
def test_streaming_component_metadata_preserves_parsed_filename_fields() -> None:
51+
backend = MetadataProbeStreamingBackend()
52+
microscope_handler = SimpleNamespace(
53+
parser=SimpleNamespace(
54+
parse_filename=lambda _filename: {"well": "A01", "channel": 1}
55+
)
56+
)
57+
58+
metadata = backend._parse_component_metadata(
59+
"A01_s001_w1_z001_t001.TIF",
60+
microscope_handler,
61+
source="Crop",
62+
)
63+
64+
assert metadata == {"well": "A01", "channel": 1, "source": "Crop"}
65+
66+
67+
def test_streaming_component_metadata_rejects_invalid_parser_result() -> None:
68+
backend = MetadataProbeStreamingBackend()
69+
microscope_handler = SimpleNamespace(
70+
parser=SimpleNamespace(parse_filename=lambda _filename: ["not", "metadata"])
71+
)
72+
73+
with pytest.raises(TypeError, match="mapping or None"):
74+
backend._parse_component_metadata(
75+
"A01_s001_w1_z001_t001.TIF",
76+
microscope_handler,
77+
source="Crop",
78+
)

0 commit comments

Comments
 (0)