diff --git a/.github/workflows/project-build-test.yaml b/.github/workflows/project-build-test.yaml index 55116aa45..55fc7ebe6 100644 --- a/.github/workflows/project-build-test.yaml +++ b/.github/workflows/project-build-test.yaml @@ -122,9 +122,18 @@ jobs: if: ${{ github.event_name == 'push' || startsWith(github.ref, 'refs/tags/') }} env: tag: ${{ env.REGISTRY }}/${{ github.repository }}/${{ matrix.project }}:${{ github.ref_name }} + # online's environment is conda-based, and its libraries need to + # come ahead of the system's. docker import drops the def's + # %environment, and apptainer ignores a docker LD_LIBRARY_PATH + # but reads PREPEND_LD_LIBRARY_PATH and prepends it. + ld_prepend: ${{ matrix.project == 'online' && '/opt/env/lib' || '' }} run: | export TAG_LC=${tag,,} - cat app.tar.gz | docker import --change "ENV PATH=${{ env.PATH }}" - $TAG_LC + changes=(--change "ENV PATH=${{ env.PATH }}") + if [ -n "$ld_prepend" ]; then + changes+=(--change "ENV PREPEND_LD_LIBRARY_PATH=$ld_prepend") + fi + cat app.tar.gz | docker import "${changes[@]}" - $TAG_LC docker push $TAG_LC diff --git a/projects/online/online/monitor/pages/summary.py b/projects/online/online/monitor/pages/summary.py index 87e607616..f026fa31e 100644 --- a/projects/online/online/monitor/pages/summary.py +++ b/projects/online/online/monitor/pages/summary.py @@ -16,6 +16,7 @@ compute_duty_cycle, current_status, downtime_breakdown, + find_overlaps, ) from online.utils.timing import gps_now @@ -68,6 +69,7 @@ def __init__( self.html_file = self.out_dir / "summary.html" self.segments = None self.stats = None + self.overlaps = None @property def plot_name_dict(self) -> dict: @@ -81,6 +83,22 @@ def plot_name_dict(self) -> dict: "duty_cycle_trend": "Hourly search duty cycle", } + def overlap_html(self) -> str: + """Writes a warning if any segment overlaps another segment""" + if self.overlaps.empty: + return "" + + date_format = "%Y-%m-%d %H:%M:%S" + first = tconvert(self.overlaps["start"].iloc[0]).strftime(date_format) + duration = format_duration(self.overlaps["overlap"].sum()) + return f""" +

+ The segment record overlaps itself, counting {duration} of + time more than once, first at {first} UTC. The duty cycle and + uptime below may be wrong. +

+ """ + def duty_cycle_html(self) -> str: """ Duty cycle stats and a table of the downtime we were @@ -106,7 +124,8 @@ def duty_cycle_html(self) -> str: ] for window, stats in windows.items() ] - html = self.html_table( + html = self.overlap_html() + html += self.html_table( [ "Window", "Duty cycle", @@ -238,5 +257,6 @@ def create(self, segments: pd.DataFrame) -> None: """ self.segments = segments self.stats = compute_duty_cycle(segments) + self.overlaps = find_overlaps(segments) self.update_summary_plots() self.write_html() diff --git a/projects/online/online/monitor/utils/plotting.py b/projects/online/online/monitor/utils/plotting.py index 92ce95d51..f4d94480e 100644 --- a/projects/online/online/monitor/utils/plotting.py +++ b/projects/online/online/monitor/utils/plotting.py @@ -148,9 +148,8 @@ def latency_plot(plotsdir: Path, df: pd.DataFrame) -> None: latency = df["aframe latency"].dropna() median = latency.median() ninetieth_percentile = latency.quantile(0.9) - bins = np.logspace( - np.log10(latency.min()), np.log10(latency.max()), num=30 - ) + # Pad the range so that the fastest and slowest events are inside the bins + bins = np.geomspace(0.95 * latency.min(), 1.05 * latency.max(), num=30) plt.hist(latency, bins=bins, alpha=0.7) plt.axvline( median, color="red", linestyle="--", label=f"Median: {median:.2f} s" diff --git a/projects/online/online/monitor/utils/segments.py b/projects/online/online/monitor/utils/segments.py index 186c763db..e1ef33027 100644 --- a/projects/online/online/monitor/utils/segments.py +++ b/projects/online/online/monitor/utils/segments.py @@ -21,6 +21,9 @@ # How long the heartbeat can go unwritten before calling the search dead. STALE_SECONDS = 30.0 +# Times are written to 5 decimal places, so ignore overlap below this +OVERLAP_TOLERANCE = 1e-3 + def segment_dir(run_dir: Path) -> Path: return run_dir / "output" / "segments" @@ -111,6 +114,26 @@ def load_segments( return df.reset_index(drop=True) +def find_overlaps(df: pd.DataFrame) -> pd.DataFrame: + """ + The segments that start before earlier ones have ended, with how + long each overlaps for in an `overlap` column. + """ + df = df.sort_values("start") + + overlapping = [] + latest_stop = float("-inf") + for row in df.to_dict("records"): + if row["start"] < latest_stop: + # this segment starts before an earlier one has ended, and + # overlaps from its start until it or the earlier one ends + overlap = min(row["stop"], latest_stop) - row["start"] + if overlap > OVERLAP_TOLERANCE: + overlapping.append({**row, "overlap": overlap}) + latest_stop = max(latest_stop, row["stop"]) + return pd.DataFrame(overlapping, columns=[*df.columns, "overlap"]) + + def _fill_gaps(df: pd.DataFrame) -> pd.DataFrame: rows = df.to_dict("records") filled = [] diff --git a/projects/online/online/utils/gdb.py b/projects/online/online/utils/gdb.py index c0427e6ce..8c1177002 100644 --- a/projects/online/online/utils/gdb.py +++ b/projects/online/online/utils/gdb.py @@ -112,6 +112,7 @@ def submit(self, event: Event): event_dir = self.write_dir / event.event_dir filename = event_dir / event.filename self.logger.info("Creating event in GraceDB") + submission_start = gps_now() response = self.create_event( group="CBC", pipeline="aframe", @@ -136,11 +137,12 @@ def submit(self, event: Event): with open(filename, "w") as f: f.write(url) + submission_end = gps_now() + # record latencies for this event; # TODO: determine underlying issue here # Handle issue where sometimes the pipeline lags, # and the frame file has already left the buffer - submission_time = gps_now() try: t_write = event.get_frame_write_time() except FileNotFoundError: @@ -157,13 +159,20 @@ def submit(self, event: Event): t_write = int(event.gpstime) # time to submit since event occured and since the file was written - total_latency = submission_time - event.gpstime + total_latency = submission_start - event.gpstime write_latency = t_write - event.gpstime - aframe_latency = submission_time - t_write + aframe_latency = submission_start - t_write + submission_duration = submission_end - submission_start latency_fname = event_dir / "latency.log" - latency = "Total Latency (s),Write Latency (s),Aframe Latency (s)\n" - latency += f"{total_latency},{write_latency},{aframe_latency}" + latency = ( + "Total Latency (s),Write Latency (s),Aframe Latency (s)," + "Submission Duration (s)\n" + ) + latency += ( + f"{total_latency},{write_latency}," + f"{aframe_latency},{submission_duration}" + ) with open(latency_fname, "w") as f: f.write(latency) diff --git a/projects/online/online/utils/segments.py b/projects/online/online/utils/segments.py index 8049ebccd..113d56d4a 100644 --- a/projects/online/online/utils/segments.py +++ b/projects/online/online/utils/segments.py @@ -99,12 +99,23 @@ def __init__( def _flush_current_file(self) -> None: with open(self.current_file, "r") as f: current_data = json.load(f) - self._append( - current_data["state"], - current_data["start"], - current_data["stop"], - current_data["ifos_ready"], + + # The last process may have been interrupted on its way out, + # after recording this segment but before removing the file + with open(self.segments_file, "r") as f: + last_state, last_start = f.read().splitlines()[-1].split(",")[:2] + recorded = ( + last_state == current_data["state"] + and last_start == f"{current_data['start']:.5f}" ) + + if not recorded: + self._append( + current_data["state"], + current_data["start"], + current_data["stop"], + current_data["ifos_ready"], + ) self.current_file.unlink() def _append(