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
33 changes: 26 additions & 7 deletions NEWS.md
Original file line number Diff line number Diff line change
@@ -1,10 +1,29 @@
## 3.0.0 (IN PROGRESS)
- Update order of fields for Title in profiles. [MODLD-1090](https://folio-org.atlassian.net/browse/MODLD-1090)
- Standalone Authority resource CRUD API [MODLD-1082](https://folio-org.atlassian.net/browse/MODLD-1082)
- Authority index message sending [MODLD-1099](https://folio-org.atlassian.net/browse/MODLD-1099)
- Update profile settings to allow for multiple settings per profile [MODLD-1040](https://folio-org.atlassian.net/browse/MODLD-1040)
- Update some Authority profile MARC tooltips [MODLD-1083](https://folio-org.atlassian.net/browse/MODLD-1083)
- Authority API Validation: Name is required field [MODLD-1107](https://folio-org.atlassian.net/browse/MODLD-1107)
## v3.0.0 (IN PROGRESS)
### Breaking changes
* Description ([ISSUE](https://folio-org.atlassian.net/browse/ISSUE))

### New APIs versions
* Provides `API_NAME vX.Y`
* Requires `API_NAME vX.Y`

### Features
* Update order of fields for Title in profiles. [MODLD-1090](https://folio-org.atlassian.net/browse/MODLD-1090)
* Standalone Authority resource CRUD API [MODLD-1082](https://folio-org.atlassian.net/browse/MODLD-1082)
* Authority index message sending [MODLD-1099](https://folio-org.atlassian.net/browse/MODLD-1099)
* Update profile settings to allow for multiple settings per profile [MODLD-1040](https://folio-org.atlassian.net/browse/MODLD-1040)
* Update some Authority profile MARC tooltips [MODLD-1083](https://folio-org.atlassian.net/browse/MODLD-1083)
* Authority API Validation: Name is required field [MODLD-1107](https://folio-org.atlassian.net/browse/MODLD-1107)

### Bug fixes
* Improve case-insensitive handling of Kafka headers ([MODLD-1117](https://folio-org.atlassian.net/browse/MODLD-1117))

### Tech Dept
* Description ([ISSUE](https://folio-org.atlassian.net/browse/ISSUE))

### Dependencies
* Bump `LIB_NAME` from `OLD_VERSION` to `NEW_VERSION`
* Add `LIB_NAME VERSION`
* Remove `LIB_NAME`

## 2.0.3 (02-06-2026)
- Exclude LIGHT_RESOURCE from reindexing [MODLD-1071](https://folio-org.atlassian.net/browse/MODLD-1071)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

import java.util.List;
import java.util.Set;
import java.util.stream.StreamSupport;
import lombok.RequiredArgsConstructor;
import lombok.extern.log4j.Log4j2;
import org.apache.kafka.clients.consumer.ConsumerRecord;
Expand Down Expand Up @@ -75,7 +76,7 @@ private void processRecord(ConsumerRecord<String, SourceRecordDomainEvent> consu

private boolean notAllRequiredHeaders(Headers headers) {
return !REQUIRED_HEADERS.stream()
.map(required -> headers.headers(required).iterator())
.allMatch(iterator -> iterator.hasNext() && iterator.next().value().length > 0);
.allMatch(required -> StreamSupport.stream(headers.spliterator(), false)
.anyMatch(h -> h.key().equalsIgnoreCase(required) && h.value().length > 0));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ public <T> void executeWithRetry(Headers headers, Retryable<T> retryable, Consum

private FolioExecutionContext folioContextFromKafkaHeadersNoToken(Headers headers) {
Map<String, Object> okapiHeaders = Arrays.stream(headers.toArray())
.filter(header -> !TOKEN.equals(header.key()))
.filter(header -> !TOKEN.equalsIgnoreCase(header.key()))
.collect(toMap(
Header::key,
Header::value,
Expand Down
9 changes: 6 additions & 3 deletions src/main/java/org/folio/linked/data/util/KafkaUtils.java
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
package org.folio.linked.data.util;

import static java.util.Optional.ofNullable;
import static org.folio.spring.integration.XOkapiHeaders.TENANT;

import java.util.Optional;
Expand All @@ -14,8 +13,12 @@
public class KafkaUtils {

public static Optional<String> getHeaderValueByName(ConsumerRecord<String, ?> consumerRecord, String headerName) {
return ofNullable(consumerRecord.headers().lastHeader(headerName))
.map(header -> new String(header.value()));
for (var header : consumerRecord.headers()) {
if (header.key().equalsIgnoreCase(headerName)) {
return Optional.of(new String(header.value()));
}
}
return Optional.empty();
}

public static <T> void handleForExistedTenant(ConsumerRecord<String, T> consumerRecord,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
package org.folio.linked.data.integration.kafka.listener;

import static org.folio.spring.integration.XOkapiHeaders.TENANT;
import static org.folio.spring.integration.XOkapiHeaders.URL;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;

import java.util.List;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.header.internals.RecordHeader;
import org.apache.kafka.common.header.internals.RecordHeaders;
import org.folio.linked.data.domain.dto.SourceRecordDomainEvent;
import org.folio.linked.data.integration.kafka.listener.handler.srs.SourceRecordDomainEventHandler;
import org.folio.linked.data.service.tenant.LinkedDataTenantService;
import org.folio.linked.data.service.tenant.TenantScopedExecutionService;
import org.folio.spring.testing.type.UnitTest;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;

@UnitTest
@ExtendWith(MockitoExtension.class)
class SourceRecordDomainEventListenerTest {

private static final String RECORD_TYPE = "folio.srs.recordType";

@InjectMocks
private SourceRecordDomainEventListener listener;

@Mock
private TenantScopedExecutionService tenantScopedExecutionService;
@Mock
private SourceRecordDomainEventHandler sourceRecordDomainEventHandler;
@Mock
private LinkedDataTenantService linkedDataTenantService;

@Test
void handleSourceRecordDomainEvent_shouldProcessRecord_whenAllRequiredHeadersPresentWithMixedCase() {
// given
var tenant = "test-tenant";
var event = new SourceRecordDomainEvent().id("1");
var consumerRecord = new ConsumerRecord<String, SourceRecordDomainEvent>("topic", 1, 1, "key", event);
var headers = new RecordHeaders(List.of(
new RecordHeader(TENANT.toUpperCase(), tenant.getBytes()),
new RecordHeader(URL.toUpperCase(), "http://okapi:9130".getBytes()),
new RecordHeader(RECORD_TYPE.toUpperCase(), "MARC_BIB".getBytes())
));
ReflectionTestUtils.setField(consumerRecord, "headers", headers);
when(linkedDataTenantService.isTenantExists(tenant)).thenReturn(true);

// when
listener.handleSourceRecordDomainEvent(List.of(consumerRecord));

// then
verify(tenantScopedExecutionService).executeWithRetry(any(), any(), any());
}

@Test
void handleSourceRecordDomainEvent_shouldIgnoreRecord_whenRequiredHeaderHasEmptyValue() {
// given
var tenant = "test-tenant";
var event = new SourceRecordDomainEvent().id("2");
var consumerRecord = new ConsumerRecord<String, SourceRecordDomainEvent>("topic", 1, 1, "key", event);
var headers = new RecordHeaders(List.of(
new RecordHeader(TENANT, tenant.getBytes()),
new RecordHeader(URL, "http://okapi:9130".getBytes()),
new RecordHeader(RECORD_TYPE, new byte[0])
));
ReflectionTestUtils.setField(consumerRecord, "headers", headers);
when(linkedDataTenantService.isTenantExists(tenant)).thenReturn(true);

// when
listener.handleSourceRecordDomainEvent(List.of(consumerRecord));

// then
verifyNoInteractions(tenantScopedExecutionService);
}

@Test
void handleSourceRecordDomainEvent_shouldIgnoreRecord_whenRequiredHeaderIsMissing() {
// given
var tenant = "test-tenant";
var event = new SourceRecordDomainEvent().id("3");
var consumerRecord = new ConsumerRecord<String, SourceRecordDomainEvent>("topic", 1, 1, "key", event);
var headers = new RecordHeaders(List.of(
new RecordHeader(TENANT, tenant.getBytes()),
new RecordHeader(URL, "http://okapi:9130".getBytes())
));
ReflectionTestUtils.setField(consumerRecord, "headers", headers);
when(linkedDataTenantService.isTenantExists(tenant)).thenReturn(true);

// when
listener.handleSourceRecordDomainEvent(List.of(consumerRecord));

// then
verifyNoInteractions(tenantScopedExecutionService);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
package org.folio.linked.data.service.tenant;

import static org.assertj.core.api.Assertions.assertThat;
import static org.folio.spring.integration.XOkapiHeaders.TENANT;
import static org.folio.spring.integration.XOkapiHeaders.TOKEN;
import static org.mockito.Mockito.when;

import java.util.List;
import org.apache.kafka.common.header.internals.RecordHeader;
import org.apache.kafka.common.header.internals.RecordHeaders;
import org.folio.spring.FolioExecutionContext;
import org.folio.spring.testing.type.UnitTest;
import org.folio.spring.tools.context.ExecutionContextBuilder;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.core.retry.RetryTemplate;
import org.springframework.messaging.MessageHeaders;

@UnitTest
@ExtendWith(MockitoExtension.class)
class TenantScopedExecutionServiceTest {

@InjectMocks
private TenantScopedExecutionService tenantScopedExecutionService;

@Mock
private RetryTemplate retryTemplate;
@Mock
private ExecutionContextBuilder contextBuilder;
@Mock
private FolioExecutionContext folioExecutionContext;

@Test
void executeWithRetry_shouldExcludeTokenHeader_whenKeyDiffersByCase() {
// given
var headers = new RecordHeaders(List.of(
new RecordHeader(TENANT, "test-tenant".getBytes()),
new RecordHeader(TOKEN.toUpperCase(), "test-token".getBytes())
));
var captor = ArgumentCaptor.forClass(MessageHeaders.class);
when(contextBuilder.forMessageHeaders(captor.capture())).thenReturn(folioExecutionContext);

// when
tenantScopedExecutionService.executeWithRetry(headers, () -> null, ex -> {});

// then
var capturedHeaders = captor.getValue();
assertThat(capturedHeaders.containsKey(TOKEN.toUpperCase())).isFalse();
assertThat(capturedHeaders).containsKey(TENANT);
}
}
18 changes: 18 additions & 0 deletions src/test/java/org/folio/linked/data/util/KafkaUtilsTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -44,4 +44,22 @@ void getHeaderValueByName_shouldReturnOptionalWithHeader_ifMessageContainsExpect
.isPresent()
.contains(headerValue);
}

@Test
void getHeaderValueByName_shouldReturnOptionalWithHeader_ifHeaderKeyDiffersByCase() {
// given
var headerKey = "headerKey";
var headerValue = UUID.randomUUID().toString();
var consumerRecord = new ConsumerRecord<>("topic", 1, 1, "key", "value");
var headers = new RecordHeaders(List.of(new RecordHeader(headerKey.toUpperCase(), headerValue.getBytes())));
ReflectionTestUtils.setField(consumerRecord, "headers", headers);

// when
var result = KafkaUtils.getHeaderValueByName(consumerRecord, headerKey);

// then
assertThat(result)
.isPresent()
.contains(headerValue);
}
}
Loading