diff --git a/configurations/default/env.yml.tmp b/configurations/default/env.yml.tmp
index 91e53f92b..9c45259db 100644
--- a/configurations/default/env.yml.tmp
+++ b/configurations/default/env.yml.tmp
@@ -17,3 +17,7 @@ GTFS_DATABASE_URL: jdbc:postgresql://localhost/catalogue # If running via docker
# GTFS_DATABASE_PASSWORD:
#MONGO_HOST: mongo-host:27017 # If running via docker, this is mongo:27017
MONGO_DB_NAME: catalogue
+
+# If set, this API key must be passed to access the /metrics endpoint.
+# Pass via header "X-API-Key" or query parameter "api_key".
+# METRICS_API_KEY: your-secret-api-key
\ No newline at end of file
diff --git a/configurations/default/server.yml.tmp b/configurations/default/server.yml.tmp
index 74b0b2338..a860628a8 100644
--- a/configurations/default/server.yml.tmp
+++ b/configurations/default/server.yml.tmp
@@ -16,6 +16,8 @@ application:
s3_region: us-east-1
gtfs_s3_bucket: bucket-name
modules:
+ metrics:
+ enabled: true
enterprise:
enabled: false
# Setting this to true will upload all feeds to S3 instead of linking to their URL
diff --git a/configurations/test/env.yml.tmp b/configurations/test/env.yml.tmp
index ee8a12604..07e7bffdc 100644
--- a/configurations/test/env.yml.tmp
+++ b/configurations/test/env.yml.tmp
@@ -24,3 +24,7 @@ MONGO_DB_NAME: catalogue
#MONGO_PASSWORD: password
#MONGO_PROTOCOL: mongodb+srv
#MONGO_USER: user
+
+# If set, this API key must be passed to access the /metrics endpoint.
+# Pass via header "X-API-Key" or query parameter "api_key".
+# METRICS_API_KEY: your-secret-api-key
\ No newline at end of file
diff --git a/configurations/test/server.yml.tmp b/configurations/test/server.yml.tmp
index c24f035fa..f88d58488 100644
--- a/configurations/test/server.yml.tmp
+++ b/configurations/test/server.yml.tmp
@@ -15,6 +15,8 @@ application:
s3_region: us-east-1
gtfs_s3_bucket: bucket-name
modules:
+ metrics:
+ enabled: true
enterprise:
enabled: false
editor:
diff --git a/pom.xml b/pom.xml
index d0be65e44..b7020e991 100644
--- a/pom.xml
+++ b/pom.xml
@@ -210,6 +210,18 @@
+
+
+
+
+ io.micrometer
+ micrometer-bom
+ 1.17.0
+ pom
+ import
+
+
+
@@ -231,6 +243,12 @@
1.2.13
+
+
+ io.micrometer
+ micrometer-registry-prometheus
+
+
org.mapdb
diff --git a/src/main/java/com/conveyal/datatools/common/status/MonitorableJob.java b/src/main/java/com/conveyal/datatools/common/status/MonitorableJob.java
index 71b7a93f7..35a110b4a 100644
--- a/src/main/java/com/conveyal/datatools/common/status/MonitorableJob.java
+++ b/src/main/java/com/conveyal/datatools/common/status/MonitorableJob.java
@@ -1,6 +1,7 @@
package com.conveyal.datatools.common.status;
import com.conveyal.datatools.manager.auth.Auth0UserProfile;
+import com.conveyal.datatools.manager.metrics.MetricsService;
import com.conveyal.datatools.manager.utils.JobUtils;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
@@ -215,6 +216,7 @@ public void run () {
} finally {
LOG.info("{} (jobId={}) {} in {} ms", type, jobId, status.error ? "errored" : "completed", status.duration);
active = false;
+ MetricsService.recordJobOutcome(type, status.error, status.duration);
}
}
diff --git a/src/main/java/com/conveyal/datatools/common/utils/Scheduler.java b/src/main/java/com/conveyal/datatools/common/utils/Scheduler.java
index b468ca89c..e46667f36 100644
--- a/src/main/java/com/conveyal/datatools/common/utils/Scheduler.java
+++ b/src/main/java/com/conveyal/datatools/common/utils/Scheduler.java
@@ -1,8 +1,10 @@
package com.conveyal.datatools.common.utils;
+import com.conveyal.datatools.manager.DataManager;
import com.conveyal.datatools.manager.auth.Auth0UserProfile;
import com.conveyal.datatools.manager.jobs.FeedExpirationNotificationJob;
import com.conveyal.datatools.manager.jobs.FetchSingleFeedJob;
+import com.conveyal.datatools.manager.metrics.MetricsService;
import com.conveyal.datatools.manager.models.FeedSource;
import com.conveyal.datatools.manager.models.FeedVersion;
import com.conveyal.datatools.manager.models.Project;
@@ -48,10 +50,22 @@ public class Scheduler {
public final static ListMultimap scheduledJobsForFeedSources =
synchronizedListMultimap(ArrayListMultimap.create());
+ // Use a separate thread to handle updating metrics inventory.
+ private static final ScheduledExecutorService metricsInventoryScheduler =
+ Executors.newSingleThreadScheduledExecutor(r -> {
+ Thread t = new Thread(r, "metrics-refresh");
+ t.setDaemon(true);
+ return t;
+ });
+
+
/**
* A method to initialize all scheduled tasks upon server startup.
*/
public static void initialize() {
+ if (DataManager.isModuleEnabled("metrics")) {
+ startInventoryRefresh(5, TimeUnit.MINUTES);
+ }
LOG.info("Scheduling recurring feed auto fetches for all projects.");
for (Project project : Persistence.projects.getAll()) {
handleAutoFeedFetch(project);
@@ -64,6 +78,17 @@ public static void initialize() {
}
}
+ /**
+ * Schedules a task to refresh the metrics inventory on a regular interval
+ * @param period refresh interval
+ * @param unit unit for the refresh interval
+ */
+ public static void startInventoryRefresh(long period, TimeUnit unit) {
+ LOG.info("Creating a separate thread for refreshing metrics inventory.");
+ metricsInventoryScheduler.scheduleAtFixedRate(MetricsService::refreshInventory,
+ 0, period, unit);
+ }
+
/**
* Convenience method for scheduling one-off jobs for a feed source.
*/
diff --git a/src/main/java/com/conveyal/datatools/manager/DataManager.java b/src/main/java/com/conveyal/datatools/manager/DataManager.java
index db1af1c90..f437c258d 100644
--- a/src/main/java/com/conveyal/datatools/manager/DataManager.java
+++ b/src/main/java/com/conveyal/datatools/manager/DataManager.java
@@ -15,6 +15,7 @@
import com.conveyal.datatools.manager.controllers.api.FeedVersionController;
import com.conveyal.datatools.manager.controllers.api.GtfsPlusController;
import com.conveyal.datatools.manager.controllers.api.LabelController;
+import com.conveyal.datatools.manager.controllers.api.MetricsController;
import com.conveyal.datatools.manager.controllers.api.NoteController;
import com.conveyal.datatools.manager.controllers.api.OrganizationController;
import com.conveyal.datatools.manager.controllers.api.ProjectController;
@@ -70,6 +71,8 @@ public class DataManager {
public static final String GTFS_PLUS_SUBDIR = "gtfsplus";
private static final Logger LOG = LoggerFactory.getLogger(DataManager.class);
+ public static long serverStartTime;
+
// These fields hold YAML files that represent the server configuration.
private static JsonNode envConfig;
private static JsonNode serverConfig;
@@ -102,7 +105,7 @@ public class DataManager {
public static final Map lastRequestForUser = new HashMap<>();
public static void main(String[] args) throws IOException {
- long serverStartTime = System.currentTimeMillis();
+ serverStartTime = System.currentTimeMillis();
initializeApplication(args);
registerRoutes();
@@ -188,7 +191,11 @@ private static void loadProperties() {
* modules and sets other core routes (e.g., 404 response) and response headers (e.g., API content type is JSON).
*/
static void registerRoutes() throws IOException {
+
CorsFilter.apply();
+ if (isModuleEnabled("metrics")) {
+ MetricsController.register();
+ }
// Initialize GTFS GraphQL API service
// FIXME: Add user permissions check to ensure user has access to feeds.
GraphQLController.initialize(GTFS_DATA_SOURCE, GTFS_API_PREFIX);
diff --git a/src/main/java/com/conveyal/datatools/manager/controllers/api/MetricsController.java b/src/main/java/com/conveyal/datatools/manager/controllers/api/MetricsController.java
new file mode 100644
index 000000000..79171e441
--- /dev/null
+++ b/src/main/java/com/conveyal/datatools/manager/controllers/api/MetricsController.java
@@ -0,0 +1,34 @@
+package com.conveyal.datatools.manager.controllers.api;
+
+import com.conveyal.datatools.manager.metrics.MetricsService;
+
+import spark.Request;
+import spark.Response;
+
+import static com.conveyal.datatools.common.utils.SparkUtils.logMessageAndHalt;
+import static com.conveyal.datatools.manager.DataManager.getConfigPropertyAsText;
+import static spark.Spark.get;
+
+public class MetricsController {
+ private static final String METRICS_API_KEY_CONFIG = "METRICS_API_KEY";
+
+ /** GET /metrics — Prometheus scrape endpoint. */
+ private static String getMetrics(Request req, Response res) {
+ String apiKey = getConfigPropertyAsText(METRICS_API_KEY_CONFIG);
+ if (apiKey != null && !apiKey.isEmpty()) {
+ String providedKey = req.headers("X-API-Key");
+ if (providedKey == null) {
+ providedKey = req.queryParams("api_key");
+ }
+ if (!apiKey.equals(providedKey)) {
+ logMessageAndHalt(req, 401, "Invalid or missing API key");
+ }
+ }
+ res.type("text/plain; version=0.0.4; charset=utf-8");
+ return MetricsService.registry().scrape();
+ }
+
+ public static void register() {
+ get("/metrics", MetricsController::getMetrics);
+ }
+}
diff --git a/src/main/java/com/conveyal/datatools/manager/metrics/MetricsService.java b/src/main/java/com/conveyal/datatools/manager/metrics/MetricsService.java
new file mode 100644
index 000000000..d6641de85
--- /dev/null
+++ b/src/main/java/com/conveyal/datatools/manager/metrics/MetricsService.java
@@ -0,0 +1,140 @@
+package com.conveyal.datatools.manager.metrics;
+
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicLong;
+
+import org.bson.Document;
+
+import com.conveyal.datatools.common.status.MonitorableJob.JobType;
+import com.conveyal.datatools.manager.DataManager;
+import com.conveyal.datatools.manager.persistence.Persistence;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.micrometer.core.instrument.Counter;
+import io.micrometer.core.instrument.Gauge;
+import io.micrometer.core.instrument.Timer;
+import io.micrometer.core.instrument.binder.jvm.ClassLoaderMetrics;
+import io.micrometer.core.instrument.binder.jvm.JvmGcMetrics;
+import io.micrometer.core.instrument.binder.jvm.JvmMemoryMetrics;
+import io.micrometer.core.instrument.binder.jvm.JvmThreadDeadlockMetrics;
+import io.micrometer.core.instrument.binder.jvm.JvmThreadMetrics;
+import io.micrometer.core.instrument.binder.system.ProcessorMetrics;
+import io.micrometer.prometheusmetrics.PrometheusConfig;
+import io.micrometer.prometheusmetrics.PrometheusMeterRegistry;
+
+public class MetricsService {
+ private static final Logger LOG = LoggerFactory.getLogger(MetricsService.class);
+
+ private static final PrometheusMeterRegistry registry;
+
+ private static final AtomicLong feedSourceCount = new AtomicLong(0);
+ private static final AtomicLong feedVersionCount = new AtomicLong(0);
+ private static final AtomicLong projectCount = new AtomicLong(0);
+ private static final AtomicLong organizationCount = new AtomicLong(0);
+ private static final AtomicLong lastRefreshEpoch = new AtomicLong(0);
+ private static final Counter refreshFailures;
+
+ private static final ConcurrentMap failedJobsByType = new ConcurrentHashMap<>();
+ private static final ConcurrentMap completedJobsByType = new ConcurrentHashMap<>();
+ private static final ConcurrentMap jobDurationByType = new ConcurrentHashMap<>();
+ private static Timer jobDuration;
+
+ static {
+ registry = new PrometheusMeterRegistry(PrometheusConfig.DEFAULT);
+ populateAndRegisterJobMetricsByType();
+
+ new ClassLoaderMetrics().bindTo(registry);
+ new JvmMemoryMetrics().bindTo(registry);
+ new JvmGcMetrics().bindTo(registry);
+ new ProcessorMetrics().bindTo(registry);
+ new JvmThreadMetrics().bindTo(registry);
+ new JvmThreadDeadlockMetrics().bindTo(registry);
+
+ Gauge.builder("datatools.feed.sources", feedSourceCount, AtomicLong::get)
+ .description("Total number of Feed Sources")
+ .strongReference(true)
+ .register(registry);
+ Gauge.builder("datatools.feed.versions", feedVersionCount, AtomicLong::get)
+ .description("Total number of Feed Versions")
+ .strongReference(true)
+ .register(registry);
+ Gauge.builder("datatools.projects", projectCount, AtomicLong::get)
+ .description("Total number of Projects")
+ .strongReference(true)
+ .register(registry);
+ Gauge.builder("datatools.organizations", organizationCount, AtomicLong::get)
+ .description("Total number of Organizations")
+ .strongReference(true)
+ .register(registry);
+ Gauge.builder("datatools.inventory.refresh.age.seconds",
+ lastRefreshEpoch, e -> e.get() == 0
+ ? Double.NaN
+ : (System.currentTimeMillis() - e.get()) / 1000.0)
+ .description("Seconds since last inventory refresh")
+ .strongReference(true)
+ .register(registry);
+ refreshFailures = Counter.builder("datatools.inventory.refresh.failures")
+ .description("Total number of failed Refreshes")
+ .register(registry);
+ Gauge.builder("datatools.uptime.seconds",
+ DataManager.serverStartTime,
+ start -> (System.currentTimeMillis() - start) / 1000.0)
+ .description("Server uptime in seconds")
+ .register(registry);
+ }
+
+ private static void populateAndRegisterJobMetricsByType() {
+ for(JobType type : JobType.values()) {
+ failedJobsByType.put(type, Counter.builder(String.format("datatools.jobs.%s.failed", type.name()))
+ .description(String.format("Number of failed jobs of type %s", type.name()))
+ .register(registry));
+ completedJobsByType.put(type, Counter.builder(String.format("datatools.jobs.%s.completed", type.name()))
+ .description(String.format("Number of completed jobs of type %s", type.name()))
+ .register(registry));
+ jobDurationByType.put(type, Timer.builder(String.format("datatools.jobs.%s.duration", type.name()))
+ .description(String.format("Execution duration of jobs of type %s", type.name()))
+ .register(registry));
+ }
+
+ Gauge.builder("datatools.jobs.failed", failedJobsByType,
+ map -> map.values().stream().mapToDouble(Counter::count).sum())
+ .description("Number of failed jobs")
+ .register(registry);
+ Gauge.builder("datatools.jobs.completed", completedJobsByType,
+ map -> map.values().stream().mapToDouble(Counter::count).sum())
+ .description("Number of completed jobs")
+ .register(registry);
+ jobDuration = Timer.builder("datatools.jobs.duration")
+ .description("Job execution duration")
+ .register(registry);
+ }
+
+ public static PrometheusMeterRegistry registry() {
+ return registry;
+ }
+
+ public static void refreshInventory() {
+ try {
+ Document allDocuments = new Document();
+ feedSourceCount.set(Persistence.feedSources.count(allDocuments));
+ feedVersionCount.set(Persistence.feedVersions.count(allDocuments));
+ projectCount.set(Persistence.projects.count(allDocuments));
+ organizationCount.set(Persistence.organizations.count(allDocuments));
+ lastRefreshEpoch.set(System.currentTimeMillis());
+ } catch (Exception e) {
+ refreshFailures.increment();
+ LOG.warn("Inventory refresh failed", e);
+ }
+ }
+
+ public static void recordJobOutcome(JobType type, boolean error, long duration) {
+ ConcurrentMap outcomeCountersByType = error ? failedJobsByType : completedJobsByType;
+
+ outcomeCountersByType.get(type).increment();
+ jobDuration.record(duration, TimeUnit.MILLISECONDS);
+ jobDurationByType.get(type).record(duration, TimeUnit.MILLISECONDS);
+ }
+}