fix(personal): require private object storage
This commit is contained in:
+38
-3
@@ -11,6 +11,7 @@ import org.dromara.aihr.personal.support.PersonalOwner;
|
|||||||
import org.dromara.common.core.exception.ServiceException;
|
import org.dromara.common.core.exception.ServiceException;
|
||||||
import org.dromara.common.oss.core.OssClient;
|
import org.dromara.common.oss.core.OssClient;
|
||||||
import org.dromara.common.oss.entity.UploadResult;
|
import org.dromara.common.oss.entity.UploadResult;
|
||||||
|
import org.dromara.common.oss.enums.AccessPolicyType;
|
||||||
import org.dromara.common.oss.factory.OssFactory;
|
import org.dromara.common.oss.factory.OssFactory;
|
||||||
import org.dromara.system.domain.vo.SysOssVo;
|
import org.dromara.system.domain.vo.SysOssVo;
|
||||||
import org.dromara.system.service.ISysOssService;
|
import org.dromara.system.service.ISysOssService;
|
||||||
@@ -60,7 +61,8 @@ public class PersonalIngestionService {
|
|||||||
PersonalKnowledgeProperties properties, ISysOssService ossService,
|
PersonalKnowledgeProperties properties, ISysOssService ossService,
|
||||||
ObjectMapper objectMapper) {
|
ObjectMapper objectMapper) {
|
||||||
this(jdbcTemplate, spaceService, properties, ossService, objectMapper,
|
this(jdbcTemplate, spaceService, properties, ossService, objectMapper,
|
||||||
new DefaultPersonalObjectStore(jdbcTemplate), IdWorker::getId);
|
new DefaultPersonalObjectStore(jdbcTemplate, properties, PersonalIngestionService::ossClient),
|
||||||
|
IdWorker::getId);
|
||||||
}
|
}
|
||||||
|
|
||||||
private PersonalIngestionService(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService,
|
private PersonalIngestionService(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService,
|
||||||
@@ -84,6 +86,12 @@ public class PersonalIngestionService {
|
|||||||
objectStore, itemIdSupplier);
|
objectStore, itemIdSupplier);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static PersonalObjectStore objectStoreForTest(JdbcTemplate jdbcTemplate,
|
||||||
|
PersonalKnowledgeProperties properties,
|
||||||
|
OssClientProvider clientProvider) {
|
||||||
|
return new DefaultPersonalObjectStore(jdbcTemplate, properties, clientProvider);
|
||||||
|
}
|
||||||
|
|
||||||
@Transactional
|
@Transactional
|
||||||
public ItemCreatedResponse createText(PersonalOwner owner, TextItemRequest request) {
|
public ItemCreatedResponse createText(PersonalOwner owner, TextItemRequest request) {
|
||||||
validateOwner(owner);
|
validateOwner(owner);
|
||||||
@@ -355,17 +363,28 @@ public class PersonalIngestionService {
|
|||||||
void deletePhysical(SysOssVo uploaded);
|
void deletePhysical(SysOssVo uploaded);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@FunctionalInterface
|
||||||
|
public interface OssClientProvider {
|
||||||
|
OssClient get(String configKey);
|
||||||
|
}
|
||||||
|
|
||||||
private static final class DefaultPersonalObjectStore implements PersonalObjectStore {
|
private static final class DefaultPersonalObjectStore implements PersonalObjectStore {
|
||||||
private final JdbcTemplate jdbcTemplate;
|
private final JdbcTemplate jdbcTemplate;
|
||||||
|
private final PersonalKnowledgeProperties properties;
|
||||||
|
private final OssClientProvider clientProvider;
|
||||||
|
|
||||||
private DefaultPersonalObjectStore(JdbcTemplate jdbcTemplate) {
|
private DefaultPersonalObjectStore(JdbcTemplate jdbcTemplate, PersonalKnowledgeProperties properties,
|
||||||
|
OssClientProvider clientProvider) {
|
||||||
this.jdbcTemplate = jdbcTemplate;
|
this.jdbcTemplate = jdbcTemplate;
|
||||||
|
this.properties = properties;
|
||||||
|
this.clientProvider = clientProvider;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public SysOssVo upload(PersonalOwner owner, long itemId, String objectKey, String suffix, String mimeType,
|
public SysOssVo upload(PersonalOwner owner, long itemId, String objectKey, String suffix, String mimeType,
|
||||||
byte[] bytes) {
|
byte[] bytes) {
|
||||||
OssClient storage = OssFactory.instance();
|
OssClient storage = clientProvider.get(normalizedConfigKey(properties.getOssConfigKey()));
|
||||||
|
requirePrivate(storage);
|
||||||
UploadResult result = storage.upload(
|
UploadResult result = storage.upload(
|
||||||
new ByteArrayInputStream(bytes), objectKey, (long) bytes.length, mimeType);
|
new ByteArrayInputStream(bytes), objectKey, (long) bytes.length, mimeType);
|
||||||
long ossId = IdWorker.getId();
|
long ossId = IdWorker.getId();
|
||||||
@@ -406,4 +425,20 @@ public class PersonalIngestionService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static OssClient ossClient(String configKey) {
|
||||||
|
return configKey == null || configKey.isBlank()
|
||||||
|
? OssFactory.instance()
|
||||||
|
: OssFactory.instance(configKey);
|
||||||
|
}
|
||||||
|
|
||||||
|
private static String normalizedConfigKey(String value) {
|
||||||
|
return value == null ? "" : value.trim();
|
||||||
|
}
|
||||||
|
|
||||||
|
private static void requirePrivate(OssClient storage) {
|
||||||
|
if (storage == null || storage.getAccessPolicy() != AccessPolicyType.PRIVATE) {
|
||||||
|
throw new ServiceException("PERSONAL_OSS_NOT_PRIVATE");
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+28
-3
@@ -4,7 +4,9 @@ import lombok.extern.slf4j.Slf4j;
|
|||||||
import org.dromara.aihr.knowledge.parse.KnowledgeDocumentParser;
|
import org.dromara.aihr.knowledge.parse.KnowledgeDocumentParser;
|
||||||
import org.dromara.aihr.knowledge.parse.ParsedDocument;
|
import org.dromara.aihr.knowledge.parse.ParsedDocument;
|
||||||
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
|
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
|
||||||
|
import org.dromara.common.core.exception.ServiceException;
|
||||||
import org.dromara.common.oss.core.OssClient;
|
import org.dromara.common.oss.core.OssClient;
|
||||||
|
import org.dromara.common.oss.enums.AccessPolicyType;
|
||||||
import org.dromara.common.oss.factory.OssFactory;
|
import org.dromara.common.oss.factory.OssFactory;
|
||||||
import org.dromara.system.domain.vo.SysOssVo;
|
import org.dromara.system.domain.vo.SysOssVo;
|
||||||
import org.dromara.system.service.ISysOssService;
|
import org.dromara.system.service.ISysOssService;
|
||||||
@@ -39,7 +41,7 @@ public class PersonalIngestionWorker {
|
|||||||
KnowledgeDocumentParser parser, PersonalKnowledgeProperties properties,
|
KnowledgeDocumentParser parser, PersonalKnowledgeProperties properties,
|
||||||
PlatformTransactionManager transactionManager) {
|
PlatformTransactionManager transactionManager) {
|
||||||
this(jdbcTemplate, parser, new TransactionTemplate(transactionManager),
|
this(jdbcTemplate, parser, new TransactionTemplate(transactionManager),
|
||||||
defaultReader(ossService), configuredMaxBytes(properties));
|
defaultReader(ossService, PersonalIngestionWorker::ossClient), configuredMaxBytes(properties));
|
||||||
}
|
}
|
||||||
|
|
||||||
private PersonalIngestionWorker(JdbcTemplate jdbcTemplate, KnowledgeDocumentParser parser,
|
private PersonalIngestionWorker(JdbcTemplate jdbcTemplate, KnowledgeDocumentParser parser,
|
||||||
@@ -60,6 +62,11 @@ public class PersonalIngestionWorker {
|
|||||||
20L * 1024 * 1024);
|
20L * 1024 * 1024);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static StoredObjectReader objectReaderForTest(ISysOssService ossService,
|
||||||
|
OssClientProvider clientProvider) {
|
||||||
|
return defaultReader(ossService, clientProvider);
|
||||||
|
}
|
||||||
|
|
||||||
@Scheduled(fixedDelayString = "${aihr.personal.ingestion-delay-ms:2000}")
|
@Scheduled(fixedDelayString = "${aihr.personal.ingestion-delay-ms:2000}")
|
||||||
public void poll() {
|
public void poll() {
|
||||||
processNext();
|
processNext();
|
||||||
@@ -156,7 +163,7 @@ public class PersonalIngestionWorker {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private static StoredObjectReader defaultReader(ISysOssService ossService) {
|
private static StoredObjectReader defaultReader(ISysOssService ossService, OssClientProvider clientProvider) {
|
||||||
return (ossId, expectedPrefix, ownerUserId, maxBytes) -> {
|
return (ossId, expectedPrefix, ownerUserId, maxBytes) -> {
|
||||||
SysOssVo object = ossService.getById(ossId);
|
SysOssVo object = ossService.getById(ossId);
|
||||||
if (object == null || object.getFileName() == null || object.getFileName().isBlank()
|
if (object == null || object.getFileName() == null || object.getFileName().isBlank()
|
||||||
@@ -165,7 +172,10 @@ public class PersonalIngestionWorker {
|
|||||||
|| !object.getFileName().startsWith(expectedPrefix)) {
|
|| !object.getFileName().startsWith(expectedPrefix)) {
|
||||||
throw new IOException("personal source object is unavailable");
|
throw new IOException("personal source object is unavailable");
|
||||||
}
|
}
|
||||||
OssClient storage = OssFactory.instance(object.getService());
|
OssClient storage = clientProvider.get(object.getService());
|
||||||
|
if (storage == null || storage.getAccessPolicy() != AccessPolicyType.PRIVATE) {
|
||||||
|
throw new ServiceException("PERSONAL_OSS_NOT_PRIVATE");
|
||||||
|
}
|
||||||
try (InputStream input = storage.getObjectContent(object.getFileName())) {
|
try (InputStream input = storage.getObjectContent(object.getFileName())) {
|
||||||
int boundedLimit = (int) Math.min(Integer.MAX_VALUE - 1L, maxBytes);
|
int boundedLimit = (int) Math.min(Integer.MAX_VALUE - 1L, maxBytes);
|
||||||
byte[] bytes = input.readNBytes(boundedLimit + 1);
|
byte[] bytes = input.readNBytes(boundedLimit + 1);
|
||||||
@@ -210,6 +220,10 @@ public class PersonalIngestionWorker {
|
|||||||
case INVALID -> new Failure("PERSONAL_PARSE_INVALID", "资料解析失败,请检查文件后重试");
|
case INVALID -> new Failure("PERSONAL_PARSE_INVALID", "资料解析失败,请检查文件后重试");
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
if (candidate instanceof ServiceException serviceError
|
||||||
|
&& "PERSONAL_OSS_NOT_PRIVATE".equals(serviceError.getMessage())) {
|
||||||
|
return new Failure("PERSONAL_OSS_NOT_PRIVATE", "个人资料存储策略不可用");
|
||||||
|
}
|
||||||
candidate = candidate.getCause();
|
candidate = candidate.getCause();
|
||||||
}
|
}
|
||||||
return new Failure("PERSONAL_PARSE_FAILED", "资料处理失败,请稍后重试");
|
return new Failure("PERSONAL_PARSE_FAILED", "资料处理失败,请稍后重试");
|
||||||
@@ -258,6 +272,11 @@ public class PersonalIngestionWorker {
|
|||||||
StoredObject read(long ossId, String expectedPrefix, long ownerUserId, long maxBytes) throws Exception;
|
StoredObject read(long ossId, String expectedPrefix, long ownerUserId, long maxBytes) throws Exception;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@FunctionalInterface
|
||||||
|
public interface OssClientProvider {
|
||||||
|
OssClient get(String configKey);
|
||||||
|
}
|
||||||
|
|
||||||
public record StoredObject(String fileName, byte[] bytes) {
|
public record StoredObject(String fileName, byte[] bytes) {
|
||||||
public StoredObject {
|
public StoredObject {
|
||||||
bytes = bytes == null ? new byte[0] : bytes.clone();
|
bytes = bytes == null ? new byte[0] : bytes.clone();
|
||||||
@@ -275,4 +294,10 @@ public class PersonalIngestionWorker {
|
|||||||
|
|
||||||
private record Failure(String code, String message) {
|
private record Failure(String code, String message) {
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static OssClient ossClient(String configKey) {
|
||||||
|
return configKey == null || configKey.isBlank()
|
||||||
|
? OssFactory.instance()
|
||||||
|
: OssFactory.instance(configKey);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+2
@@ -15,4 +15,6 @@ public class PersonalKnowledgeProperties {
|
|||||||
private int maxItems = 1000;
|
private int maxItems = 1000;
|
||||||
private int downloadUrlMinutes = 5;
|
private int downloadUrlMinutes = 5;
|
||||||
private String qdrantCollection = "aihr_personal_knowledge";
|
private String qdrantCollection = "aihr_personal_knowledge";
|
||||||
|
/** Optional sys_oss_config key. Blank selects the system default client. */
|
||||||
|
private String ossConfigKey = "";
|
||||||
}
|
}
|
||||||
|
|||||||
+56
@@ -10,6 +10,9 @@ import org.dromara.aihr.personal.service.PersonalSpaceService;
|
|||||||
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
|
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
|
||||||
import org.dromara.aihr.personal.support.PersonalOwner;
|
import org.dromara.aihr.personal.support.PersonalOwner;
|
||||||
import org.dromara.common.core.exception.ServiceException;
|
import org.dromara.common.core.exception.ServiceException;
|
||||||
|
import org.dromara.common.oss.core.OssClient;
|
||||||
|
import org.dromara.common.oss.entity.UploadResult;
|
||||||
|
import org.dromara.common.oss.enums.AccessPolicyType;
|
||||||
import org.dromara.system.domain.vo.SysOssVo;
|
import org.dromara.system.domain.vo.SysOssVo;
|
||||||
import org.dromara.system.service.ISysOssService;
|
import org.dromara.system.service.ISysOssService;
|
||||||
import org.junit.jupiter.api.Tag;
|
import org.junit.jupiter.api.Tag;
|
||||||
@@ -249,6 +252,59 @@ class PersonalIngestionServiceTest {
|
|||||||
verifyNoInteractions(fixture.spaces, fixture.store, fixture.jdbc, fixture.oss);
|
verifyNoInteractions(fixture.spaces, fixture.store, fixture.jdbc, fixture.oss);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void defaultPublicOssClientIsRejectedBeforeUploadOrMetadataWrite() {
|
||||||
|
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||||
|
PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties();
|
||||||
|
PersonalIngestionService.OssClientProvider clients =
|
||||||
|
mock(PersonalIngestionService.OssClientProvider.class);
|
||||||
|
OssClient publicClient = mock(OssClient.class);
|
||||||
|
when(clients.get("")).thenReturn(publicClient);
|
||||||
|
when(publicClient.getAccessPolicy()).thenReturn(AccessPolicyType.PUBLIC);
|
||||||
|
PersonalObjectStore store = PersonalIngestionService.objectStoreForTest(jdbc, properties, clients);
|
||||||
|
|
||||||
|
ServiceException error = assertThrows(ServiceException.class, () -> store.upload(
|
||||||
|
OWNER, 100L, "personal/000000/101/100/a.txt", "txt", "text/plain", new byte[]{1}));
|
||||||
|
|
||||||
|
assertEquals("PERSONAL_OSS_NOT_PRIVATE", error.getMessage());
|
||||||
|
verify(clients).get("");
|
||||||
|
verify(publicClient, never()).upload(any(java.io.InputStream.class), anyString(), anyLong(), anyString());
|
||||||
|
verifyNoInteractions(jdbc);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void configuredPrivateOssClientIsSelectedAndWritesSystemMetadata() {
|
||||||
|
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||||
|
PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties();
|
||||||
|
properties.setOssConfigKey(" personal-private ");
|
||||||
|
PersonalIngestionService.OssClientProvider clients =
|
||||||
|
mock(PersonalIngestionService.OssClientProvider.class);
|
||||||
|
OssClient privateClient = mock(OssClient.class);
|
||||||
|
when(clients.get("personal-private")).thenReturn(privateClient);
|
||||||
|
when(privateClient.getAccessPolicy()).thenReturn(AccessPolicyType.PRIVATE);
|
||||||
|
when(privateClient.getConfigKey()).thenReturn("personal-private");
|
||||||
|
when(privateClient.upload(any(java.io.InputStream.class), anyString(), anyLong(), eq("text/plain")))
|
||||||
|
.thenReturn(UploadResult.builder()
|
||||||
|
.filename("personal/000000/101/100/a.txt")
|
||||||
|
.url("https://private.invalid/personal/000000/101/100/a.txt")
|
||||||
|
.build());
|
||||||
|
when(jdbc.update(contains("insert into sys_oss"), any(), any(), any(), any(), any(), any(), any(), any(), any()))
|
||||||
|
.thenReturn(1);
|
||||||
|
PersonalObjectStore store = PersonalIngestionService.objectStoreForTest(jdbc, properties, clients);
|
||||||
|
|
||||||
|
SysOssVo uploaded = store.upload(
|
||||||
|
OWNER, 100L, "personal/000000/101/100/a.txt", "txt", "text/plain", new byte[]{1});
|
||||||
|
|
||||||
|
assertEquals("personal-private", uploaded.getService());
|
||||||
|
verify(clients).get("personal-private");
|
||||||
|
verify(privateClient).upload(any(java.io.InputStream.class),
|
||||||
|
eq("personal/000000/101/100/a.txt"), eq(1L), eq("text/plain"));
|
||||||
|
verify(jdbc).update(contains("insert into sys_oss"), any(), eq("000000"),
|
||||||
|
eq("personal/000000/101/100/a.txt"), eq("a.txt"), eq(".txt"),
|
||||||
|
eq("https://private.invalid/personal/000000/101/100/a.txt"), eq(101L), eq(101L),
|
||||||
|
eq("personal-private"));
|
||||||
|
}
|
||||||
|
|
||||||
private static Fixture fixture(long itemId) {
|
private static Fixture fixture(long itemId) {
|
||||||
return fixture(() -> itemId);
|
return fixture(() -> itemId);
|
||||||
}
|
}
|
||||||
|
|||||||
+31
@@ -4,6 +4,9 @@ import org.dromara.aihr.knowledge.parse.KnowledgeDocumentParser;
|
|||||||
import org.dromara.aihr.knowledge.parse.ParsedDocument;
|
import org.dromara.aihr.knowledge.parse.ParsedDocument;
|
||||||
import org.dromara.aihr.personal.service.PersonalIngestionWorker;
|
import org.dromara.aihr.personal.service.PersonalIngestionWorker;
|
||||||
import org.dromara.system.service.ISysOssService;
|
import org.dromara.system.service.ISysOssService;
|
||||||
|
import org.dromara.system.domain.vo.SysOssVo;
|
||||||
|
import org.dromara.common.oss.core.OssClient;
|
||||||
|
import org.dromara.common.oss.enums.AccessPolicyType;
|
||||||
import org.junit.jupiter.api.Tag;
|
import org.junit.jupiter.api.Tag;
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
|
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
|
||||||
@@ -94,6 +97,34 @@ class PersonalIngestionWorkerTest {
|
|||||||
eq("资料解析失败,请检查文件后重试"), eq("000000"), eq(101L), eq(9L));
|
eq("资料解析失败,请检查文件后重试"), eq("000000"), eq(101L), eq(9L));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void workerRejectsPublicPolicyBeforeReadingObjectContent() throws Exception {
|
||||||
|
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||||
|
ISysOssService ossService = mock(ISysOssService.class);
|
||||||
|
OssClient publicClient = mock(OssClient.class);
|
||||||
|
SysOssVo object = new SysOssVo();
|
||||||
|
object.setOssId(81L);
|
||||||
|
object.setFileName("personal/000000/101/9/a.txt");
|
||||||
|
object.setService("public-client");
|
||||||
|
object.setCreateBy(101L);
|
||||||
|
when(ossService.getById(81L)).thenReturn(object);
|
||||||
|
when(publicClient.getAccessPolicy()).thenReturn(AccessPolicyType.PUBLIC);
|
||||||
|
PersonalIngestionWorker.OssClientProvider clients =
|
||||||
|
mock(PersonalIngestionWorker.OssClientProvider.class);
|
||||||
|
when(clients.get("public-client")).thenReturn(publicClient);
|
||||||
|
when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item()));
|
||||||
|
when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1);
|
||||||
|
PersonalIngestionWorker worker = PersonalIngestionWorker.forTest(
|
||||||
|
jdbc, ossService, mock(KnowledgeDocumentParser.class), immediateTransactions(),
|
||||||
|
PersonalIngestionWorker.objectReaderForTest(ossService, clients));
|
||||||
|
|
||||||
|
assertTrue(worker.processNext());
|
||||||
|
|
||||||
|
verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_OSS_NOT_PRIVATE"),
|
||||||
|
eq("个人资料存储策略不可用"), eq("000000"), eq(101L), eq(9L));
|
||||||
|
verify(publicClient, never()).getObjectContent(any(String.class));
|
||||||
|
}
|
||||||
|
|
||||||
private static Map<String, Object> item() {
|
private static Map<String, Object> item() {
|
||||||
return Map.of(
|
return Map.of(
|
||||||
"id", 9L,
|
"id", 9L,
|
||||||
|
|||||||
+1
@@ -46,6 +46,7 @@ class PersonalSpaceServiceTest {
|
|||||||
assertEquals(1000, properties.getMaxItems());
|
assertEquals(1000, properties.getMaxItems());
|
||||||
assertEquals(5, properties.getDownloadUrlMinutes());
|
assertEquals(5, properties.getDownloadUrlMinutes());
|
||||||
assertEquals("aihr_personal_knowledge", properties.getQdrantCollection());
|
assertEquals("aihr_personal_knowledge", properties.getQdrantCollection());
|
||||||
|
assertEquals("", properties.getOssConfigKey());
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
Reference in New Issue
Block a user