Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -362,19 +362,30 @@ def get_r_point_data(r_point_data: str) -> list[dict[str, Any]]:
return r_data_list


def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
def decode_data_streaming(named_file_contents: NamedFileContents) -> Any: # Generator
"""
Decodes the proprietary file into a structured dict
Decodes the proprietary file into a structured dict, streaming cycle data.
:param named_file_contents: The named file contents containing the input file details
:return: structured dictionary of decoded data
:return: Iterator yielding metadata first, then cycle data dicts
"""
intermediate_json: dict[str, Any] = {}
content = ole.OleFileIO(named_file_contents.get_bytes_stream())
streams = content.listdir()
sensorgram_df_list = []
report_point_data = {}

# Group streams by cycle number
cycle_streams: dict[int, list[Any]] = {}
report_point_df = None

for stream in streams:
# Check if this is a cycle-specific stream first
if cycle_match := cycle_pattern.search(stream[0]):
cycle_number = int(cycle_match.group(1))
if cycle_number not in cycle_streams:
cycle_streams[cycle_number] = []
cycle_streams[cycle_number].append(stream)
continue

# Process metadata streams
stream_content = content.openstream(stream).read()
if stream == ["Environment"]:
intermediate_json["system_information"] = extract_stream_data(
Expand All @@ -384,7 +395,6 @@ def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
intermediate_json["chip"] = extract_stream_data(
stream_content.decode("utf-8")
)

elif stream == ["AppData", "Dip"]:
intermediate_json["dip"] = get_dip_data(stream_content.decode("utf-8"))
elif stream == ["AppData", "ApplicationTemplate"]:
Expand All @@ -396,76 +406,11 @@ def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
intermediate_json["sample_data"] = sample_data
elif stream == ["RPoint Table"]:
r_point_data = get_r_point_data(stream_content.decode("utf-8"))
r_point_dataframe = pd.DataFrame(r_point_data)
report_point_data = {
group["Cycle"].iloc[0]: group
for _, group in r_point_dataframe.groupby("Cycle")
}
report_point_df = pd.DataFrame(r_point_data)

if not (cycle_match := cycle_pattern.search(stream[0])):
continue
cycle_number = int(cycle_match.group(1))

if (curve_match := curve_pattern.search(str(stream))) and (
window_match := window_pattern.search(str(stream))
):
curve_number = curve_match.group(1)
window_number = window_match.group(1)

# get the flow cell number
if "Labels" in str(stream):
label_list = []
if len(stream_content) != 0:
decoded_stream_content = stream_content.decode("utf-8")
for line in decoded_stream_content.strip().split("\n"):
label_list.append(line)
if "Fc" in line:
flow_cell = label_list[0].split("=")[1]

# gets the xy data as data frame
elif "XYData" in str(stream):
xy_data_list = list(
struct.unpack("f" * (len(stream_content) // 4), stream_content)
)
indexed_xy_data_list = xy_data_list[3:]
length_of_xydata = len(xy_data_list[3:])
xy_result_list = indexed_xy_data_list[int(length_of_xydata / 2) :]
time_values = indexed_xy_data_list[: int(length_of_xydata / 2)]
length_of_xy = len(xy_result_list)

xy_data_frame = pd.DataFrame(
{
"Flow Cell Number": [flow_cell] * length_of_xy,
"Cycle Number": [cycle_number] * length_of_xy,
"Curve Number": [curve_number] * length_of_xy,
"Window Number": [window_number] * length_of_xy,
"Sensorgram (RU)": xy_result_list,
"Time (s)": time_values,
}
)

sensorgram_df_list.append(xy_data_frame)

# gets the segment data
elif "Segment" in str(stream):
segment_data_list = list(
struct.unpack("f" * (len(stream_content) // 4), stream_content)
)
length_of_segment = len(segment_data_list[11:])
segment_data_frame = pd.DataFrame(
{
"Flow Cell Number": [flow_cell] * length_of_segment,
"Cycle Number": [cycle_number] * length_of_segment,
"Curve Number": [curve_number] * length_of_segment,
"Window Number": [window_number] * length_of_segment,
"Sensorgram (RU)": segment_data_list[11:],
}
)
sensorgram_df_list.append(segment_data_frame)

combined_sensorgram_df = pd.concat(sensorgram_df_list, ignore_index=True)

sensorgram_data = {}
# Get total cycles and data collection rate
total_cycles = max(cycle_streams.keys()) if cycle_streams else 0
intermediate_json["total_cycles"] = total_cycles

data_collection_rate = None
if intermediate_json.get("application_template_details", {}).get(
Expand All @@ -475,10 +420,82 @@ def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
"DataCollectionRate"
]["value"]

grouped_df = combined_sensorgram_df.groupby("Cycle Number")
for cycle_num, cycle_group in grouped_df:
group = cycle_group.copy()
# Yield metadata first
yield intermediate_json

# Process cycles one at a time
for cycle_num in sorted(cycle_streams.keys()):
sensorgram_df_list = []
flow_cell = "1" # default

for stream in cycle_streams[cycle_num]:
stream_content = content.openstream(stream).read()

if (curve_match := curve_pattern.search(str(stream))) and (
window_match := window_pattern.search(str(stream))
):
curve_number = curve_match.group(1)
window_number = window_match.group(1)

# get the flow cell number
if "Labels" in str(stream):
label_list = []
if len(stream_content) != 0:
decoded_stream_content = stream_content.decode("utf-8")
for line in decoded_stream_content.strip().split("\n"):
label_list.append(line)
if "Fc" in line:
flow_cell = label_list[0].split("=")[1]

# gets the xy data as data frame
elif "XYData" in str(stream):
xy_data_list = list(
struct.unpack("f" * (len(stream_content) // 4), stream_content)
)
indexed_xy_data_list = xy_data_list[3:]
length_of_xydata = len(xy_data_list[3:])
xy_result_list = indexed_xy_data_list[int(length_of_xydata / 2) :]
time_values = indexed_xy_data_list[: int(length_of_xydata / 2)]
length_of_xy = len(xy_result_list)

xy_data_frame = pd.DataFrame(
{
"Flow Cell Number": [flow_cell] * length_of_xy,
"Cycle Number": [cycle_num] * length_of_xy,
"Curve Number": [curve_number] * length_of_xy,
"Window Number": [window_number] * length_of_xy,
"Sensorgram (RU)": xy_result_list,
"Time (s)": time_values,
}
)

sensorgram_df_list.append(xy_data_frame)

# gets the segment data
elif "Segment" in str(stream):
segment_data_list = list(
struct.unpack("f" * (len(stream_content) // 4), stream_content)
)
length_of_segment = len(segment_data_list[11:])
segment_data_frame = pd.DataFrame(
{
"Flow Cell Number": [flow_cell] * length_of_segment,
"Cycle Number": [cycle_num] * length_of_segment,
"Curve Number": [curve_number] * length_of_segment,
"Window Number": [window_number] * length_of_segment,
"Sensorgram (RU)": segment_data_list[11:],
}
)
sensorgram_df_list.append(segment_data_frame)

# Process this cycle's sensorgram data
if not sensorgram_df_list:
continue

combined_sensorgram_df = pd.concat(sensorgram_df_list, ignore_index=True)
group = combined_sensorgram_df.copy()

# Time normalization logic
if "Time (s)" in group.columns:
max_flow_cell = group["Flow Cell Number"].max()
max_flow_cell_mask = group["Flow Cell Number"] == max_flow_cell
Expand All @@ -495,9 +512,9 @@ def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
)
reference_times = max_flow_cell_data["Time (s)"].values

for flow_cell in unique_flow_cells:
if flow_cell != max_flow_cell:
flow_cell_mask = group["Flow Cell Number"] == flow_cell
for fc in unique_flow_cells:
if fc != max_flow_cell:
flow_cell_mask = group["Flow Cell Number"] == fc
flow_cell_indices = group[flow_cell_mask].index

num_points = len(flow_cell_indices)
Expand All @@ -516,7 +533,7 @@ def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
else:
group["Time (s)"] = group.groupby("Flow Cell Number").cumcount() + 1

sensorgram_data[str(cycle_num)] = group[
sensorgram_data = group[
[
"Flow Cell Number",
"Cycle Number",
Expand All @@ -527,16 +544,33 @@ def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
]
]

# In some cases the "RPoint Table" stream may not have data for the first cycle,
# but we get streams with sensorgram data (XYData or Segment streams) for all the cycles.
intermediate_json["cycle_data"] = [
{
"cycle_number": cycle,
"report_point_data": report_point_data.get(cycle),
"sensorgram_data": sensorgram_df,
# Get report point data for this cycle
report_point_data = None
if report_point_df is not None:
cycle_report = report_point_df[report_point_df["Cycle"] == str(cycle_num)]
if not cycle_report.empty:
report_point_data = cycle_report

# Yield this cycle's data
yield {
"cycle_number": str(cycle_num),
"report_point_data": report_point_data,
"sensorgram_data": sensorgram_data,
}
for cycle, sensorgram_df in sensorgram_data.items()
]
intermediate_json["total_cycles"] = cycle_number


def decode_data(named_file_contents: NamedFileContents) -> dict[str, Any]:
"""
Decodes the proprietary file into a structured dict
:param named_file_contents: The named file contents containing the input file details
:return: structured dictionary of decoded data
"""
# Use streaming decoder but collect all data
stream_gen = decode_data_streaming(named_file_contents)
intermediate_json: dict[str, Any] = next(stream_gen) # Get metadata

# Collect all cycle data
cycle_data = list(stream_gen)
intermediate_json["cycle_data"] = cycle_data

return intermediate_json
Loading