Skip to content

Commit a6b2f4f

Browse files
committed
use trigger instead of 3 second timer
1 parent 51853c7 commit a6b2f4f

2 files changed

Lines changed: 19 additions & 8 deletions

File tree

‎operators/read-tem-data/operator.json‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,14 @@
2929
"required": true
3030
}
3131
],
32+
"triggers": [
33+
{
34+
"name": "read_now",
35+
"label": "Read data now",
36+
"description": "Manually read and emit the data without waiting.",
37+
"mode": "single"
38+
}
39+
],
3240
"parallel_config": {
3341
"type": "none"
3442
}

‎operators/read-tem-data/run.py‎

Lines changed: 11 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -15,11 +15,11 @@
1515

1616
@operator
1717
def read_data_ncempy(
18-
inputs: BytesMessage | None, parameters: dict[str, Any]
18+
inputs: BytesMessage | None, parameters: dict[str, Any], trigger=None
1919
) -> BytesMessage | None:
20-
"""This reads data from disk and sends it on."""
21-
22-
# This operator does not require inputs
20+
"""Read and emit TEM data from a specified file when triggered."""
21+
if trigger is None:
22+
return None
2323

2424
# Extract parameters
2525
directory = parameters.get("raw_data_dir", "/test_data")
@@ -40,11 +40,14 @@ def read_data_ncempy(
4040
except Exception as e:
4141
logger.info(f"Problem loading file. Error: {e}")
4242
return None
43-
finally:
44-
# TODO: Use trigger instead of time.sleep to control when data is sent
45-
time.sleep(3.0) # wait to read it again.
4643

4744
# Process and return result if the data was loaded successfully
4845
data_bytes = data.tobytes()
49-
header = MessageHeader(subject=MessageSubject.BYTES, meta={'shape': data.shape, 'dtype': str(data.dtype)})
46+
header = MessageHeader(
47+
subject=MessageSubject.BYTES,
48+
meta={
49+
"shape": data.shape,
50+
"dtype": str(data.dtype),
51+
},
52+
)
5053
return BytesMessage(header=header, data=data_bytes)

0 commit comments

Comments
 (0)