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
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
import com.juliashtal.devanalytics.ai.model.MetricsSummaryDto;
import com.juliashtal.devanalytics.ai.service.MetricsAiService;
import com.juliashtal.devanalytics.config.SystemClock;
import com.juliashtal.devanalytics.metrics.service.MetricWriteGate;
import com.juliashtal.devanalytics.metrics.service.MetricsService;
import com.juliashtal.devanalytics.notification.NotificationDispatchService;
import com.juliashtal.devanalytics.user.repository.UserRepository;
Expand All @@ -28,6 +29,7 @@ public class MetricsSummaryScheduler {
private final MetricsAiService metricsAiService;
private final NotificationDispatchService notificationDispatch;
private final SystemClock systemClock;
private final MetricWriteGate writeGate;

@Scheduled(cron = "0 0 8 * * MON", zone = "UTC")
public void generateWeeklySummaries() {
Expand All @@ -39,7 +41,13 @@ public void generateWeeklySummaries() {
userRepository.findAll().forEach(user -> {
try {
// Compute the week before summarising it: the nightly job may not have covered this ISO week.
metricsService.calculateDailyMetrics(user.getId(), from, to);
// Only the refresh is gated — a summary over stored snapshots still beats no summary
// for a week, which is what skipping the whole run would cost at this cadence.
boolean refreshed = writeGate.runExclusively(
() -> metricsService.calculateDailyMetrics(user.getId(), from, to));
if (!refreshed) {
log.info("Metric refresh skipped for userId={}; summarising stored snapshots", user.getId());
}

MetricsSummaryDto summary = metricsAiService.generateSummary(user, from, to, null);
log.debug("Generated weekly summary for userId={}", user.getId());
Expand Down
Original file line number Diff line number Diff line change
@@ -1,22 +1,30 @@
package com.juliashtal.devanalytics.config;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.converter.json.MappingJackson2HttpMessageConverter;
import org.springframework.http.client.SimpleClientHttpRequestFactory;
import org.springframework.web.client.RestTemplate;

/**
* Provides the shared RestTemplate used for outbound HTTP calls.
*
* <p>Timeouts bound each hop because the default factory waits forever, and
* {@code JiraProjectService.discoverProjectsFromJira} runs on the request thread — an unreachable
* Jira host would otherwise hold a Tomcat worker until the client gave up. They bound a single
* request, not a paged loop over many.</p>
*/
@Configuration
public class RestTemplateConfig {

@Bean
public RestTemplate restTemplate() {
RestTemplate restTemplate = new RestTemplate();
restTemplate.getMessageConverters().add(new MappingJackson2HttpMessageConverter());
return restTemplate;
public RestTemplate restTemplate(
@Value("${app.http.connect-timeout-ms}") int connectTimeoutMs,
@Value("${app.http.read-timeout-ms}") int readTimeoutMs) {
SimpleClientHttpRequestFactory factory = new SimpleClientHttpRequestFactory();
factory.setConnectTimeout(connectTimeoutMs);
factory.setReadTimeout(readTimeoutMs);
return new RestTemplate(factory);
}

}

Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import com.juliashtal.devanalytics.security.service.CustomUserDetailsService;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.HttpMethod;
import org.springframework.http.HttpStatus;
import org.springframework.security.authentication.AuthenticationManager;
import org.springframework.security.authentication.AuthenticationProvider;
Expand Down Expand Up @@ -99,6 +100,10 @@ public SecurityFilterChain securityFilterChain(HttpSecurity http) throws Excepti
.authorizeHttpRequests(auth -> auth
.requestMatchers(PUBLIC_ENDPOINTS).permitAll()
.requestMatchers("/api/admin/**").hasRole("ADMIN")
// Ordered before the team rule below: the first matching pattern wins. Named
// in full rather than as a subtree — "me" is a literal where every sibling
// route takes a {teamId}, so a wildcard here would shadow all of them.
.requestMatchers(HttpMethod.GET, "/api/teams/me/memberships").authenticated()
.requestMatchers("/api/teams/**").hasAnyRole("MANAGER", "ADMIN")
.anyRequest().authenticated()
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,9 @@
* Drives {@link MetricBackfillService} over every user once a night, filling days inside each
* user's collected history that have never been calculated.
*
* <p>Separate from {@link MetricsScheduler}, which only computes yesterday. Runs at 03:00 UTC;
* ordering against that job is not a correctness requirement, because both write through the
* same {@link MetricsService#calculateDailyMetrics} upsert guard.</p>
* <p>Separate from {@link MetricsScheduler}, which only computes yesterday. Runs at 03:00 UTC and
* takes {@link MetricWriteGate} first: the upsert guard both jobs share reads before it writes, so
* it orders writes but does not make concurrent ones safe.</p>
*/
@Component
@RequiredArgsConstructor
Expand All @@ -22,9 +22,17 @@ public class MetricBackfillScheduler {

private final UserRepository userRepository;
private final MetricBackfillService backfillService;
private final MetricWriteGate writeGate;

@Scheduled(cron = "0 0 3 * * ?", zone = "UTC")
public void backfillAll() {
boolean ran = writeGate.runExclusively(this::backfillAllUsers);
if (!ran) {
log.info("History backfill job skipped: another metric writer is running");
}
}

private void backfillAllUsers() {
log.info("History backfill job started");

userRepository.findAll().forEach(user -> {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
package com.juliashtal.devanalytics.metrics.service;

import org.springframework.stereotype.Component;

import java.util.concurrent.locks.ReentrantLock;

/**
* Serialises the scheduled jobs that write {@code metric_snapshots}; request threads, the
* {@code collect-} pool and the attribution listener still write outside it.
*
* <p>The table carries no unique key, so the guard in {@code MetricSnapshotWriter} is a
* read-then-write: two writers computing the same window both miss {@code findExisting} and both
* insert. A process-local lock narrows that to the scheduled writers; a unique index over the
* {@code findExisting} columns is what would close it for every caller.</p>
*/
@Component
public class MetricWriteGate {

private final ReentrantLock lock = new ReentrantLock();

/**
* Runs {@code job} only while no other writer holds the gate.
*
* <p>Skips rather than queues: each caller recomputes from persisted state, so the next run
* covers whatever this one declined.</p>
*
* @return false when the job was skipped because another writer was running
*/
public boolean runExclusively(Runnable job) {
if (!lock.tryLock()) {
return false;
}
try {
job.run();
return true;
} finally {
lock.unlock();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -24,10 +24,18 @@ public class MetricsScheduler {
private final MetricsService metricsService;
private final UserRepository userRepository;
private final SystemClock systemClock;
private final MetricWriteGate writeGate;

/** Runs every day at 01:00 UTC. */
@Scheduled(cron = "0 0 1 * * ?", zone = "UTC")
public void calculateYesterday() {
boolean ran = writeGate.runExclusively(this::calculateYesterdayForAllUsers);
if (!ran) {
log.info("Nightly metrics scheduler skipped: another metric writer is running");
}
}

private void calculateYesterdayForAllUsers() {
LocalDate yesterday = systemClock.yesterday();
log.info("Nightly metrics scheduler started for {}", yesterday);

Expand Down
10 changes: 10 additions & 0 deletions dev-analytics/src/main/resources/application.yml
Original file line number Diff line number Diff line change
@@ -1,4 +1,10 @@
spring:
task:
scheduling:
pool:
# Enough for the 2-minute enrichment pass to overrun without delaying the three
# nightly crons, which fire an hour apart. The default of 1 queues them behind it.
size: ${SCHEDULING_POOL_SIZE:4}
servlet:
multipart:
max-file-size: 2MB
Expand Down Expand Up @@ -41,6 +47,10 @@ app:
refresh-expiration: 604800000 # 7 days
password-reset:
expiration: 3600000 # 1 hour
http:
# Bounds one outbound hop on the shared RestTemplate (Jira discovery and collection).
connect-timeout-ms: ${HTTP_CONNECT_TIMEOUT_MS:5000}
read-timeout-ms: ${HTTP_READ_TIMEOUT_MS:30000}
base-url: ${APP_BASE_URL:http://localhost:8080}
frontend-url: ${APP_FRONTEND_URL:http://localhost:5173}
encryption:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import com.juliashtal.devanalytics.ai.scheduler.MetricsSummaryScheduler;
import com.juliashtal.devanalytics.ai.service.MetricsAiService;
import com.juliashtal.devanalytics.config.SystemClock;
import com.juliashtal.devanalytics.metrics.service.MetricWriteGate;
import com.juliashtal.devanalytics.metrics.service.MetricsService;
import com.juliashtal.devanalytics.notification.NotificationDispatchService;
import com.juliashtal.devanalytics.user.model.User;
Expand Down Expand Up @@ -61,7 +62,7 @@ void restoreZone() {
@BeforeEach
void setUp() {
scheduler = new MetricsSummaryScheduler(userRepository, metricsService, metricsAiService,
notificationDispatch, new SystemClock(Clock.fixed(FIXED, ZoneOffset.UTC)));
notificationDispatch, new SystemClock(Clock.fixed(FIXED, ZoneOffset.UTC)), new MetricWriteGate());
}

@Test
Expand Down Expand Up @@ -117,6 +118,27 @@ void generateWeeklySummaries_oneUsersCalculationFailing_doesNotStopTheRest() {
verify(metricsAiService, org.mockito.Mockito.never()).generateSummary(eq(failing), any(), any(), any());
}

/**
* A weekly job that skipped entirely on contention would leave the user with no brief for
* seven days, so only the refresh is allowed to be declined.
*/
@Test
void generateWeeklySummaries_anotherWriterHoldsGate_stillSummarisesStoredSnapshots() {
User user = user(1L);
MetricWriteGate busyGate = org.mockito.Mockito.mock(MetricWriteGate.class);
when(busyGate.runExclusively(any())).thenReturn(false);
scheduler = new MetricsSummaryScheduler(userRepository, metricsService, metricsAiService,
notificationDispatch, new SystemClock(Clock.fixed(FIXED, ZoneOffset.UTC)), busyGate);
when(userRepository.findAll()).thenReturn(List.of(user));
when(metricsAiService.generateSummary(any(), any(), any(), any()))
.thenReturn(MetricsSummaryDto.builder().headline("h").build());

scheduler.generateWeeklySummaries();

org.mockito.Mockito.verifyNoInteractions(metricsService);
verify(metricsAiService).generateSummary(user, FROM, TO, null);
}

/**
* The job fires at 08:00 UTC, an hour at which Auckland has already entered the next day and
* Los Angeles is still in the previous one; the summarised week must not move with either.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
package com.juliashtal.devanalytics.config;

import org.junit.jupiter.api.Test;
import org.springframework.http.client.ClientHttpRequestFactory;
import org.springframework.http.client.SimpleClientHttpRequestFactory;
import org.springframework.test.util.ReflectionTestUtils;
import org.springframework.web.client.RestTemplate;

import static org.assertj.core.api.Assertions.assertThat;

/**
* Pins that the configured timeouts reach the request factory, so the shared RestTemplate cannot
* wait forever on an unresponsive host.
*
* <p>The factory exposes no getters, so its fields are the only reachable evidence; a rename
* there fails this test loudly rather than silently restoring the unlimited default.</p>
*/
class RestTemplateConfigTest {

@Test
void restTemplate_configuredTimeouts_reachTheRequestFactory() {
RestTemplate restTemplate = new RestTemplateConfig().restTemplate(1_500, 9_000);

ClientHttpRequestFactory factory = restTemplate.getRequestFactory();
assertThat(factory).isInstanceOf(SimpleClientHttpRequestFactory.class);
assertThat(ReflectionTestUtils.getField(factory, "connectTimeout")).isEqualTo(1_500);
assertThat(ReflectionTestUtils.getField(factory, "readTimeout")).isEqualTo(9_000);
}

@Test
void restTemplate_applicationDefaults_areFiniteAndPositive() {
RestTemplate restTemplate = new RestTemplateConfig().restTemplate(5_000, 30_000);

ClientHttpRequestFactory factory = restTemplate.getRequestFactory();
assertThat((Integer) ReflectionTestUtils.getField(factory, "connectTimeout")).isPositive();
assertThat((Integer) ReflectionTestUtils.getField(factory, "readTimeout")).isPositive();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import com.juliashtal.devanalytics.metrics.model.BackfillResult;
import com.juliashtal.devanalytics.metrics.service.MetricBackfillScheduler;
import com.juliashtal.devanalytics.metrics.service.MetricWriteGate;
import com.juliashtal.devanalytics.metrics.service.MetricBackfillService;
import com.juliashtal.devanalytics.user.model.User;
import com.juliashtal.devanalytics.user.repository.UserRepository;
Expand Down Expand Up @@ -33,12 +34,23 @@ void backfillAll_everyUser_isBackfilled() {
when(backfillService.backfillUser(anyLong()))
.thenReturn(new BackfillResult(0, 0, LocalDate.now(), LocalDate.now()));

new MetricBackfillScheduler(userRepository, backfillService).backfillAll();
new MetricBackfillScheduler(userRepository, backfillService, new MetricWriteGate()).backfillAll();

verify(backfillService).backfillUser(1L);
verify(backfillService).backfillUser(2L);
}

/** Without this, removing the gate from the scheduler leaves every other test here green. */
@Test
void backfillAll_anotherWriterHoldsGate_backfillsNothing() {
MetricWriteGate busyGate = mock(MetricWriteGate.class);
when(busyGate.runExclusively(any())).thenReturn(false);

new MetricBackfillScheduler(userRepository, backfillService, busyGate).backfillAll();

verifyNoInteractions(backfillService, userRepository);
}

@Test
void backfillAll_oneUserFails_othersStillProcessed() {
when(userRepository.findAll())
Expand All @@ -49,7 +61,7 @@ void backfillAll_oneUserFails_othersStillProcessed() {
when(backfillService.backfillUser(3L))
.thenReturn(new BackfillResult(0, 0, LocalDate.now(), LocalDate.now()));

new MetricBackfillScheduler(userRepository, backfillService).backfillAll();
new MetricBackfillScheduler(userRepository, backfillService, new MetricWriteGate()).backfillAll();

verify(backfillService).backfillUser(1L);
verify(backfillService).backfillUser(3L);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
package com.juliashtal.devanalytics.metrics;

import com.juliashtal.devanalytics.metrics.service.MetricWriteGate;
import org.junit.jupiter.api.Test;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;

import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

/**
* Pins that only one metric writer runs at a time, and that a declined run does not execute.
*
* <p>Contention is produced with a real second thread rather than a re-entrant call: the gate
* holds a {@link java.util.concurrent.locks.ReentrantLock}, so a same-thread attempt would be
* admitted and the test would pass without proving anything.</p>
*/
class MetricWriteGateTest {

@Test
void runExclusively_noContention_runsJobAndReturnsTrue() {
MetricWriteGate gate = new MetricWriteGate();
AtomicBoolean ran = new AtomicBoolean(false);

boolean result = gate.runExclusively(() -> ran.set(true));

assertThat(result).isTrue();
assertThat(ran).isTrue();
}

@Test
void runExclusively_otherWriterHolding_skipsJobAndReturnsFalse() throws Exception {
MetricWriteGate gate = new MetricWriteGate();
CountDownLatch holding = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
AtomicBoolean secondJobRan = new AtomicBoolean(false);

Thread holder = new Thread(() -> gate.runExclusively(() -> {
holding.countDown();
try {
release.await(5, TimeUnit.SECONDS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}));
holder.start();
assertThat(holding.await(5, TimeUnit.SECONDS)).isTrue();

boolean result = gate.runExclusively(() -> secondJobRan.set(true));

release.countDown();
holder.join(5_000);

assertThat(result).isFalse();
assertThat(secondJobRan).isFalse();
}

@Test
void runExclusively_jobThrows_releasesGateForTheNextRun() {
MetricWriteGate gate = new MetricWriteGate();

assertThatThrownBy(() -> gate.runExclusively(() -> {
throw new IllegalStateException("boom");
})).isInstanceOf(IllegalStateException.class);

AtomicBoolean ran = new AtomicBoolean(false);
assertThat(gate.runExclusively(() -> ran.set(true))).isTrue();
assertThat(ran).isTrue();
}
}
Loading
Loading