Skip to content

Commit c3cdcf7

Browse files
committed
reformate with black
1 parent 5f53124 commit c3cdcf7

13 files changed

Lines changed: 366 additions & 202 deletions

connectors/dnanexus.py

Lines changed: 51 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -22,13 +22,15 @@
2222
# Shared data structures
2323
# ---------------------------------------------------------------------------
2424

25+
2526
@dataclass
2627
class RemoteFile:
2728
"""Represents a file on a remote platform before transfer."""
29+
2830
platform: str
2931
file_id: str
3032
file_name: str
31-
file_type: str # FASTQ, BAM, VCF, etc.
33+
file_type: str # FASTQ, BAM, VCF, etc.
3234
file_size_bytes: int
3335
source_path: str
3436
sample_id: str
@@ -39,6 +41,7 @@ class RemoteFile:
3941
# Base connector interface
4042
# ---------------------------------------------------------------------------
4143

44+
4245
class BaseConnector(ABC):
4346
"""
4447
Abstract base class for all platform connectors.
@@ -51,7 +54,9 @@ def list_files(self, project_id: str, file_type: str = None) -> list[RemoteFile]
5154
...
5255

5356
@abstractmethod
54-
def stream_to_s3(self, remote_file: RemoteFile, s3_bucket: str, s3_key: str) -> dict:
57+
def stream_to_s3(
58+
self, remote_file: RemoteFile, s3_bucket: str, s3_key: str
59+
) -> dict:
5560
"""
5661
Stream a file directly from the platform to S3 without
5762
materializing the full file in memory.
@@ -69,6 +74,7 @@ def get_metadata(self, file_id: str) -> dict:
6974
# DNAnexus connector
7075
# ---------------------------------------------------------------------------
7176

77+
7278
class DNAnexusConnector(BaseConnector):
7379
"""
7480
Connector for DNAnexus platform.
@@ -77,9 +83,12 @@ class DNAnexusConnector(BaseConnector):
7783
"""
7884

7985
SUPPORTED_EXTENSIONS = {
80-
".fastq": "FASTQ", ".fastq.gz": "FASTQ",
81-
".bam": "BAM", ".bam.bai": "BAM",
82-
".vcf": "VCF", ".vcf.gz": "VCF",
86+
".fastq": "FASTQ",
87+
".fastq.gz": "FASTQ",
88+
".bam": "BAM",
89+
".bam.bai": "BAM",
90+
".vcf": "VCF",
91+
".vcf.gz": "VCF",
8392
".cram": "CRAM",
8493
}
8594

@@ -119,23 +128,25 @@ def list_files(self, project_id: str, file_type: str = None) -> list[RemoteFile]
119128
props = desc.get("properties", {})
120129
sample_id = props.get("sample_id", props.get("sample", "unknown"))
121130

122-
results.append(RemoteFile(
123-
platform="DNAnexus",
124-
file_id=item["id"],
125-
file_name=fname,
126-
file_type=inferred_type,
127-
file_size_bytes=desc.get("size", 0),
128-
source_path=f"{project_id}:{desc.get('folder','/')}/{fname}",
129-
sample_id=sample_id,
130-
metadata={
131-
"dx_project": project_id,
132-
"dx_folder": desc.get("folder", "/"),
133-
"dx_tags": desc.get("tags", []),
134-
"dx_properties": props,
135-
"dx_created": desc.get("created"),
136-
"dx_modified": desc.get("modified"),
137-
}
138-
))
131+
results.append(
132+
RemoteFile(
133+
platform="DNAnexus",
134+
file_id=item["id"],
135+
file_name=fname,
136+
file_type=inferred_type,
137+
file_size_bytes=desc.get("size", 0),
138+
source_path=f"{project_id}:{desc.get('folder','/')}/{fname}",
139+
sample_id=sample_id,
140+
metadata={
141+
"dx_project": project_id,
142+
"dx_folder": desc.get("folder", "/"),
143+
"dx_tags": desc.get("tags", []),
144+
"dx_properties": props,
145+
"dx_created": desc.get("created"),
146+
"dx_modified": desc.get("modified"),
147+
},
148+
)
149+
)
139150

140151
except dxpy.exceptions.DXAPIError as e:
141152
logger.error(f"DNAnexus API error listing {project_id}: {e}")
@@ -144,7 +155,9 @@ def list_files(self, project_id: str, file_type: str = None) -> list[RemoteFile]
144155
logger.info(f"Found {len(results)} files in {project_id}")
145156
return results
146157

147-
def stream_to_s3(self, remote_file: RemoteFile, s3_bucket: str, s3_key: str) -> dict:
158+
def stream_to_s3(
159+
self, remote_file: RemoteFile, s3_bucket: str, s3_key: str
160+
) -> dict:
148161
"""
149162
Stream file from DNAnexus directly to S3 using chunked reads.
150163
Uses DNAnexus download URL + boto3 multipart upload.
@@ -160,7 +173,9 @@ def stream_to_s3(self, remote_file: RemoteFile, s3_bucket: str, s3_key: str) ->
160173
sha256 = hashlib.sha256()
161174
bytes_transferred = 0
162175

163-
logger.info(f"Starting stream: {remote_file.file_name} → s3://{s3_bucket}/{s3_key}")
176+
logger.info(
177+
f"Starting stream: {remote_file.file_name} → s3://{s3_bucket}/{s3_key}"
178+
)
164179

165180
# Multipart upload config — 100MB parts, 4 parallel threads
166181
config = TransferConfig(
@@ -207,10 +222,17 @@ def chunked_stream() -> Iterator[bytes]:
207222
def get_metadata(self, file_id: str) -> dict:
208223
"""Fetch full metadata for a single DNAnexus file."""
209224
try:
210-
desc = dxpy.DXFile(file_id).describe(fields={
211-
"name": True, "size": True, "properties": True,
212-
"tags": True, "folder": True, "created": True, "modified": True,
213-
})
225+
desc = dxpy.DXFile(file_id).describe(
226+
fields={
227+
"name": True,
228+
"size": True,
229+
"properties": True,
230+
"tags": True,
231+
"folder": True,
232+
"created": True,
233+
"modified": True,
234+
}
235+
)
214236
return desc
215237
except dxpy.exceptions.DXAPIError as e:
216238
logger.error(f"Failed to get metadata for {file_id}: {e}")
@@ -219,6 +241,7 @@ def get_metadata(self, file_id: str) -> dict:
219241

220242
class _IterableToFileObj:
221243
"""Adapter to make a generator look like a file object for boto3 upload_fileobj."""
244+
222245
def __init__(self, iterable):
223246
self._iter = iterable
224247
self._buffer = b""

connectors/healthomics.py

Lines changed: 30 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -50,30 +50,36 @@ def list_files(self, project_id: str, file_type: str = None) -> list[RemoteFile]
5050
paginator = self.omics.get_paginator("list_read_sets")
5151
for page in paginator.paginate(sequenceStoreId=project_id):
5252
for rs in page.get("readSets", []):
53-
inferred_type = self.FILE_TYPE_MAP.get(rs.get("fileType", ""), "OTHER")
53+
inferred_type = self.FILE_TYPE_MAP.get(
54+
rs.get("fileType", ""), "OTHER"
55+
)
5456

5557
if file_type and inferred_type != file_type:
5658
continue
5759

5860
sample_id = rs.get("sampleId", rs.get("name", "unknown"))
5961

60-
results.append(RemoteFile(
61-
platform="HealthOmics",
62-
file_id=rs["id"],
63-
file_name=rs.get("name", rs["id"]),
64-
file_type=inferred_type,
65-
file_size_bytes=rs.get("sequenceInformation", {}).get("totalBaseCount", 0),
66-
source_path=f"omics://{project_id}/readSet/{rs['id']}",
67-
sample_id=sample_id,
68-
metadata={
69-
"store_id": project_id,
70-
"subject_id": rs.get("subjectId"),
71-
"reference_arn": rs.get("referenceArn"),
72-
"created_time": str(rs.get("creationTime")),
73-
"status": rs.get("status"),
74-
"tags": rs.get("tags", {}),
75-
}
76-
))
62+
results.append(
63+
RemoteFile(
64+
platform="HealthOmics",
65+
file_id=rs["id"],
66+
file_name=rs.get("name", rs["id"]),
67+
file_type=inferred_type,
68+
file_size_bytes=rs.get("sequenceInformation", {}).get(
69+
"totalBaseCount", 0
70+
),
71+
source_path=f"omics://{project_id}/readSet/{rs['id']}",
72+
sample_id=sample_id,
73+
metadata={
74+
"store_id": project_id,
75+
"subject_id": rs.get("subjectId"),
76+
"reference_arn": rs.get("referenceArn"),
77+
"created_time": str(rs.get("creationTime")),
78+
"status": rs.get("status"),
79+
"tags": rs.get("tags", {}),
80+
},
81+
)
82+
)
7783

7884
except ClientError as e:
7985
logger.error(f"HealthOmics API error for store {project_id}: {e}")
@@ -82,7 +88,9 @@ def list_files(self, project_id: str, file_type: str = None) -> list[RemoteFile]
8288
logger.info(f"Found {len(results)} ReadSets in store {project_id}")
8389
return results
8490

85-
def stream_to_s3(self, remote_file: RemoteFile, s3_bucket: str, s3_key: str) -> dict:
91+
def stream_to_s3(
92+
self, remote_file: RemoteFile, s3_bucket: str, s3_key: str
93+
) -> dict:
8694
"""
8795
Export a HealthOmics ReadSet to S3.
8896
Uses HealthOmics StartReadSetExportJob for large files,
@@ -91,7 +99,9 @@ def stream_to_s3(self, remote_file: RemoteFile, s3_bucket: str, s3_key: str) ->
9199
# For large files, use HealthOmics native export to S3
92100
store_id = remote_file.metadata.get("store_id")
93101

94-
logger.info(f"Exporting ReadSet {remote_file.file_id} → s3://{s3_bucket}/{s3_key}")
102+
logger.info(
103+
f"Exporting ReadSet {remote_file.file_id} → s3://{s3_bucket}/{s3_key}"
104+
)
95105

96106
md5 = hashlib.md5()
97107
sha256 = hashlib.sha256()

0 commit comments

Comments
 (0)