Coverage for event_normalizer/parsers.py: 100%
19 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-21 20:04 +0000
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-21 20:04 +0000
1"""Parsers for different data sources."""
3from datetime import UTC, datetime
4from typing import Any
6from .models import (
7 EventBurst,
8 EventCoinc,
9 EventLabels,
10 EventLinks,
11 EventMLyBurst,
12 EventSingleInspiral,
13 NormalizedEvent,
14)
15from .transformers import Transformers
17# ============================================================================
18# GraceDB, Kafka and File Parsers
19# ============================================================================
22def parse_gracedb_event(data: dict[str, Any]) -> NormalizedEvent:
23 """
24 Parse GraceDB REST API event response.
26 Extracts common fields and pipeline-specific data from extra_attributes.
27 """
29 uid = data.get("graceid")
30 if not uid:
31 raise ValueError("Missing 'graceid' in event data")
33 # Initialize a Transformer to deal with nested extra_attributes
34 transformer = Transformers(
35 group=data.get("group"),
36 pipeline=data.get("pipeline"),
37 search=data.get("search"),
38 event=data,
39 )
41 # Get channel names
42 channels_dict = transformer.transform_channels()
44 # Get entries from CoincInspiral table
45 coinc_dict = transformer.transform_coinc_table(
46 [
47 "mass",
48 "mchirp",
49 "minimum_duration",
50 "snr",
51 "ifos",
52 "end_time",
53 "end_time_ns",
54 "false_alarm_rate",
55 "combined_far",
56 ]
57 )
59 # Get entries from MultiBurst table
60 burst_dict = transformer.transform_burst_table(
61 [
62 "mchirp",
63 "snr",
64 "duration",
65 "ifos",
66 "start_time",
67 "start_time_ns",
68 "strain",
69 "peak_time",
70 "peak_time_ns",
71 "central_freq",
72 "bandwidth",
73 "amplitude",
74 "confidence",
75 "false_alarm_rate",
76 "ligo_axis_ra",
77 "ligo_axis_dec",
78 "ligo_angle",
79 "ligo_angle_sig",
80 "single_ifo_times",
81 "code",
82 ]
83 )
85 # Get entries from MLyBurst table
86 mly_dict = transformer.transform_mly_table(
87 [
88 "central_freq",
89 "bandwidth",
90 "central_time",
91 "detection_statistic",
92 "bbh",
93 "sglf",
94 "sghf",
95 "background",
96 "glitch",
97 "freq_correlation",
98 "mass1",
99 "mass2",
100 "mtotal",
101 "mchirp",
102 "spin1z",
103 "spin2z",
104 "end_time",
105 "end_time_ns",
106 "template_duration",
107 "SNR",
108 "scores.coherency",
109 "scores.coincidence",
110 "scores.combined",
111 ]
112 )
114 # Get SingleInspiral data per IFO
115 single_inspiral_dict = transformer.transform_single_inspiral(
116 [
117 "search",
118 "end_time",
119 "end_time_ns",
120 "end_time_gmst",
121 "impulse_time",
122 "impulse_time_ns",
123 "template_duration",
124 "event_duration",
125 "amplitude",
126 "eff_distance",
127 "coa_phase",
128 "mass1",
129 "mass2",
130 "mchirp",
131 "mtotal",
132 "eta",
133 "kappa",
134 "chi",
135 "tau0",
136 "tau2",
137 "tau3",
138 "tau4",
139 "tau5",
140 "ttotal",
141 "psi0",
142 "psi3",
143 "alpha",
144 "alpha1",
145 "alpha2",
146 "alpha3",
147 "alpha4",
148 "alpha5",
149 "alpha6",
150 "beta",
151 "f_final",
152 "snr",
153 "chisq",
154 "chisq_dof",
155 "bank_chisq",
156 "bank_chisq_dof",
157 "cont_chisq",
158 "cont_chisq_dof",
159 "sigmasq",
160 "rsqveto_duration",
161 "Gamma0",
162 "Gamma1",
163 "Gamma2",
164 "Gamma3",
165 "Gamma4",
166 "Gamma5",
167 "Gamma6",
168 "Gamma7",
169 "Gamma8",
170 "Gamma9",
171 "spin1x",
172 "spin1y",
173 "spin1z",
174 "spin2x",
175 "spin2y",
176 "spin2z",
177 ]
178 )
180 # Get labels
181 labels = transformer.transform_labels()
183 # Get instruments (convert None to empty list for model validation)
184 instruments = transformer.transform_instruments(data.get("instruments")) or []
186 # Get links
187 links = transformer.transform_links()
189 # Define NormalizedEvent
190 event = NormalizedEvent(
191 # Mandatory
192 uid=uid,
193 group=data.get("group"),
194 pipeline=data.get("pipeline"),
195 search=data.get("search"),
196 far=data.get("far"),
197 far_is_upper_limit=data.get("far_is_upper_limit"),
198 instruments=instruments,
199 H1_channel=channels_dict.get("H1_channel", "None"),
200 L1_channel=channels_dict.get("L1_channel", "None"),
201 V1_channel=channels_dict.get("V1_channel", "None"),
202 K1_channel=channels_dict.get("K1_channel", "None"),
203 reporting_latency=data.get("reporting_latency"),
204 gpstime=data.get("gpstime"),
205 last_updated=datetime.now(UTC),
206 # Optional root fields
207 alert_type=data.get("alert_type"),
208 submitter=data.get("submitter"),
209 offline=data.get("offline"),
210 nevents=data.get("nevents"),
211 likelihood=data.get("likelihood"),
212 superevent=data.get("superevent"),
213 created=data.get("created"),
214 processing_status=data.get("processing_status"),
215 content_id=data.get("content_id"),
216 # CBC specific
217 coinc=EventCoinc(**coinc_dict) if coinc_dict else None,
218 # Burst specific
219 burst=EventBurst(**burst_dict) if burst_dict else None,
220 # MLy specific
221 mly=EventMLyBurst(**mly_dict) if mly_dict else None,
222 # Single Inspiral per IFO
223 single_H1=EventSingleInspiral(**single_inspiral_dict["H1"])
224 if "H1" in single_inspiral_dict
225 else None,
226 single_L1=EventSingleInspiral(**single_inspiral_dict["L1"])
227 if "L1" in single_inspiral_dict
228 else None,
229 single_V1=EventSingleInspiral(**single_inspiral_dict["V1"])
230 if "V1" in single_inspiral_dict
231 else None,
232 single_K1=EventSingleInspiral(**single_inspiral_dict["K1"])
233 if "K1" in single_inspiral_dict
234 else None,
235 # Labels
236 labels=EventLabels(**labels) if labels else None,
237 # Links
238 links=EventLinks(**links) if links else None,
239 )
241 return event
244# TODO: parse_kafka_alert, parse_pastro, parse_embright