378 lines
16 KiB
Java
378 lines
16 KiB
Java
package com.magistr.app.service;
|
||
|
||
import com.magistr.app.dto.ScheduleRuleDto;
|
||
import com.magistr.app.dto.ScheduleRuleSlotDto;
|
||
import com.magistr.app.model.ScheduleParity;
|
||
import com.magistr.app.model.ScheduleRule;
|
||
import com.magistr.app.repository.ScheduleRuleRepository;
|
||
import org.junit.jupiter.api.AfterEach;
|
||
import org.junit.jupiter.api.Test;
|
||
import org.springframework.aop.support.AopUtils;
|
||
import org.springframework.beans.factory.annotation.Autowired;
|
||
import org.springframework.boot.SpringBootConfiguration;
|
||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||
import org.springframework.boot.autoconfigure.domain.EntityScan;
|
||
import org.springframework.boot.test.context.SpringBootTest;
|
||
import org.springframework.boot.test.context.TestConfiguration;
|
||
import org.springframework.context.annotation.Bean;
|
||
import org.springframework.context.annotation.Import;
|
||
import org.springframework.data.jpa.repository.config.EnableJpaRepositories;
|
||
import org.springframework.jdbc.core.JdbcTemplate;
|
||
import org.springframework.test.context.DynamicPropertyRegistry;
|
||
import org.springframework.test.context.DynamicPropertySource;
|
||
import org.testcontainers.containers.PostgreSQLContainer;
|
||
import org.testcontainers.junit.jupiter.Container;
|
||
import org.testcontainers.junit.jupiter.Testcontainers;
|
||
|
||
import javax.sql.DataSource;
|
||
import java.sql.Connection;
|
||
import java.sql.PreparedStatement;
|
||
import java.sql.ResultSet;
|
||
import java.sql.SQLException;
|
||
import java.util.ArrayList;
|
||
import java.util.List;
|
||
import java.util.concurrent.CountDownLatch;
|
||
import java.util.concurrent.ExecutorService;
|
||
import java.util.concurrent.Executors;
|
||
import java.util.concurrent.Future;
|
||
import java.util.concurrent.TimeUnit;
|
||
|
||
import static org.assertj.core.api.Assertions.assertThat;
|
||
import static org.mockito.Mockito.mock;
|
||
import static org.mockito.Mockito.times;
|
||
import static org.mockito.Mockito.verify;
|
||
|
||
@Testcontainers
|
||
@SpringBootTest(
|
||
classes = ScheduleRuleServiceConcurrencyIntegrationTest.TestApplication.class,
|
||
properties = {
|
||
"spring.jpa.open-in-view=false",
|
||
"management.endpoint.health.validate-group-membership=false"
|
||
}
|
||
)
|
||
class ScheduleRuleServiceConcurrencyIntegrationTest {
|
||
|
||
@Container
|
||
static final PostgreSQLContainer<?> POSTGRES =
|
||
new PostgreSQLContainer<>(com.magistr.app.testing.TestContainerImages.POSTGRES);
|
||
|
||
@DynamicPropertySource
|
||
static void databaseProperties(DynamicPropertyRegistry registry) {
|
||
registry.add("spring.datasource.url", POSTGRES::getJdbcUrl);
|
||
registry.add("spring.datasource.username", POSTGRES::getUsername);
|
||
registry.add("spring.datasource.password", POSTGRES::getPassword);
|
||
}
|
||
|
||
@Autowired
|
||
private ScheduleRuleService service;
|
||
|
||
@Autowired
|
||
private ScheduleGeneratorService scheduleGeneratorService;
|
||
|
||
@Autowired
|
||
private JdbcTemplate jdbcTemplate;
|
||
|
||
@Autowired
|
||
private DataSource dataSource;
|
||
|
||
private ExecutorService executor;
|
||
|
||
@AfterEach
|
||
void stopExecutor() throws InterruptedException {
|
||
if (executor == null) {
|
||
return;
|
||
}
|
||
executor.shutdownNow();
|
||
assertThat(executor.awaitTermination(5, TimeUnit.SECONDS)).isTrue();
|
||
}
|
||
|
||
@Test
|
||
void serializesConcurrentCreateInsideOneSemesterAcrossSpringTransactions() throws Exception {
|
||
assertThat(AopUtils.isAopProxy(service)).isTrue();
|
||
|
||
SeedIds ids = createIsolatedSemesterAndLoadSeedIds();
|
||
ScheduleRuleDto request = request(ids);
|
||
executor = Executors.newFixedThreadPool(2);
|
||
CountDownLatch ready = new CountDownLatch(2);
|
||
CountDownLatch start = new CountDownLatch(1);
|
||
List<Future<CreateOutcome>> futures = new ArrayList<>();
|
||
|
||
try (Connection lockConnection = dataSource.getConnection()) {
|
||
lockConnection.setAutoCommit(false);
|
||
int holderPid = lockSemesterRow(lockConnection, ids.semesterId());
|
||
futures.add(executor.submit(() -> createAfterSignal(request, ready, start)));
|
||
futures.add(executor.submit(() -> createAfterSignal(request, ready, start)));
|
||
awaitReady(ready);
|
||
start.countDown();
|
||
awaitTwoBlockedCreates(futures, holderPid);
|
||
lockConnection.commit();
|
||
}
|
||
|
||
List<CreateOutcome> outcomes = List.of(
|
||
futures.get(0).get(10, TimeUnit.SECONDS),
|
||
futures.get(1).get(10, TimeUnit.SECONDS)
|
||
);
|
||
|
||
assertThat(outcomes).filteredOn(outcome -> outcome.result() != null).hasSize(1);
|
||
assertThat(outcomes).filteredOn(outcome -> outcome.failure() != null).singleElement()
|
||
.satisfies(outcome -> {
|
||
assertThat(outcome.failure()).isInstanceOf(ScheduleRuleConflictException.class);
|
||
assertThat(outcome.failure())
|
||
.hasMessage("Невозможно сохранить правило: слот занят");
|
||
});
|
||
assertThat(queryCount("SELECT count(*) FROM schedule_rules WHERE semester_id = ?", ids.semesterId()))
|
||
.isEqualTo(1L);
|
||
assertThat(queryCount("""
|
||
SELECT count(*)
|
||
FROM schedule_rule_slots slot
|
||
JOIN schedule_rules rule ON rule.id = slot.schedule_rule_id
|
||
WHERE rule.semester_id = ?
|
||
""", ids.semesterId())).isEqualTo(1L);
|
||
assertThat(queryCount("""
|
||
SELECT count(*)
|
||
FROM schedule_rule_groups rule_group
|
||
JOIN schedule_rules rule ON rule.id = rule_group.schedule_rule_id
|
||
WHERE rule.semester_id = ?
|
||
""", ids.semesterId())).isEqualTo(1L);
|
||
verify(scheduleGeneratorService, times(1)).clearCache();
|
||
}
|
||
|
||
private CreateOutcome createAfterSignal(ScheduleRuleDto request,
|
||
CountDownLatch ready,
|
||
CountDownLatch start) {
|
||
ready.countDown();
|
||
try {
|
||
if (!start.await(5, TimeUnit.SECONDS)) {
|
||
return new CreateOutcome(null, new IllegalStateException(
|
||
"Конкурентные операции не получили общий сигнал запуска"
|
||
));
|
||
}
|
||
return new CreateOutcome(service.create(request), null);
|
||
} catch (Throwable failure) {
|
||
return new CreateOutcome(null, failure);
|
||
}
|
||
}
|
||
|
||
private void awaitReady(CountDownLatch ready) {
|
||
try {
|
||
if (!ready.await(5, TimeUnit.SECONDS)) {
|
||
throw new IllegalStateException("Конкурентные операции не успели подготовиться");
|
||
}
|
||
} catch (InterruptedException exception) {
|
||
Thread.currentThread().interrupt();
|
||
throw new IllegalStateException("Ожидание конкурентных операций прервано", exception);
|
||
}
|
||
}
|
||
|
||
private int lockSemesterRow(Connection connection, long semesterId) throws SQLException {
|
||
try (PreparedStatement statement = connection.prepareStatement(
|
||
"SELECT id, pg_backend_pid() FROM semesters WHERE id = ? FOR UPDATE"
|
||
)) {
|
||
statement.setLong(1, semesterId);
|
||
try (ResultSet result = statement.executeQuery()) {
|
||
assertThat(result.next()).as("Строка тестового семестра должна существовать").isTrue();
|
||
assertThat(result.getLong(1)).isEqualTo(semesterId);
|
||
int holderPid = result.getInt(2);
|
||
assertThat(holderPid).isPositive();
|
||
return holderPid;
|
||
}
|
||
}
|
||
}
|
||
|
||
private void awaitTwoBlockedCreates(List<Future<CreateOutcome>> futures, int holderPid) {
|
||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
|
||
do {
|
||
failOnEarlyCompletion(futures);
|
||
Long blocked = jdbcTemplate.queryForObject("""
|
||
SELECT count(DISTINCT activity.pid)
|
||
FROM pg_stat_activity activity
|
||
WHERE activity.datname = current_database()
|
||
AND activity.pid <> ?
|
||
AND activity.state = 'active'
|
||
AND activity.wait_event_type = 'Lock'
|
||
AND lower(activity.query) LIKE '%semester%'
|
||
AND EXISTS (
|
||
SELECT 1
|
||
FROM pg_locks relation_lock
|
||
WHERE relation_lock.pid = activity.pid
|
||
AND relation_lock.relation = 'semesters'::regclass
|
||
)
|
||
""", Long.class, holderPid);
|
||
if (blocked != null && blocked >= 2L) {
|
||
return;
|
||
}
|
||
try {
|
||
Thread.sleep(25);
|
||
} catch (InterruptedException exception) {
|
||
Thread.currentThread().interrupt();
|
||
throw new IllegalStateException("Ожидание блокировки операций прервано", exception);
|
||
}
|
||
} while (System.nanoTime() < deadline);
|
||
|
||
failOnEarlyCompletion(futures);
|
||
throw new IllegalStateException("Обе операции создания не заблокировались на строке одного семестра. "
|
||
+ "Снимок PostgreSQL: " + lockWaitSnapshot(holderPid));
|
||
}
|
||
|
||
private void failOnEarlyCompletion(List<Future<CreateOutcome>> futures) {
|
||
for (int index = 0; index < futures.size(); index++) {
|
||
Future<CreateOutcome> future = futures.get(index);
|
||
if (!future.isDone()) {
|
||
continue;
|
||
}
|
||
try {
|
||
CreateOutcome outcome = future.get(0, TimeUnit.MILLISECONDS);
|
||
String detail = outcome.failure() == null
|
||
? "операция неожиданно завершилась успешно"
|
||
: outcome.failure().getClass().getSimpleName() + ": " + outcome.failure().getMessage();
|
||
throw new IllegalStateException(
|
||
"Конкурентная операция " + (index + 1) + " завершилась до снятия внешнего барьера: " + detail
|
||
);
|
||
} catch (IllegalStateException exception) {
|
||
throw exception;
|
||
} catch (Exception exception) {
|
||
throw new IllegalStateException(
|
||
"Не удалось прочитать ранний результат конкурентной операции " + (index + 1),
|
||
exception
|
||
);
|
||
}
|
||
}
|
||
}
|
||
|
||
private String lockWaitSnapshot(int holderPid) {
|
||
List<String> rows = jdbcTemplate.query("""
|
||
SELECT concat(
|
||
'pid=', activity.pid,
|
||
', state=', activity.state,
|
||
', wait=', coalesce(activity.wait_event_type, '-'), '/', coalesce(activity.wait_event, '-'),
|
||
', blockers=', pg_blocking_pids(activity.pid)::text,
|
||
', query=', left(regexp_replace(activity.query, '\\s+', ' ', 'g'), 180)
|
||
)
|
||
FROM pg_stat_activity activity
|
||
WHERE activity.datname = current_database()
|
||
AND activity.pid <> ?
|
||
ORDER BY activity.pid
|
||
""", (result, rowNumber) -> result.getString(1), holderPid);
|
||
return rows.isEmpty() ? "активные подключения не найдены" : String.join(" | ", rows);
|
||
}
|
||
|
||
private SeedIds createIsolatedSemesterAndLoadSeedIds() {
|
||
Long academicYearId = jdbcTemplate.queryForObject("""
|
||
INSERT INTO academic_years (title, start_date, end_date)
|
||
VALUES ('2098-2099', DATE '2098-09-01', DATE '2099-06-30')
|
||
RETURNING id
|
||
""", Long.class);
|
||
Long semesterId = jdbcTemplate.queryForObject("""
|
||
INSERT INTO semesters (academic_year_id, semester_type, start_date, end_date)
|
||
VALUES (?, 'autumn', DATE '2098-09-01', DATE '2099-01-31')
|
||
RETURNING id
|
||
""", Long.class, academicYearId);
|
||
Long scheduleVersionId = jdbcTemplate.queryForObject("""
|
||
INSERT INTO schedule_versions (semester_id, version_number, name, status)
|
||
VALUES (?, 1, 'Конкурентный черновик', 'DRAFT')
|
||
RETURNING id
|
||
""", Long.class, semesterId);
|
||
|
||
return new SeedIds(
|
||
semesterId,
|
||
scheduleVersionId,
|
||
queryId("SELECT id FROM subjects WHERE name = 'Высшая математика'"),
|
||
queryId("SELECT id FROM student_groups WHERE name = 'ИВТ-21-1'"),
|
||
queryId("""
|
||
SELECT slot.id
|
||
FROM time_slots slot
|
||
JOIN time_slot_scopes scope ON scope.id = slot.time_slot_scope_id
|
||
WHERE scope.code = 'default' AND slot.order_number = 1
|
||
"""),
|
||
queryId("SELECT id FROM users WHERE username = 'Тестовый преподаватель'"),
|
||
queryId("SELECT id FROM classrooms WHERE name = '101 Ленинская'"),
|
||
queryId("SELECT id FROM lesson_types WHERE name = 'Лекция'")
|
||
);
|
||
}
|
||
|
||
private long queryId(String sql) {
|
||
Long value = jdbcTemplate.queryForObject(sql, Long.class);
|
||
assertThat(value).as("Тестовый seed-идентификатор должен существовать: %s", sql).isNotNull();
|
||
return value;
|
||
}
|
||
|
||
private long queryCount(String sql, Object... arguments) {
|
||
Long value = jdbcTemplate.queryForObject(sql, Long.class, arguments);
|
||
assertThat(value).isNotNull();
|
||
return value;
|
||
}
|
||
|
||
private ScheduleRuleDto request(SeedIds ids) {
|
||
ScheduleRuleSlotDto slot = new ScheduleRuleSlotDto(
|
||
null,
|
||
1,
|
||
null,
|
||
ScheduleParity.BOTH,
|
||
ids.timeSlotId(),
|
||
null,
|
||
null,
|
||
null,
|
||
null,
|
||
List.of(),
|
||
List.of(),
|
||
ids.teacherId(),
|
||
null,
|
||
ids.classroomId(),
|
||
null,
|
||
ids.lessonTypeId(),
|
||
null,
|
||
"Очно"
|
||
);
|
||
return new ScheduleRuleDto(
|
||
null,
|
||
ids.subjectId(),
|
||
null,
|
||
ids.semesterId(),
|
||
null,
|
||
null,
|
||
2,
|
||
0,
|
||
0,
|
||
1,
|
||
1,
|
||
1,
|
||
List.of(ids.groupId()),
|
||
List.of(),
|
||
List.of(slot),
|
||
ids.scheduleVersionId()
|
||
);
|
||
}
|
||
|
||
@SpringBootConfiguration
|
||
@EnableAutoConfiguration
|
||
@EntityScan(basePackageClasses = ScheduleRule.class)
|
||
@EnableJpaRepositories(basePackageClasses = ScheduleRuleRepository.class)
|
||
@Import({ScheduleRuleService.class, TestDependencies.class})
|
||
static class TestApplication {
|
||
}
|
||
|
||
@TestConfiguration(proxyBeanMethods = false)
|
||
static class TestDependencies {
|
||
|
||
@Bean
|
||
ScheduleGeneratorService scheduleGeneratorService() {
|
||
return mock(ScheduleGeneratorService.class);
|
||
}
|
||
}
|
||
|
||
private record SeedIds(
|
||
long semesterId,
|
||
long scheduleVersionId,
|
||
long subjectId,
|
||
long groupId,
|
||
long timeSlotId,
|
||
long teacherId,
|
||
long classroomId,
|
||
long lessonTypeId
|
||
) {
|
||
}
|
||
|
||
private record CreateOutcome(ScheduleRuleDto result, Throwable failure) {
|
||
}
|
||
}
|