From 32655e1454a25388b5cc7142bb899330738cce79 Mon Sep 17 00:00:00 2001 From: "slominskir-coding-agent[bot]" <335168828+slominskir-coding-agent[bot]@users.noreply.github.com> Date: Sat, 3 Oct 2026 23:41:53 -0400 Subject: [PATCH] Detect frozen PVs with an independent CA context IOC and network restarts have historically left PVs "frozen": a monitor that stops delivering while the IOC still serves the PV, which a restart of epics2web fixes. Nothing reported that, so Tomcat is restarted periodically by cron. FrozenPvDetector gives suspicious PVs a short-lived subscription in a second CA context, with its own virtual circuits, made like the monitor's so the IOC applies the same deadband to both. Comparing values instead would mistake a change within a deadband for a missed update. A PV is frozen when its monitor has been disconnected, or still connecting, past the grace period while the independent subscription connects, or when a monitor that used to update has gone quiet and the independent subscription receives a change it doesn't. Only PVs not connected past the grace period, and quiet PVs that have changed before, are probed, at most 20 at a time and each at most once every 30 check intervals. A PV stops being frozen when its monitor gets an update. Detection is report-only: frozen PVs are logged as warnings and listed by /healthcheck, and the default and strict responses are unchanged. ?frozen=true answers 503 when a PV is frozen, for restart automation once it's trusted. FROZEN_CHECK_SECONDS (default 10) sets the check interval, from which the other timings derive, and FROZEN_PV_CHECK=false turns detection off. FrozenPvDetectorTest covers the decisions with a controlled clock. FrozenPvDetectionTest runs the detector against the test IOC with two real CA contexts, and simulates a lost subscription by clearing a monitor's CAJ subscription behind its back. HealthcheckTest checks that working monitors and a stopped IOC aren't reported frozen. Part of #28 Co-Authored-By: Claude Opus 5.5 --- README.md | 5 + build.yaml | 1 + .../jlab/epics2web/FrozenPvDetectionTest.java | 137 +++++++++ .../org/jlab/epics2web/HealthcheckTest.java | 26 ++ .../java/org/jlab/epics2web/Application.java | 78 +++++ .../epics2web/controller/Healthcheck.java | 43 ++- .../jlab/epics2web/epics/CaProbeFactory.java | 125 ++++++++ .../jlab/epics2web/epics/ChannelMonitor.java | 9 + .../epics2web/epics/FrozenPvDetector.java | 291 ++++++++++++++++++ .../epics2web/epics/FrozenPvDetectorTest.java | 250 +++++++++++++++ 10 files changed, 960 insertions(+), 5 deletions(-) create mode 100644 src/integration/java/org/jlab/epics2web/FrozenPvDetectionTest.java create mode 100644 src/main/java/org/jlab/epics2web/epics/CaProbeFactory.java create mode 100644 src/main/java/org/jlab/epics2web/epics/FrozenPvDetector.java create mode 100644 src/test/java/org/jlab/epics2web/epics/FrozenPvDetectorTest.java diff --git a/README.md b/README.md index af8384e..35e0773 100644 --- a/README.md +++ b/README.md @@ -69,6 +69,11 @@ A client that stops reading makes writes to it block once the network buffers fi By default the response is 200 whenever the server is up, which suits load balancers: an IOC being down affects every instance alike. With `?strict=true` the response is 503 when a PV that was connected has disconnected, which suits monitoring that should alert on that, such as Nagios. PVs that never connected are listed but don't fail strict mode, since they may simply not exist. +#### Frozen PVs +A PV is frozen when its monitor has stopped working while the IOC still serves it, so restarting epics2web would be expected to fix it. epics2web looks for them every **FROZEN_CHECK_SECONDS** (default 10) by giving suspicious PVs a short-lived subscription in a second, independent CA context with its own connections: PVs disconnected for longer than the grace period, and PVs that used to update but have been quiet for 6 check intervals. A PV is frozen if its monitor stays disconnected while the independent subscription connects, or if the independent subscription receives a change the monitor doesn't. An IOC that's down doesn't make its PVs frozen. Through a gateway, both contexts reach the PV through the gateway, so this finds problems between epics2web and the gateway, not inside it. + +Frozen PVs are logged as warnings and listed by `/healthcheck` with `frozen`, `frozen_minutes` and `frozen_reason`. They don't change the default or strict response; with `?frozen=true` the response is 503 when a PV is frozen, for automation that restarts epics2web. **FROZEN_PV_CHECK**=false turns detection off. + ### Logging This app is designed to run on Tomcat so [Tomcat logging configuration](https://tomcat.apache.org/tomcat-9.0-doc/logging.html) applies. We use the built-in JVM logging library, which Tomcat uses with some slight modifications to support separate classloaders. In the past we bundled an application [logging.properites](https://github.com/JeffersonLab/epics2web/blob/956894699ef1b303907a04720aeb50260ffa72b1/src/main/resources/logging.properties) inside the epics2web.war file. We no longer do that because it then appears to require repackaging/rebuilding a new version of the app to modify the logging config as the app bundled config overrides the global Tomcat config at conf/logging.properties. The recommend logging strategy is to now make configuration in the global Tomcat config so as to make it easy to modify logging levels. An app specific handler can be created. The global configuration location is generally set by the Tomcat default start script via JVM system properties. The system properties should look something like: - `-Djava.util.logging.config.file=/usr/share/tomcat/conf/logging.properties` diff --git a/build.yaml b/build.yaml index 38abd3d..37ee8f7 100644 --- a/build.yaml +++ b/build.yaml @@ -15,6 +15,7 @@ services: WEBSOCKET_SEND_TIMEOUT_SECONDS: 3 # Short, so HealthcheckTest runs in seconds HEALTHCHECK_GRACE_SECONDS: 2 + FROZEN_CHECK_SECONDS: 1 build: context: . dockerfile: Dockerfile diff --git a/src/integration/java/org/jlab/epics2web/FrozenPvDetectionTest.java b/src/integration/java/org/jlab/epics2web/FrozenPvDetectionTest.java new file mode 100644 index 0000000..9f8992d --- /dev/null +++ b/src/integration/java/org/jlab/epics2web/FrozenPvDetectionTest.java @@ -0,0 +1,137 @@ +package org.jlab.epics2web; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import com.cosylab.epics.caj.CAJContext; +import gov.aps.jca.JCALibrary; +import gov.aps.jca.Monitor; +import gov.aps.jca.configuration.DefaultConfiguration; +import gov.aps.jca.dbr.DBR; +import gov.aps.jca.dbr.DBRType; +import java.lang.reflect.Field; +import java.time.Clock; +import java.time.Duration; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import org.jlab.epics2web.epics.CaProbeFactory; +import org.jlab.epics2web.epics.ChannelManager; +import org.jlab.epics2web.epics.ChannelMonitor; +import org.jlab.epics2web.epics.FrozenPvDetector; +import org.jlab.epics2web.epics.PvListener; +import org.junit.After; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.Timeout; + +/** + * Runs FrozenPvDetector in this JVM against the test IOC, through its published ports, with two + * real CA contexts. A lost subscription is simulated by clearing a monitor's CAJ subscription + * without telling the monitor, as a CA client bug might. + */ +public class FrozenPvDetectionTest { + + private static final String MISSING = "epics2web:test:frozen:missing"; + + @Rule public Timeout globalTimeout = Timeout.seconds(90); + + private CAJContext monitorContext; + private CAJContext probeContext; + private ScheduledExecutorService timeoutExecutor; + private ExecutorService callbackExecutor; + private ChannelManager manager; + private FrozenPvDetector detector; + private final PvListener listener = new NoopListener(); + + @BeforeClass + public static void disableRepeater() { + System.setProperty("CA_DISABLE_REPEATER", "true"); // As in the unit tests (#32) + } + + @Before + public void setUp() throws Exception { + monitorContext = newContext(); + probeContext = newContext(); + timeoutExecutor = Executors.newSingleThreadScheduledExecutor(); + callbackExecutor = Executors.newCachedThreadPool(); + manager = new ChannelManager(monitorContext, timeoutExecutor, callbackExecutor); + detector = + new FrozenPvDetector( + () -> FrozenPvDetector.viewsOf(manager.getMonitorMap()), + new CaProbeFactory(probeContext), + Clock.systemUTC(), + Duration.ofSeconds(2), + Duration.ofSeconds(1), + 20); + } + + @After + public void tearDown() throws Exception { + detector.close(); + manager.removeAll(listener); + timeoutExecutor.shutdownNow(); + callbackExecutor.shutdownNow(); + monitorContext.destroy(); + probeContext.destroy(); + } + + @Test + public void lostSubscriptionIsFrozenAndNothingElseIs() throws Exception { + manager.addPv(listener, "HELLO"); // Changes every 0.2 s + manager.addPv(listener, "channel1"); // Never changes + manager.addPv(listener, MISSING); // Never connects + + // Working monitors, a PV that doesn't change, and a PV that doesn't exist aren't frozen. 10 s + // is longer than the 6 s after which a PV counts as quiet, and the 2 s grace period. + checkFor(10); + assertEquals(Map.of(), detector.getFrozen()); + assertTrue(manager.getMonitorMap().get("HELLO").getUpdateCount() > 10); + + loseSubscription(manager.getMonitorMap().get("HELLO")); + + long deadline = System.currentTimeMillis() + 30_000; + while (!detector.getFrozen().containsKey("HELLO") && System.currentTimeMillis() < deadline) { + checkFor(1); + } + assertEquals(Set.of("HELLO"), detector.getFrozen().keySet()); + assertTrue( + detector.getFrozen().get("HELLO").reason(), + detector.getFrozen().get("HELLO").reason().contains("didn't receive")); + } + + private void checkFor(int seconds) throws InterruptedException { + for (int i = 0; i < seconds; i++) { + Thread.sleep(1_000); + detector.check(); + } + } + + /** Cancel the monitor's CAJ subscription without telling the monitor. */ + private static void loseSubscription(ChannelMonitor monitor) throws Exception { + Field field = ChannelMonitor.class.getDeclaredField("monitor"); + field.setAccessible(true); + ((Monitor) field.get(monitor)).clear(); + } + + private static CAJContext newContext() throws Exception { + DefaultConfiguration config = new DefaultConfiguration("test"); + config.setAttribute("class", JCALibrary.CHANNEL_ACCESS_JAVA); + config.setAttribute("addr_list", "127.0.0.1"); + config.setAttribute("auto_addr_list", "false"); + return (CAJContext) JCALibrary.getInstance().createContext(config); + } + + private static class NoopListener implements PvListener { + @Override + public void notifyPvInfo( + String pv, boolean couldConnect, DBRType type, Integer count, String[] enumLabels) {} + + @Override + public void notifyPvUpdate(String pv, DBR dbr) {} + } +} diff --git a/src/integration/java/org/jlab/epics2web/HealthcheckTest.java b/src/integration/java/org/jlab/epics2web/HealthcheckTest.java index 53e14f3..10104fa 100644 --- a/src/integration/java/org/jlab/epics2web/HealthcheckTest.java +++ b/src/integration/java/org/jlab/epics2web/HealthcheckTest.java @@ -1,6 +1,7 @@ package org.jlab.epics2web; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -77,6 +78,12 @@ public void disconnectedPvFailsStrictModeOnly() throws Exception { assertTrue("Reported after only " + reportedAfter + " ms", reportedAfter >= 2_000); assertEquals(200, get("healthcheck").statusCode()); assertEquals(503, get("healthcheck?strict=true").statusCode()); + + // The IOC is down, so restarting epics2web wouldn't help: not frozen. With + // FROZEN_CHECK_SECONDS 1, the detector has probed it within 2 s of the grace period ending. + Thread.sleep(3_000); + assertFalse(entry(get("healthcheck"), "channel1").containsKey("frozen")); + assertEquals(200, get("healthcheck?frozen=true").statusCode()); } finally { docker("start", "softioc"); waitForEntry("channel1", false); // Reconnected @@ -85,6 +92,25 @@ public void disconnectedPvFailsStrictModeOnly() throws Exception { assertEquals(200, get("healthcheck?strict=true").statusCode()); } + /** + * Working monitors aren't frozen. HELLO changes every 0.2 s; channel1 never changes. 10 s is + * longer than the 6 s after which a quiet PV is probed, with FROZEN_CHECK_SECONDS 1. + */ + @Test + public void workingMonitorsAreNotFrozen() throws Exception { + WebSocket hello = monitor("HELLO"); + WebSocket channel1 = monitor("channel1"); + try { + Thread.sleep(10_000); + assertEquals(200, get("healthcheck?frozen=true").statusCode()); + assertNull(entry(get("healthcheck"), "HELLO")); + assertNull(entry(get("healthcheck"), "channel1")); + } finally { + hello.abort(); + channel1.abort(); + } + } + private static WebSocket monitor(String pv) { WebSocket socket = HTTP.newWebSocketBuilder() diff --git a/src/main/java/org/jlab/epics2web/Application.java b/src/main/java/org/jlab/epics2web/Application.java index bc52e39..f701322 100644 --- a/src/main/java/org/jlab/epics2web/Application.java +++ b/src/main/java/org/jlab/epics2web/Application.java @@ -15,6 +15,7 @@ import jakarta.websocket.SendResult; import jakarta.websocket.Session; import java.io.IOException; +import java.time.Clock; import java.time.Duration; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ExecutorService; @@ -26,8 +27,10 @@ import java.util.concurrent.locks.LockSupport; import java.util.logging.Level; import java.util.logging.Logger; +import org.jlab.epics2web.epics.CaProbeFactory; import org.jlab.epics2web.epics.ChannelManager; import org.jlab.epics2web.epics.ContextFactory; +import org.jlab.epics2web.epics.FrozenPvDetector; import org.jlab.epics2web.websocket.WebSocketSessionManager; import org.jlab.epics2web.websocket.WriteQueue; import org.jlab.epics2web.websocket.WriteStrategy; @@ -46,6 +49,9 @@ public class Application implements ServletContextListener { public static ChannelManager channelManager = null; public static WebSocketSessionManager sessionManager = null; + /** Null if frozen PV detection is off. */ + public static FrozenPvDetector frozenPvDetector = null; + private static final int TIMEOUT_EXECUTOR_POOL_SIZE = 1; private static final Logger LOGGER = Logger.getLogger(Application.class.getName()); @@ -67,11 +73,23 @@ public class Application implements ServletContextListener { public static final long HEALTHCHECK_GRACE_SECONDS = getSecondsFromEnv("HEALTHCHECK_GRACE_SECONDS", 30); + /** Whether to look for frozen PVs; FROZEN_PV_CHECK=false turns it off. */ + static final boolean FROZEN_PV_CHECK = + !"false".equalsIgnoreCase(System.getenv("FROZEN_PV_CHECK")); + + /** How often to look for frozen PVs; FrozenPvDetector derives its other timings from it. */ + static final long FROZEN_CHECK_SECONDS = getSecondsFromEnv("FROZEN_CHECK_SECONDS", 10); + + /** The most PVs probed at once in the independent context. */ + private static final int FROZEN_CHECK_MAX_PROBES = 20; + private static ScheduledExecutorService timeoutExecutor = null; private static ExecutorService callbackExecutor = null; private static ExecutorService writerExecutor = null; private static ScheduledExecutorService sessionCheckExecutor = null; private static ExecutorService pingExecutor = null; + private static ScheduledExecutorService frozenCheckExecutor = null; + private static volatile CAJContext probeContext = null; private static ContextFactory factory = null; private static volatile CAJContext context = null; @@ -205,6 +223,10 @@ public void contextInitialized(ServletContextEvent sce) { LOGGER.log(Level.SEVERE, "Unable to register context callbacks", e); } + if (FROZEN_PV_CHECK) { + startFrozenPvDetection(); + } + if (WRITE_STRATEGY == WriteStrategy.ASYNC_QUEUE) { writerExecutor.execute( new Runnable() { @@ -296,6 +318,22 @@ public void contextDestroyed(ServletContextEvent sce) { pingExecutor.shutdownNow(); } + if (frozenCheckExecutor != null) { + frozenCheckExecutor.shutdownNow(); + } + + if (frozenPvDetector != null) { + frozenPvDetector.close(); + } + + if (probeContext != null) { + try { + probeContext.destroy(); + } catch (CAException e) { + LOGGER.log(Level.WARNING, "Unable to destroy probe context", e); + } + } + if (timeoutExecutor != null) { try { if (!timeoutExecutor.awaitTermination(5, TimeUnit.SECONDS)) { @@ -329,6 +367,46 @@ public void contextDestroyed(ServletContextEvent sce) { } } + /** + * Look for frozen PVs with subscriptions in a second CA context, which has its own virtual + * circuits. See FrozenPvDetector. + */ + private void startFrozenPvDetection() { + try { + probeContext = factory.newContext(); + } catch (Exception e) { + LOGGER.log(Level.SEVERE, "Unable to create CA context for frozen PV detection", e); + return; + } + + frozenPvDetector = + new FrozenPvDetector( + () -> FrozenPvDetector.viewsOf(channelManager.getMonitorMap()), + new CaProbeFactory(probeContext), + Clock.systemUTC(), + Duration.ofSeconds(HEALTHCHECK_GRACE_SECONDS), + Duration.ofSeconds(FROZEN_CHECK_SECONDS), + FROZEN_CHECK_MAX_PROBES); + + // Its own thread: closing a probe can wait briefly for the IOC + frozenCheckExecutor = + Executors.newSingleThreadScheduledExecutor( + new CustomPrefixThreadFactory("Frozen-PV-Check-")); + frozenCheckExecutor.scheduleWithFixedDelay( + () -> { + try { + frozenPvDetector.check(); + } catch (RuntimeException e) { // An exception would cancel the schedule + LOGGER.log(Level.WARNING, "Unable to check for frozen PVs", e); + } + }, + FROZEN_CHECK_SECONDS, + FROZEN_CHECK_SECONDS, + TimeUnit.SECONDS); + + LOGGER.log(Level.INFO, "Checking for frozen PVs every {0} s", FROZEN_CHECK_SECONDS); + } + private void registerContextListeners(CAJContext c) throws CAException { c.addContextExceptionListener( new ContextExceptionListener() { diff --git a/src/main/java/org/jlab/epics2web/controller/Healthcheck.java b/src/main/java/org/jlab/epics2web/controller/Healthcheck.java index 573621e..b3647d8 100644 --- a/src/main/java/org/jlab/epics2web/controller/Healthcheck.java +++ b/src/main/java/org/jlab/epics2web/controller/Healthcheck.java @@ -12,12 +12,14 @@ import java.io.PrintWriter; import java.time.Duration; import java.time.Instant; +import java.util.LinkedHashMap; import java.util.Map; import java.util.logging.Level; import java.util.logging.Logger; import org.jlab.epics2web.Application; import org.jlab.epics2web.epics.ChannelManager; import org.jlab.epics2web.epics.ChannelMonitor; +import org.jlab.epics2web.epics.FrozenPvDetector; /** * Controller for Healthcheck page. Return 200 OK, for healthy Return 503 Service Unavailable for @@ -48,8 +50,11 @@ protected void doGet(HttpServletRequest request, HttpServletResponse response) // Strict mode answers 503 when a PV is reported, for monitoring that alerts on PVs. The default // answers 200 whenever the server is up, for load balancers: an IOC being down affects every // instance alike, and restarting the server doesn't bring it back. - String strictParam = request.getParameter("strict"); - boolean strict = strictParam != null && !"false".equalsIgnoreCase(strictParam); + boolean strict = isSet(request.getParameter("strict")); + + // Frozen mode answers 503 when a PV is frozen: its monitor stopped working while the IOC still + // serves it, which a restart is expected to fix. See FrozenPvDetector. + boolean frozenMode = isSet(request.getParameter("frozen")); boolean healthy = true; @@ -57,7 +62,7 @@ protected void doGet(HttpServletRequest request, HttpServletResponse response) Instant now = Instant.now(); - JsonArrayBuilder unhealthyChannelArray = Json.createArrayBuilder(); + Map entries = new LinkedHashMap<>(); for (Map.Entry entry : monitorMap.entrySet()) { String pv = entry.getKey(); @@ -83,17 +88,41 @@ protected void doGet(HttpServletRequest request, HttpServletResponse response) unhealthyChannel.add("state", state.name()); unhealthyChannel.add( "disconnected_minutes", String.format("%.1f", notConnected.toSeconds() / 60.0)); - unhealthyChannelArray.add(unhealthyChannel); + entries.put(pv, unhealthyChannel); } } + Map frozen = + Application.frozenPvDetector == null ? Map.of() : Application.frozenPvDetector.getFrozen(); + + for (Map.Entry entry : frozen.entrySet()) { + String pv = entry.getKey(); + ChannelMonitor monitor = monitorMap.get(pv); + JsonObjectBuilder frozenChannel = + entries.computeIfAbsent( + pv, + k -> + Json.createObjectBuilder() + .add("name", pv) + .add("state", monitor == null ? "UNKNOWN" : monitor.getState().name())); + frozenChannel.add("frozen", true); + frozenChannel.add( + "frozen_minutes", + String.format( + "%.1f", Duration.between(entry.getValue().since(), now).toSeconds() / 60.0)); + frozenChannel.add("frozen_reason", entry.getValue().reason()); + } + + JsonArrayBuilder unhealthyChannelArray = Json.createArrayBuilder(); + entries.values().forEach(unhealthyChannelArray::add); + response.setContentType("application/json"); PrintWriter pw = response.getWriter(); response.setStatus(HttpServletResponse.SC_OK); - if (strict && !healthy) { + if ((strict && !healthy) || (frozenMode && !frozen.isEmpty())) { response.setStatus(HttpServletResponse.SC_SERVICE_UNAVAILABLE); } @@ -109,4 +138,8 @@ protected void doGet(HttpServletRequest request, HttpServletResponse response) LOGGER.log(Level.SEVERE, "PrintWriter Error"); } } + + private static boolean isSet(String param) { + return param != null && !"false".equalsIgnoreCase(param); + } } diff --git a/src/main/java/org/jlab/epics2web/epics/CaProbeFactory.java b/src/main/java/org/jlab/epics2web/epics/CaProbeFactory.java new file mode 100644 index 0000000..4ca5379 --- /dev/null +++ b/src/main/java/org/jlab/epics2web/epics/CaProbeFactory.java @@ -0,0 +1,125 @@ +package org.jlab.epics2web.epics; + +import com.cosylab.epics.caj.CAJChannel; +import com.cosylab.epics.caj.CAJContext; +import gov.aps.jca.CAException; +import gov.aps.jca.Channel; +import gov.aps.jca.Monitor; +import gov.aps.jca.event.MonitorEvent; +import gov.aps.jca.event.MonitorListener; +import java.time.Instant; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.logging.Level; +import java.util.logging.Logger; + +/** + * Probes PVs for FrozenPvDetector with subscriptions in their own CA context, separate from the one + * monitors use. Each subscription is made like ChannelMonitor's, so the IOC applies the same + * deadband to both. + */ +public class CaProbeFactory implements FrozenPvDetector.ProbeFactory { + + private static final Logger LOGGER = Logger.getLogger(CaProbeFactory.class.getName()); + + private final CAJContext context; + + /** + * Create a new CaProbeFactory. + * + * @param context A CA context used only for probes + */ + public CaProbeFactory(CAJContext context) { + this.context = context; + } + + @Override + public FrozenPvDetector.Probe start(String pv) throws CAException { + CAJChannel channel = (CAJChannel) context.createChannel(pv); + context.flushIO(); + return new CaProbe(pv, channel); + } + + private class CaProbe implements FrozenPvDetector.Probe, MonitorListener { + private final String pv; + private final CAJChannel channel; + private final CountDownLatch firstUpdate = new CountDownLatch(1); + private Monitor monitor; // Only used from the detector's thread + private volatile Instant connectedAt; + private volatile Instant firstChangeAt; + + private CaProbe(String pv, CAJChannel channel) { + this.pv = pv; + this.channel = channel; + } + + /** + * Subscribe once connected. CAJ's own callbacks mustn't call back into it, so it's done here. + */ + @Override + public void poll() { + if (monitor == null && channel.getConnectionState() == Channel.ConnectionState.CONNECTED) { + try { + // As ChannelMonitor does: arrays aren't handled, except BYTE[] as a long string + int count = 1; + if (channel.getFieldType().isBYTE() && channel.getElementCount() > 1) { + count = channel.getElementCount(); + } + monitor = channel.addMonitor(channel.getFieldType(), count, Monitor.VALUE, this); + context.flushIO(); + } catch (CAException | IllegalStateException e) { + LOGGER.log(Level.FINE, "Unable to subscribe probe of " + pv, e); + } + } + } + + @Override + public void monitorChanged(MonitorEvent event) { + if (event.getStatus() == null || !event.getStatus().isSuccessful()) { + return; + } + Instant now = Instant.now(); + if (connectedAt == null) { + connectedAt = now; // The initial value sent in reply to the subscription + firstUpdate.countDown(); + } else if (firstChangeAt == null) { + firstChangeAt = now; + } + } + + @Override + public Instant connectedAt() { + return connectedAt; + } + + @Override + public Instant firstChangeAt() { + return firstChangeAt; + } + + @Override + public void close() { + if (monitor != null) { + // As in ChannelMonitor.close(): CAJ sends a cancel at once but queues the add, so let the + // IOC confirm the subscription first, or it may get the cancel first and drop the circuit + try { + if (channel.getConnectionState() == Channel.ConnectionState.CONNECTED) { + firstUpdate.await(ChannelMonitor.TIMEOUT_MILLIS, TimeUnit.MILLISECONDS); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + try { + monitor.clear(); + } catch (CAException | IllegalStateException e) { + LOGGER.log(Level.FINE, "Unable to clear probe of " + pv, e); + } + } + try { + context.destroyChannel(channel, false); + } catch (CAException | IllegalStateException e) { + LOGGER.log(Level.FINE, "Unable to destroy probe channel of " + pv, e); + } + } + } +} diff --git a/src/main/java/org/jlab/epics2web/epics/ChannelMonitor.java b/src/main/java/org/jlab/epics2web/epics/ChannelMonitor.java index ef3cb6f..849b6c8 100644 --- a/src/main/java/org/jlab/epics2web/epics/ChannelMonitor.java +++ b/src/main/java/org/jlab/epics2web/epics/ChannelMonitor.java @@ -20,6 +20,7 @@ import java.util.Date; import java.util.Set; import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.logging.Level; import java.util.logging.Logger; @@ -37,6 +38,9 @@ public class ChannelMonitor implements Closeable { private volatile DBR lastDbr = null; + /** The number of updates received, including the initial value sent on each (re)subscription. */ + private final AtomicLong updateCount = new AtomicLong(); + /** * We don't use TIME typed DBR, so we just track 'received' timestamp (which may differ from IOC * 'generated' timestamp) @@ -208,6 +212,10 @@ public Date getLastTimestamp() { return lastTimestamp; } + public long getUpdateCount() { + return updateCount.get(); + } + /** * Close the ChannelMonitor. * @@ -498,6 +506,7 @@ public void monitorChanged(MonitorEvent me) { DBR dbr = me.getDBR(); lastDbr = dbr; + updateCount.incrementAndGet(); lastTimestamp = new Date(); // Make sure handlers do not call back into CA lib on this callback thread. diff --git a/src/main/java/org/jlab/epics2web/epics/FrozenPvDetector.java b/src/main/java/org/jlab/epics2web/epics/FrozenPvDetector.java new file mode 100644 index 0000000..599cdc9 --- /dev/null +++ b/src/main/java/org/jlab/epics2web/epics/FrozenPvDetector.java @@ -0,0 +1,291 @@ +package org.jlab.epics2web.epics; + +import java.io.Closeable; +import java.time.Clock; +import java.time.Duration; +import java.time.Instant; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Supplier; +import java.util.logging.Level; +import java.util.logging.Logger; + +/** + * Finds frozen PVs: PVs whose monitor has stopped working while the IOC still serves them, which + * restarting epics2web would be expected to fix. A PV whose IOC is down isn't frozen. + * + *

Suspicious PVs get a short-lived subscription in a second, independent CA context, with its + * own virtual circuits. The subscription is made like the monitor's, so the IOC applies the same + * deadband to both, and comparing them doesn't mistake a change within the deadband for a missed + * update. A PV is frozen when either: + * + *

    + *
  • its monitor has been disconnected, or still connecting, for longer than the grace period, + * but the independent subscription connects and receives the PV; or + *
  • its monitor, which used to get updates, has gone quiet, and the independent subscription + * receives a change the monitor doesn't. + *
+ * + *

A PV stops being frozen when its monitor receives an update. Through a gateway, both contexts + * reach the PV through the gateway, so this finds problems between epics2web and the gateway, not + * inside it. + */ +public class FrozenPvDetector implements Closeable { + + private static final Logger LOGGER = Logger.getLogger(FrozenPvDetector.class.getName()); + + /** What the detector needs to know about a monitor. */ + public record MonitorView( + ChannelMonitor.MonitorState state, + Instant stateChanged, + Instant lastUpdate, + long updateCount) {} + + /** A subscription to a PV in the independent CA context. */ + public interface Probe extends Closeable { + + /** Called before each evaluation, from the detector's thread. */ + default void poll() {} + + /** + * @return When the first update arrived, which shows the IOC serves the PV, or null + */ + Instant connectedAt(); + + /** + * @return When the first update after the initial one arrived, which shows the PV changed, or + * null + */ + Instant firstChangeAt(); + + @Override + void close(); + } + + /** Starts probes. */ + public interface ProbeFactory { + Probe start(String pv) throws Exception; + } + + /** A frozen PV: since when, and why. */ + public record FrozenPv(Instant since, String reason) {} + + private enum Kind { + NOT_CONNECTED, + QUIET + } + + private record ProbeSession(Probe probe, Kind kind, Instant started, long updateCountAtStart) {} + + private final Supplier> monitors; + private final ProbeFactory probes; + private final Clock clock; + private final Duration grace; + private final Duration quiet; + private final Duration window; + private final Duration margin; + private final Duration reprobe; + private final int maxProbes; + + private final Map active = new HashMap<>(); + private final Map lastProbed = new HashMap<>(); + private final Map updateCountWhenFrozen = new HashMap<>(); + private final Map frozen = new ConcurrentHashMap<>(); + + /** + * Create a new FrozenPvDetector. Timings derive from the check interval: a PV is quiet after 6 + * intervals without an update, a probe lasts at most 6 intervals, a PV is probed at most once + * every 30 intervals, and a probe's evidence must be half an interval old to count. + * + * @param monitors Supplies the monitors to check, by PV + * @param probes Starts subscriptions in the independent context + * @param clock The clock + * @param grace How long a monitor may be disconnected before it's checked + * @param interval How often check is called + * @param maxProbes The most probes open at once + */ + public FrozenPvDetector( + Supplier> monitors, + ProbeFactory probes, + Clock clock, + Duration grace, + Duration interval, + int maxProbes) { + this.monitors = monitors; + this.probes = probes; + this.clock = clock; + this.grace = grace; + this.quiet = interval.multipliedBy(6); + this.window = interval.multipliedBy(6); + this.margin = interval.dividedBy(2); + this.reprobe = interval.multipliedBy(30); + this.maxProbes = maxProbes; + } + + /** + * Views of the given monitors, for the detector. + * + * @param monitorMap The monitors, by PV + * @return The views, by PV + */ + public static Map viewsOf(Map monitorMap) { + Map views = new HashMap<>(); + for (Map.Entry entry : monitorMap.entrySet()) { + ChannelMonitor monitor = entry.getValue(); + views.put( + entry.getKey(), + new MonitorView( + monitor.getState(), + monitor.getStateChanged(), + monitor.getLastTimestamp() == null ? null : monitor.getLastTimestamp().toInstant(), + monitor.getUpdateCount())); + } + return views; + } + + /** Run one check: release recovered PVs, evaluate open probes, and start new ones. */ + public synchronized void check() { + Instant now = clock.instant(); + Map views = monitors.get(); + + releaseRecovered(views, now); + evaluateProbes(views, now); + startProbes(views, now); + + lastProbed.keySet().retainAll(views.keySet()); + } + + /** + * @return The frozen PVs, by PV + */ + public Map getFrozen() { + return Map.copyOf(frozen); + } + + @Override + public synchronized void close() { + for (ProbeSession session : active.values()) { + session.probe().close(); + } + active.clear(); + } + + private void releaseRecovered(Map views, Instant now) { + for (Iterator it = frozen.keySet().iterator(); it.hasNext(); ) { + String pv = it.next(); + MonitorView view = views.get(pv); + if (view == null) { // No longer monitored + updateCountWhenFrozen.remove(pv); + it.remove(); + } else if (view.updateCount() > updateCountWhenFrozen.get(pv)) { + LOGGER.log( + Level.INFO, + "PV {0} recovered after being frozen for {1} s", + new Object[] {pv, Duration.between(frozen.get(pv).since(), now).toSeconds()}); + updateCountWhenFrozen.remove(pv); + it.remove(); + } + } + } + + private void evaluateProbes(Map views, Instant now) { + for (Iterator> it = active.entrySet().iterator(); + it.hasNext(); ) { + Map.Entry entry = it.next(); + String pv = entry.getKey(); + ProbeSession session = entry.getValue(); + Probe probe = session.probe(); + MonitorView view = views.get(pv); + + probe.poll(); + + String reason = null; + boolean done; + if (view == null || view.updateCount() > session.updateCountAtStart()) { + done = true; // No longer monitored, or the monitor is getting updates + } else if (session.kind() == Kind.NOT_CONNECTED + && view.state() != ChannelMonitor.MonitorState.CONNECTED + && isOlderThanMargin(probe.connectedAt(), now)) { + reason = + "The IOC serves it, but its monitor has been " + + view.state() + + " since " + + view.stateChanged(); + done = true; + } else if (session.kind() == Kind.QUIET && isOlderThanMargin(probe.firstChangeAt(), now)) { + reason = + "The IOC sent a change its monitor didn't receive; last update " + view.lastUpdate(); + done = true; + } else { + done = !now.isBefore(session.started().plus(window)); // Inconclusive + } + + if (reason != null) { + LOGGER.log(Level.WARNING, "PV {0} is frozen: {1}", new Object[] {pv, reason}); + frozen.put(pv, new FrozenPv(now, reason)); + updateCountWhenFrozen.put(pv, view.updateCount()); + } + if (done) { + probe.close(); + it.remove(); + } + } + } + + private void startProbes(Map views, Instant now) { + List> candidates = new ArrayList<>(); + for (Map.Entry entry : views.entrySet()) { + String pv = entry.getKey(); + Instant probed = lastProbed.get(pv); + if (frozen.containsKey(pv) + || active.containsKey(pv) + || (probed != null && now.isBefore(probed.plus(reprobe)))) { + continue; + } + Kind kind = kindOf(entry.getValue(), now); + if (kind != null) { + candidates.add(Map.entry(pv, kind)); + } + } + + // Least recently probed first, so every candidate gets its turn + candidates.sort(Comparator.comparing(c -> lastProbed.getOrDefault(c.getKey(), Instant.MIN))); + + for (Map.Entry candidate : candidates) { + if (active.size() >= maxProbes) { + break; + } + String pv = candidate.getKey(); + lastProbed.put(pv, now); + try { + Probe probe = probes.start(pv); + active.put( + pv, new ProbeSession(probe, candidate.getValue(), now, views.get(pv).updateCount())); + } catch (Exception e) { + LOGGER.log(Level.FINE, "Unable to probe " + pv, e); + } + } + } + + private Kind kindOf(MonitorView view, Instant now) { + if (view.state() != ChannelMonitor.MonitorState.CONNECTED) { + return now.isBefore(view.stateChanged().plus(grace)) ? null : Kind.NOT_CONNECTED; + } + // Only PVs that have changed since subscribing can show updates going missing + if (view.updateCount() >= 2 + && view.lastUpdate() != null + && !now.isBefore(view.lastUpdate().plus(quiet))) { + return Kind.QUIET; + } + return null; + } + + private boolean isOlderThanMargin(Instant time, Instant now) { + return time != null && !now.isBefore(time.plus(margin)); + } +} diff --git a/src/test/java/org/jlab/epics2web/epics/FrozenPvDetectorTest.java b/src/test/java/org/jlab/epics2web/epics/FrozenPvDetectorTest.java new file mode 100644 index 0000000..08930da --- /dev/null +++ b/src/test/java/org/jlab/epics2web/epics/FrozenPvDetectorTest.java @@ -0,0 +1,250 @@ +package org.jlab.epics2web.epics; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +import java.time.Clock; +import java.time.Duration; +import java.time.Instant; +import java.time.ZoneId; +import java.time.ZoneOffset; +import java.util.HashMap; +import java.util.Map; +import java.util.Set; +import org.jlab.epics2web.epics.ChannelMonitor.MonitorState; +import org.jlab.epics2web.epics.FrozenPvDetector.MonitorView; +import org.junit.Test; + +/** + * Tests FrozenPvDetector's decisions with a controlled clock, monitors and probes. With a 10 s + * interval: quiet after 60 s, probes last 60 s, re-probe after 300 s, evidence counts after 5 s. + */ +public class FrozenPvDetectorTest { + + private static final Duration GRACE = Duration.ofSeconds(30); + private static final Duration INTERVAL = Duration.ofSeconds(10); + + private final MutableClock clock = new MutableClock(); + private final Map monitors = new HashMap<>(); + private final Map probes = new HashMap<>(); + private final FrozenPvDetector detector = + new FrozenPvDetector(() -> Map.copyOf(monitors), this::startProbe, clock, GRACE, INTERVAL, 2); + + /** A PV that used to update, went quiet, and changed for the probe but not the monitor. */ + @Test + public void quietPvThatChangedForProbeIsFrozen() { + connected("pv1", 5); + advance(60); + detector.check(); + assertEquals(Set.of("pv1"), probes.keySet()); + + probes.get("pv1").connectedAt = clock.instant(); + probes.get("pv1").firstChangeAt = clock.instant(); + advance(4); // Evidence not yet 5 s old: the monitor may still be about to get it + detector.check(); + assertTrue(detector.getFrozen().isEmpty()); + + advance(1); + detector.check(); + assertEquals(Set.of("pv1"), detector.getFrozen().keySet()); + assertTrue(probes.get("pv1").closed); + } + + @Test + public void frozenPvRecoversWhenItsMonitorGetsAnUpdate() { + quietPvThatChangedForProbeIsFrozen(); + + advance(10); + connected("pv1", 6); + detector.check(); + + assertTrue(detector.getFrozen().isEmpty()); + } + + /** A quiet PV that doesn't change for the probe either just isn't changing. */ + @Test + public void quietPvThatDoesNotChangeIsNotFrozen() { + connected("pv1", 5); + advance(60); + detector.check(); + probes.get("pv1").connectedAt = clock.instant(); + + advance(60); + detector.check(); + + assertTrue(detector.getFrozen().isEmpty()); + assertTrue("Inconclusive probe left open", probes.get("pv1").closed); + + probes.clear(); + advance(200); // 260 s since it was probed + detector.check(); + assertTrue("Probed again too soon", probes.isEmpty()); + + advance(40); + detector.check(); + assertEquals(Set.of("pv1"), probes.keySet()); + } + + @Test + public void monitorThatGetsAnUpdateDuringProbeIsNotFrozen() { + connected("pv1", 5); + advance(60); + detector.check(); + + connected("pv1", 6); + probes.get("pv1").connectedAt = clock.instant(); + probes.get("pv1").firstChangeAt = clock.instant(); + advance(10); + detector.check(); + + assertTrue(detector.getFrozen().isEmpty()); + assertTrue(probes.get("pv1").closed); + } + + /** Only PVs that changed since subscribing can show updates going missing. */ + @Test + public void pvThatNeverChangedIsNotProbed() { + connected("pv1", 1); + advance(600); + detector.check(); + + assertTrue(probes.isEmpty()); + } + + @Test + public void disconnectedPvTheIocServesIsFrozen() { + monitors.put("pv1", new MonitorView(MonitorState.DISCONNECTED, clock.instant(), null, 3)); + advance(29); + detector.check(); + assertTrue("Probed within the grace period", probes.isEmpty()); + + advance(1); + detector.check(); + probes.get("pv1").connectedAt = clock.instant(); + advance(5); + detector.check(); + + assertEquals(Set.of("pv1"), detector.getFrozen().keySet()); + } + + @Test + public void neverConnectedPvTheIocServesIsFrozen() { + monitors.put("pv1", new MonitorView(MonitorState.CONNECTING, clock.instant(), null, 0)); + advance(30); + detector.check(); + probes.get("pv1").connectedAt = clock.instant(); + advance(5); + detector.check(); + + assertEquals(Set.of("pv1"), detector.getFrozen().keySet()); + } + + /** If the probe can't reach the PV either, the IOC is down or the PV doesn't exist. */ + @Test + public void disconnectedPvTheIocDoesNotServeIsNotFrozen() { + monitors.put("pv1", new MonitorView(MonitorState.DISCONNECTED, clock.instant(), null, 3)); + advance(30); + detector.check(); + + advance(60); + detector.check(); + + assertTrue(detector.getFrozen().isEmpty()); + assertTrue(probes.get("pv1").closed); + } + + @Test + public void atMostMaxProbesAtOnce() { + connected("pv1", 5); + connected("pv2", 5); + connected("pv3", 5); + advance(60); + + detector.check(); + + assertEquals(2, probes.size()); + } + + @Test + public void pvNoLongerMonitoredIsDropped() { + quietPvThatChangedForProbeIsFrozen(); + connected("pv2", 5); + advance(60); + detector.check(); // Starts probing pv2 + + monitors.clear(); + detector.check(); + + assertTrue(detector.getFrozen().isEmpty()); + assertTrue(probes.get("pv2").closed); + } + + @Test + public void closeClosesOpenProbes() { + connected("pv1", 5); + advance(60); + detector.check(); + + detector.close(); + + assertTrue(probes.get("pv1").closed); + assertFalse(detector.getFrozen().containsKey("pv1")); + } + + /** A connected monitor whose latest update is now. */ + private void connected(String pv, long updateCount) { + monitors.put( + pv, new MonitorView(MonitorState.CONNECTED, Instant.EPOCH, clock.instant(), updateCount)); + } + + private void advance(long seconds) { + clock.now = clock.now.plusSeconds(seconds); + } + + private FrozenPvDetector.Probe startProbe(String pv) { + FakeProbe probe = new FakeProbe(); + probes.put(pv, probe); + return probe; + } + + private static class FakeProbe implements FrozenPvDetector.Probe { + Instant connectedAt; + Instant firstChangeAt; + boolean closed; + + @Override + public Instant connectedAt() { + return connectedAt; + } + + @Override + public Instant firstChangeAt() { + return firstChangeAt; + } + + @Override + public void close() { + closed = true; + } + } + + private static class MutableClock extends Clock { + Instant now = Instant.parse("2026-10-04T12:00:00Z"); + + @Override + public ZoneId getZone() { + return ZoneOffset.UTC; + } + + @Override + public Clock withZone(ZoneId zone) { + return this; + } + + @Override + public Instant instant() { + return now; + } + } +}