Skip to content
Merged
Show file tree
Hide file tree
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
11 changes: 10 additions & 1 deletion .github/workflows/project-build-test.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
22 changes: 21 additions & 1 deletion projects/online/online/monitor/pages/summary.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
compute_duty_cycle,
current_status,
downtime_breakdown,
find_overlaps,
)
from online.utils.timing import gps_now

Expand Down Expand Up @@ -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:
Expand All @@ -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"""
<p class="red" style="max-width: 820px; margin: 0 auto;">
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.
</p>
"""

def duty_cycle_html(self) -> str:
"""
Duty cycle stats and a table of the downtime we were
Expand All @@ -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",
Expand Down Expand Up @@ -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()
5 changes: 2 additions & 3 deletions projects/online/online/monitor/utils/plotting.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
23 changes: 23 additions & 0 deletions projects/online/online/monitor/utils/segments.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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 = []
Expand Down
19 changes: 14 additions & 5 deletions projects/online/online/utils/gdb.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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:
Expand All @@ -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)

Expand Down
21 changes: 16 additions & 5 deletions projects/online/online/utils/segments.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Loading