feat(personal): expose authenticated API and deletion workflow
This commit is contained in:
+148
@@ -0,0 +1,148 @@
|
||||
package org.dromara.aihr.personal.controller;
|
||||
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.AskRequest;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.AskResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.DownloadUrlResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.ItemCreatedResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.ItemResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.PageResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.PersonalSearchRequest;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.PersonalSearchResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.SessionDetailResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.SessionResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.SpaceResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.TextItemRequest;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.UrlItemRequest;
|
||||
import org.dromara.aihr.personal.service.PersonalAnswerService;
|
||||
import org.dromara.aihr.personal.service.PersonalCleanupService;
|
||||
import org.dromara.aihr.personal.service.PersonalIngestionService;
|
||||
import org.dromara.aihr.personal.service.PersonalRetrievalService;
|
||||
import org.dromara.aihr.personal.service.PersonalSpaceService;
|
||||
import org.dromara.aihr.personal.service.PersonalUrlFetchService;
|
||||
import org.dromara.aihr.personal.support.PersonalOwner;
|
||||
import org.dromara.aihr.personal.support.PersonalOwnerProvider;
|
||||
import org.dromara.common.core.domain.R;
|
||||
import org.springframework.format.annotation.DateTimeFormat;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.web.bind.annotation.DeleteMapping;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.PathVariable;
|
||||
import org.springframework.web.bind.annotation.PostMapping;
|
||||
import org.springframework.web.bind.annotation.RequestBody;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RequestParam;
|
||||
import org.springframework.web.bind.annotation.RequestPart;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
|
||||
import java.time.LocalDate;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
@RequiredArgsConstructor
|
||||
@RestController
|
||||
@RequestMapping("/api/aihr/personal-assistant")
|
||||
public class PersonalAssistantController {
|
||||
|
||||
private final PersonalOwnerProvider ownerProvider;
|
||||
private final PersonalSpaceService spaceService;
|
||||
private final PersonalIngestionService ingestionService;
|
||||
private final PersonalUrlFetchService urlFetchService;
|
||||
private final PersonalRetrievalService retrievalService;
|
||||
private final PersonalAnswerService answerService;
|
||||
private final PersonalCleanupService cleanupService;
|
||||
|
||||
@GetMapping("/space")
|
||||
public R<SpaceResponse> space() {
|
||||
return R.ok(spaceService.space(owner()));
|
||||
}
|
||||
|
||||
@GetMapping("/items")
|
||||
public R<PageResponse<ItemResponse>> items(@RequestParam(required = false) Integer pageNum,
|
||||
@RequestParam(required = false) Integer pageSize,
|
||||
@RequestParam(required = false) String status,
|
||||
@RequestParam(required = false) String sourceType,
|
||||
@RequestParam(required = false)
|
||||
@DateTimeFormat(iso = DateTimeFormat.ISO.DATE) LocalDate dateFrom,
|
||||
@RequestParam(required = false)
|
||||
@DateTimeFormat(iso = DateTimeFormat.ISO.DATE) LocalDate dateTo,
|
||||
@RequestParam(required = false) String keyword) {
|
||||
return R.ok(spaceService.items(owner(), pageNum, pageSize, status, sourceType, dateFrom, dateTo, keyword));
|
||||
}
|
||||
|
||||
@PostMapping("/items/text")
|
||||
public R<ItemCreatedResponse> createText(@RequestBody TextItemRequest request) {
|
||||
return R.ok(ingestionService.createText(owner(), request));
|
||||
}
|
||||
|
||||
@PostMapping(value = "/items/file", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
|
||||
public R<ItemCreatedResponse> createFile(@RequestPart("file") MultipartFile file,
|
||||
@RequestParam(required = false) String title,
|
||||
@RequestParam(required = false)
|
||||
@DateTimeFormat(iso = DateTimeFormat.ISO.DATE_TIME)
|
||||
LocalDateTime capturedAt) {
|
||||
return R.ok(ingestionService.createFile(owner(), file, title, capturedAt));
|
||||
}
|
||||
|
||||
@PostMapping("/items/url")
|
||||
public R<ItemCreatedResponse> createUrl(@RequestBody UrlItemRequest request) {
|
||||
PersonalOwner owner = owner();
|
||||
PersonalUrlFetchService.FetchResult fetched = urlFetchService.fetch(request == null ? null : request.url());
|
||||
return R.ok(ingestionService.createUrl(owner, request, fetched));
|
||||
}
|
||||
|
||||
@GetMapping("/items/{id}")
|
||||
public R<ItemResponse> item(@PathVariable long id) {
|
||||
return R.ok(spaceService.itemResponse(owner(), id));
|
||||
}
|
||||
|
||||
@PostMapping("/items/{id}/retry")
|
||||
public R<ItemResponse> retry(@PathVariable long id) {
|
||||
PersonalOwner owner = owner();
|
||||
ingestionService.retry(owner, id);
|
||||
return R.ok(spaceService.itemResponse(owner, id));
|
||||
}
|
||||
|
||||
@DeleteMapping("/items/{id}")
|
||||
public R<Map<String, Long>> deleteItem(@PathVariable long id) {
|
||||
return R.ok(Map.of("cleanupJobId", cleanupService.requestDelete(owner(), id)));
|
||||
}
|
||||
|
||||
@GetMapping("/items/{id}/download-url")
|
||||
public R<DownloadUrlResponse> downloadUrl(@PathVariable long id) {
|
||||
return R.ok(spaceService.downloadUrl(owner(), id));
|
||||
}
|
||||
|
||||
@PostMapping("/search")
|
||||
public R<PersonalSearchResponse> search(@RequestBody PersonalSearchRequest request) {
|
||||
return R.ok(new PersonalSearchResponse(request == null ? null : request.queryText(),
|
||||
retrievalService.search(owner(), request)));
|
||||
}
|
||||
|
||||
@PostMapping("/ask")
|
||||
public R<AskResponse> ask(@RequestBody AskRequest request) {
|
||||
return R.ok(answerService.ask(owner(), request));
|
||||
}
|
||||
|
||||
@GetMapping("/sessions")
|
||||
public R<List<SessionResponse>> sessions() {
|
||||
return R.ok(spaceService.sessions(owner()));
|
||||
}
|
||||
|
||||
@GetMapping("/sessions/{id}")
|
||||
public R<SessionDetailResponse> session(@PathVariable long id) {
|
||||
return R.ok(spaceService.session(owner(), id));
|
||||
}
|
||||
|
||||
@DeleteMapping("/sessions/{id}")
|
||||
public R<Void> deleteSession(@PathVariable long id) {
|
||||
spaceService.deleteSession(owner(), id);
|
||||
return R.ok();
|
||||
}
|
||||
|
||||
private PersonalOwner owner() {
|
||||
return ownerProvider.current();
|
||||
}
|
||||
}
|
||||
+227
@@ -0,0 +1,227 @@
|
||||
package org.dromara.aihr.personal.service;
|
||||
|
||||
import com.baomidou.mybatisplus.core.toolkit.IdWorker;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
|
||||
import org.dromara.aihr.personal.support.PersonalOwner;
|
||||
import org.dromara.common.core.exception.ServiceException;
|
||||
import org.dromara.system.service.ISysOssService;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.dao.EmptyResultDataAccessException;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.PlatformTransactionManager;
|
||||
import org.springframework.transaction.support.TransactionTemplate;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
@Slf4j
|
||||
@Service
|
||||
public class PersonalCleanupService {
|
||||
|
||||
private static final String ITEM_NOT_FOUND = "PERSONAL_ITEM_NOT_FOUND";
|
||||
private final JdbcTemplate jdbc;
|
||||
private final PersonalVectorStore vectorStore;
|
||||
private final OssCleanup ossCleanup;
|
||||
private final LongSupplier idSupplier;
|
||||
private final DbPhase dbPhase;
|
||||
private final int batchSize;
|
||||
|
||||
@Autowired
|
||||
public PersonalCleanupService(JdbcTemplate jdbc, PersonalVectorStore vectorStore, ISysOssService ossService,
|
||||
PlatformTransactionManager transactionManager,
|
||||
PersonalKnowledgeProperties properties) {
|
||||
TransactionTemplate transaction = new TransactionTemplate(transactionManager);
|
||||
this.jdbc = jdbc;
|
||||
this.vectorStore = vectorStore;
|
||||
this.ossCleanup = ossId -> ossService.deleteWithValidByIds(List.of(ossId), false);
|
||||
this.idSupplier = IdWorker::getId;
|
||||
this.dbPhase = action -> transaction.execute(status -> action.get());
|
||||
this.batchSize = Math.max(1, Math.min(100, properties.getCleanupBatchSize()));
|
||||
}
|
||||
|
||||
private PersonalCleanupService(JdbcTemplate jdbc, PersonalVectorStore vectorStore, OssCleanup ossCleanup,
|
||||
LongSupplier idSupplier, DbPhase dbPhase, int batchSize) {
|
||||
this.jdbc = jdbc;
|
||||
this.vectorStore = vectorStore;
|
||||
this.ossCleanup = ossCleanup;
|
||||
this.idSupplier = idSupplier;
|
||||
this.dbPhase = dbPhase;
|
||||
this.batchSize = batchSize;
|
||||
}
|
||||
|
||||
public static PersonalCleanupService forTest(JdbcTemplate jdbc, PersonalVectorStore vectorStore,
|
||||
OssCleanup ossCleanup, LongSupplier idSupplier,
|
||||
DbPhase dbPhase) {
|
||||
return new PersonalCleanupService(jdbc, vectorStore, ossCleanup, idSupplier, dbPhase, 20);
|
||||
}
|
||||
|
||||
public long requestDelete(PersonalOwner owner, long itemId) {
|
||||
requireOwner(owner);
|
||||
if (itemId <= 0) throw new ServiceException(ITEM_NOT_FOUND);
|
||||
return inDb(() -> {
|
||||
Map<String, Object> item;
|
||||
try {
|
||||
item = jdbc.queryForMap("""
|
||||
select i.id, i.space_id, i.size_bytes, i.status, i.oss_id, o.oss_id owned_oss_id
|
||||
from aihr_personal_item i
|
||||
left join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id
|
||||
and o.create_by = i.owner_user_id
|
||||
and cast(json_unquote(json_extract(o.ext1, '$.itemId')) as unsigned) = i.id
|
||||
where binary i.tenant_id = binary ? and i.owner_user_id = ? and i.id = ?
|
||||
and i.status in ('QUEUED','PARSING','READY','FAILED')
|
||||
for update
|
||||
""", owner.tenantId(), owner.userId(), itemId);
|
||||
} catch (EmptyResultDataAccessException ex) {
|
||||
throw new ServiceException(ITEM_NOT_FOUND);
|
||||
}
|
||||
if (item.get("oss_id") != null && item.get("owned_oss_id") == null) {
|
||||
throw new ServiceException(ITEM_NOT_FOUND);
|
||||
}
|
||||
String status = String.valueOf(item.get("status"));
|
||||
int hidden = jdbc.update("""
|
||||
update aihr_personal_item
|
||||
set status = 'DELETING', update_time = now()
|
||||
where tenant_id = ? and owner_user_id = ? and id = ? and status = ?
|
||||
""", owner.tenantId(), owner.userId(), itemId, status);
|
||||
if (hidden != 1) throw new ServiceException(ITEM_NOT_FOUND);
|
||||
long jobId = positiveId(idSupplier.getAsLong());
|
||||
int inserted = jdbc.update("""
|
||||
insert into aihr_personal_cleanup_job
|
||||
(id, tenant_id, owner_user_id, item_id, status, attempt_count, create_time, update_time)
|
||||
values (?, ?, ?, ?, 'PENDING', 0, now(), now())
|
||||
""", jobId, owner.tenantId(), owner.userId(), itemId);
|
||||
if (inserted != 1) throw new ServiceException("PERSONAL_CLEANUP_CREATE_FAILED");
|
||||
return jobId;
|
||||
});
|
||||
}
|
||||
|
||||
/** Executes external cleanup outside the database transaction. Every step is safe to repeat. */
|
||||
public void cleanup(long cleanupJobId) {
|
||||
if (cleanupJobId <= 0) return;
|
||||
List<Map<String, Object>> rows = jdbc.queryForList("""
|
||||
select j.id job_id, j.tenant_id, j.owner_user_id, j.item_id, j.status job_status,
|
||||
i.space_id, i.size_bytes, i.oss_id
|
||||
from aihr_personal_cleanup_job j
|
||||
join aihr_personal_item i
|
||||
on i.id = j.item_id and binary i.tenant_id = binary j.tenant_id
|
||||
and i.owner_user_id = j.owner_user_id
|
||||
where j.id = ? and j.status in ('PENDING','RETRY') and i.status = 'DELETING'
|
||||
limit 1
|
||||
""", cleanupJobId);
|
||||
if (rows.isEmpty()) return;
|
||||
CleanupItem item = cleanupItem(rows.get(0));
|
||||
try {
|
||||
vectorStore.deleteItem(item.owner(), item.itemId());
|
||||
jdbc.update("""
|
||||
delete from aihr_personal_fragment
|
||||
where tenant_id = ? and owner_user_id = ? and item_id = ?
|
||||
""", item.owner().tenantId(), item.owner().userId(), item.itemId());
|
||||
if (item.ossId() != null && item.ossId() > 0) {
|
||||
ossCleanup.delete(item.ossId());
|
||||
}
|
||||
inDb(() -> finalizeDeletion(item));
|
||||
} catch (RuntimeException ex) {
|
||||
jdbc.update("""
|
||||
update aihr_personal_cleanup_job
|
||||
set status = 'RETRY', attempt_count = attempt_count + 1,
|
||||
last_error = ?, update_time = now()
|
||||
where id = ? and tenant_id = ? and owner_user_id = ? and item_id = ? and status <> 'DONE'
|
||||
""", safeError(ex), item.jobId(), item.owner().tenantId(), item.owner().userId(), item.itemId());
|
||||
log.warn("event=personal_cleanup_retry jobId={} itemId={} exception={}", item.jobId(), item.itemId(),
|
||||
ex.getClass().getSimpleName());
|
||||
throw new ServiceException("PERSONAL_CLEANUP_RETRY_PENDING");
|
||||
}
|
||||
}
|
||||
|
||||
@Scheduled(fixedDelayString = "${aihr.personal.cleanup-delay-ms:60000}", scheduler = "personalTaskScheduler")
|
||||
public void poll() {
|
||||
List<Long> jobs = jdbc.query("""
|
||||
select id from aihr_personal_cleanup_job
|
||||
where status in ('PENDING','RETRY')
|
||||
order by update_time, id limit ?
|
||||
""", (rs, rowNum) -> rs.getLong("id"), batchSize);
|
||||
for (Long jobId : jobs) {
|
||||
try {
|
||||
cleanup(jobId);
|
||||
} catch (RuntimeException ignored) {
|
||||
// cleanup() persisted the retry state; later polls resume it.
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private Void finalizeDeletion(CleanupItem item) {
|
||||
int deleted = jdbc.update("""
|
||||
update aihr_personal_item
|
||||
set status = 'DELETED', summary = null, original_url = null, oss_id = null,
|
||||
error_code = null, error_message = null, deleted_at = now(), update_time = now()
|
||||
where tenant_id = ? and owner_user_id = ? and id = ? and status = 'DELETING'
|
||||
""", item.owner().tenantId(), item.owner().userId(), item.itemId());
|
||||
if (deleted == 1) {
|
||||
jdbc.update("""
|
||||
update aihr_personal_space
|
||||
set used_bytes = greatest(0, used_bytes - ?), item_count = greatest(0, item_count - 1),
|
||||
update_time = now()
|
||||
where tenant_id = ? and owner_user_id = ? and id = ?
|
||||
""", item.sizeBytes(), item.owner().tenantId(), item.owner().userId(), item.spaceId());
|
||||
}
|
||||
jdbc.update("""
|
||||
update aihr_personal_cleanup_job
|
||||
set status = 'DONE', attempt_count = attempt_count + 1, last_error = null,
|
||||
completed_at = now(), update_time = now()
|
||||
where id = ? and tenant_id = ? and owner_user_id = ? and item_id = ? and status <> 'DONE'
|
||||
""", item.jobId(), item.owner().tenantId(), item.owner().userId(), item.itemId());
|
||||
return null;
|
||||
}
|
||||
|
||||
private static CleanupItem cleanupItem(Map<String, Object> row) {
|
||||
PersonalOwner owner = new PersonalOwner(String.valueOf(row.get("tenant_id")), number(row, "owner_user_id"), null);
|
||||
Object oss = row.get("oss_id");
|
||||
return new CleanupItem(number(row, "job_id"), owner, number(row, "item_id"), number(row, "space_id"),
|
||||
number(row, "size_bytes"), oss instanceof Number value ? value.longValue() : null);
|
||||
}
|
||||
|
||||
private static long number(Map<String, Object> row, String key) {
|
||||
if (row.get(key) instanceof Number value) return value.longValue();
|
||||
throw new ServiceException("PERSONAL_CLEANUP_STATE_INVALID");
|
||||
}
|
||||
|
||||
private static void requireOwner(PersonalOwner owner) {
|
||||
if (owner == null || owner.tenantId() == null || owner.tenantId().isBlank() || owner.userId() <= 0) {
|
||||
throw new ServiceException("PERSONAL_OWNER_INVALID");
|
||||
}
|
||||
}
|
||||
|
||||
private static long positiveId(long value) {
|
||||
if (value <= 0) throw new ServiceException("PERSONAL_CLEANUP_CREATE_FAILED");
|
||||
return value;
|
||||
}
|
||||
|
||||
private static String safeError(RuntimeException ex) {
|
||||
String value = ex.getClass().getSimpleName();
|
||||
return value.length() <= 80 ? value : value.substring(0, 80);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private <T> T inDb(Supplier<T> action) {
|
||||
return (T) dbPhase.execute(action);
|
||||
}
|
||||
|
||||
@FunctionalInterface
|
||||
public interface OssCleanup {
|
||||
void delete(Long ossId);
|
||||
}
|
||||
|
||||
@FunctionalInterface
|
||||
public interface DbPhase {
|
||||
Object execute(Supplier<?> action);
|
||||
}
|
||||
|
||||
private record CleanupItem(long jobId, PersonalOwner owner, long itemId, long spaceId, long sizeBytes,
|
||||
Long ossId) {
|
||||
}
|
||||
}
|
||||
+50
-6
@@ -6,6 +6,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.ItemCreatedResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.TextItemRequest;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.UrlItemRequest;
|
||||
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
|
||||
import org.dromara.aihr.personal.support.PersonalOwner;
|
||||
import org.dromara.common.core.exception.ServiceException;
|
||||
@@ -49,6 +50,8 @@ public class PersonalIngestionService {
|
||||
private static final Set<String> SUPPORTED_FILE_SUFFIXES = Set.of(
|
||||
"txt", "md", "markdown", "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx"
|
||||
);
|
||||
private static final Set<String> INTERNAL_CAPTURE_SUFFIXES = Set.of("html", "txt", "pdf", "doc", "docx",
|
||||
"xls", "xlsx", "ppt", "pptx");
|
||||
|
||||
private final JdbcTemplate jdbcTemplate;
|
||||
private final PersonalSpaceService spaceService;
|
||||
@@ -103,7 +106,7 @@ public class PersonalIngestionService {
|
||||
byte[] bytes = request.content().getBytes(StandardCharsets.UTF_8);
|
||||
validateSize(bytes.length);
|
||||
return create(owner, "TEXT", cleanTitle(request.title(), "文字资料"), "txt", "text/plain", bytes,
|
||||
request.capturedAt(), request.tags());
|
||||
null, request.capturedAt(), request.tags());
|
||||
}
|
||||
|
||||
@Transactional(propagation = Propagation.NOT_SUPPORTED)
|
||||
@@ -120,7 +123,22 @@ public class PersonalIngestionService {
|
||||
validateSize(bytes.length);
|
||||
String originalName = safeFileName(file.getOriginalFilename());
|
||||
return create(owner, "FILE", cleanTitle(title, originalName), suffix(originalName),
|
||||
cleanMime(file.getContentType()), bytes, capturedAt, List.of());
|
||||
cleanMime(file.getContentType()), bytes, null, capturedAt, List.of());
|
||||
}
|
||||
|
||||
@Transactional(propagation = Propagation.NOT_SUPPORTED)
|
||||
public ItemCreatedResponse createUrl(PersonalOwner owner, UrlItemRequest request,
|
||||
PersonalUrlFetchService.FetchResult fetched) {
|
||||
validateOwner(owner);
|
||||
if (request == null || fetched == null || fetched.finalUri() == null || fetched.body() == null) {
|
||||
throw new ServiceException("PERSONAL_URL_FETCH_FAILED");
|
||||
}
|
||||
validateSize(fetched.body().length);
|
||||
String suffix = captureSuffix(fetched.contentType());
|
||||
String fallbackTitle = fetched.finalUri().getHost() == null ? "网页收藏" : fetched.finalUri().getHost();
|
||||
return create(owner, "URL", cleanTitle(request.title(), fallbackTitle), suffix,
|
||||
cleanMime(fetched.contentType()), fetched.body(), fetched.finalUri().toString(), request.capturedAt(),
|
||||
List.of());
|
||||
}
|
||||
|
||||
public void retry(PersonalOwner owner, long itemId) {
|
||||
@@ -141,7 +159,8 @@ public class PersonalIngestionService {
|
||||
}
|
||||
|
||||
private ItemCreatedResponse create(PersonalOwner owner, String sourceType, String title, String suffix,
|
||||
String mimeType, byte[] bytes, LocalDateTime capturedAt, List<String> tags) {
|
||||
String mimeType, byte[] bytes, String originalUrl, LocalDateTime capturedAt,
|
||||
List<String> tags) {
|
||||
String serviceKey = objectStore.requirePrivateService();
|
||||
long itemId = positiveId(idSupplier.getAsLong());
|
||||
long ossId = positiveId(idSupplier.getAsLong());
|
||||
@@ -152,7 +171,7 @@ public class PersonalIngestionService {
|
||||
String hash = sha256(bytes);
|
||||
|
||||
PhaseOne phaseOne = phaseTransaction.execute(status -> phaseOne(
|
||||
draft, sourceType, title, hash, capturedAt, tags));
|
||||
draft, sourceType, title, hash, originalUrl, capturedAt, tags));
|
||||
if (phaseOne == null) {
|
||||
throw new ServiceException("PERSONAL_ITEM_CREATE_FAILED");
|
||||
}
|
||||
@@ -177,7 +196,7 @@ public class PersonalIngestionService {
|
||||
}
|
||||
}
|
||||
|
||||
private PhaseOne phaseOne(UploadIntent draft, String sourceType, String title, String hash,
|
||||
private PhaseOne phaseOne(UploadIntent draft, String sourceType, String title, String hash, String originalUrl,
|
||||
LocalDateTime capturedAt, List<String> tags) {
|
||||
long spaceId = spaceService.reserve(draft.owner(), draft.sizeBytes());
|
||||
ItemCreatedResponse duplicate = duplicate(draft.owner(), spaceId, hash);
|
||||
@@ -202,6 +221,13 @@ public class PersonalIngestionService {
|
||||
""", intent.itemId(), intent.owner().tenantId(), spaceId, intent.owner().userId(), sourceType,
|
||||
title, intent.ossId(), intent.mimeType(), intent.sizeBytes(), hash, tagsJson(tags),
|
||||
capturedAt == null ? LocalDateTime.now() : capturedAt);
|
||||
if (originalUrl != null) {
|
||||
int linked = jdbcTemplate.update("""
|
||||
update aihr_personal_item set original_url = ?
|
||||
where tenant_id = ? and owner_user_id = ? and id = ? and status = 'QUEUED'
|
||||
""", originalUrl, intent.owner().tenantId(), intent.owner().userId(), intent.itemId());
|
||||
if (linked != 1) throw new ServiceException("PERSONAL_ITEM_CREATE_FAILED");
|
||||
}
|
||||
int counterUpdated = jdbcTemplate.update("""
|
||||
update aihr_personal_space
|
||||
set used_bytes = used_bytes + ?, item_count = item_count + 1, update_time = now()
|
||||
@@ -565,7 +591,9 @@ public class PersonalIngestionService {
|
||||
private static String objectKey(PersonalOwner owner, long itemId, String suffix) {
|
||||
validateOwner(owner); positiveId(itemId);
|
||||
String safeSuffix = suffix == null ? "" : suffix.toLowerCase(Locale.ROOT);
|
||||
if (!SUPPORTED_FILE_SUFFIXES.contains(safeSuffix)) throw new ServiceException("PERSONAL_FILE_UNSUPPORTED");
|
||||
if (!SUPPORTED_FILE_SUFFIXES.contains(safeSuffix) && !INTERNAL_CAPTURE_SUFFIXES.contains(safeSuffix)) {
|
||||
throw new ServiceException("PERSONAL_FILE_UNSUPPORTED");
|
||||
}
|
||||
return "personal/" + owner.tenantId() + "/" + owner.userId() + "/" + itemId + "/"
|
||||
+ UUID.randomUUID().toString().replace("-", "") + "." + safeSuffix;
|
||||
}
|
||||
@@ -610,6 +638,22 @@ public class PersonalIngestionService {
|
||||
int separator = mime.indexOf(';'); return separator < 0 ? mime : mime.substring(0, separator).trim();
|
||||
}
|
||||
|
||||
private static String captureSuffix(String mimeType) {
|
||||
return switch (cleanMime(mimeType)) {
|
||||
case "text/html" -> "html";
|
||||
case "text/plain", "text/markdown" -> "txt";
|
||||
case "application/pdf" -> "pdf";
|
||||
case "application/msword" -> "doc";
|
||||
case "application/vnd.ms-excel" -> "xls";
|
||||
case "application/vnd.ms-powerpoint" -> "ppt";
|
||||
case "application/vnd.openxmlformats-officedocument.wordprocessingml.document" -> "docx";
|
||||
case "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet" -> "xlsx";
|
||||
case "application/vnd.openxmlformats-officedocument.presentationml.presentation" -> "pptx";
|
||||
default -> throw new ServiceException("PERSONAL_URL_CONTENT_TYPE_UNSUPPORTED");
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
public interface PersonalObjectStore {
|
||||
String requirePrivateService();
|
||||
String uploadPhysical(String serviceKey, String objectKey, String mimeType, byte[] bytes);
|
||||
|
||||
+273
-2
@@ -1,27 +1,141 @@
|
||||
package org.dromara.aihr.personal.service;
|
||||
|
||||
import com.fasterxml.jackson.core.type.TypeReference;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.ChatMessageResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.CitationResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.DownloadUrlResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.ItemResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.PageResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.SessionDetailResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.SessionResponse;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.SpaceResponse;
|
||||
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
|
||||
import org.dromara.aihr.personal.support.PersonalOwner;
|
||||
import org.dromara.common.core.exception.ServiceException;
|
||||
import org.dromara.common.oss.core.OssClient;
|
||||
import org.dromara.common.oss.enums.AccessPolicyType;
|
||||
import org.dromara.common.oss.factory.OssFactory;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.dao.EmptyResultDataAccessException;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Propagation;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
import java.time.DateTimeException;
|
||||
import java.time.Duration;
|
||||
import java.time.LocalDate;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
@Service
|
||||
public class PersonalSpaceService {
|
||||
|
||||
private static final String QUOTA_EXCEEDED = "PERSONAL_SPACE_QUOTA_EXCEEDED";
|
||||
private static final String ITEM_NOT_FOUND = "PERSONAL_ITEM_NOT_FOUND";
|
||||
private static final String SESSION_NOT_FOUND = "PERSONAL_SESSION_NOT_FOUND";
|
||||
private static final Set<String> ITEM_STATUSES = Set.of("QUEUED", "PARSING", "READY", "FAILED");
|
||||
private static final Set<String> SOURCE_TYPES = Set.of("TEXT", "FILE", "IMAGE", "URL");
|
||||
|
||||
private final JdbcTemplate jdbcTemplate;
|
||||
private final PersonalKnowledgeProperties properties;
|
||||
private final ObjectMapper objectMapper;
|
||||
private final DownloadSigner downloadSigner;
|
||||
|
||||
public PersonalSpaceService(JdbcTemplate jdbcTemplate, PersonalKnowledgeProperties properties) {
|
||||
this(jdbcTemplate, properties, new ObjectMapper(), PersonalSpaceService::presign);
|
||||
}
|
||||
|
||||
@Autowired
|
||||
public PersonalSpaceService(JdbcTemplate jdbcTemplate, PersonalKnowledgeProperties properties,
|
||||
ObjectMapper objectMapper) {
|
||||
this(jdbcTemplate, properties, objectMapper, PersonalSpaceService::presign);
|
||||
}
|
||||
|
||||
private PersonalSpaceService(JdbcTemplate jdbcTemplate, PersonalKnowledgeProperties properties,
|
||||
ObjectMapper objectMapper, DownloadSigner downloadSigner) {
|
||||
this.jdbcTemplate = jdbcTemplate;
|
||||
this.properties = properties;
|
||||
this.objectMapper = objectMapper;
|
||||
this.downloadSigner = downloadSigner;
|
||||
}
|
||||
|
||||
public static PersonalSpaceService forTest(JdbcTemplate jdbcTemplate, PersonalKnowledgeProperties properties,
|
||||
ObjectMapper objectMapper, DownloadSigner downloadSigner) {
|
||||
return new PersonalSpaceService(jdbcTemplate, properties, objectMapper, downloadSigner);
|
||||
}
|
||||
|
||||
public SpaceResponse space(PersonalOwner owner) {
|
||||
requireOwner(owner);
|
||||
List<Map<String, Object>> rows = jdbcTemplate.queryForList("""
|
||||
select id, status, quota_bytes, used_bytes, item_count
|
||||
from aihr_personal_space
|
||||
where tenant_id = ? and owner_user_id = ?
|
||||
limit 1
|
||||
""", owner.tenantId(), owner.userId());
|
||||
if (rows.isEmpty()) {
|
||||
return new SpaceResponse(0L, "ACTIVE", defaultQuotaBytes(), 0L, 0);
|
||||
}
|
||||
Map<String, Object> row = rows.get(0);
|
||||
return new SpaceResponse(number(row, "id"), text(row, "status"), number(row, "quota_bytes"),
|
||||
number(row, "used_bytes"), Math.toIntExact(number(row, "item_count")));
|
||||
}
|
||||
|
||||
public PageResponse<ItemResponse> items(PersonalOwner owner, Integer pageNum, Integer pageSize, String status,
|
||||
String sourceType, LocalDate dateFrom, LocalDate dateTo, String keyword) {
|
||||
requireOwner(owner);
|
||||
int page = pageNum == null ? 1 : Math.max(1, pageNum);
|
||||
int size = pageSize == null ? 20 : Math.max(1, Math.min(100, pageSize));
|
||||
String normalizedStatus = normalized(status, ITEM_STATUSES, "PERSONAL_ITEM_STATUS_INVALID");
|
||||
String normalizedSource = normalized(sourceType, SOURCE_TYPES, "PERSONAL_ITEM_SOURCE_INVALID");
|
||||
validateDates(dateFrom, dateTo);
|
||||
StringBuilder where = new StringBuilder("""
|
||||
from aihr_personal_item
|
||||
where tenant_id = ? and owner_user_id = ?
|
||||
and status not in ('DELETING','DELETED')
|
||||
""");
|
||||
List<Object> args = new ArrayList<>(List.of(owner.tenantId(), owner.userId()));
|
||||
if (normalizedStatus != null) {
|
||||
where.append(" and status = ?");
|
||||
args.add(normalizedStatus);
|
||||
}
|
||||
if (normalizedSource != null) {
|
||||
where.append(" and source_type = ?");
|
||||
args.add(normalizedSource);
|
||||
}
|
||||
if (dateFrom != null) {
|
||||
where.append(" and captured_at >= ?");
|
||||
args.add(dateFrom.atStartOfDay());
|
||||
}
|
||||
if (dateTo != null) {
|
||||
where.append(" and captured_at < ?");
|
||||
args.add(dateTo.plusDays(1).atStartOfDay());
|
||||
}
|
||||
if (keyword != null && !keyword.isBlank()) {
|
||||
String value = keyword.trim();
|
||||
if (value.length() > 200) throw new ServiceException("PERSONAL_ITEM_FILTER_INVALID");
|
||||
where.append(" and title like ? escape '\\\\'");
|
||||
args.add("%" + escapeLike(value) + "%");
|
||||
}
|
||||
Long total = jdbcTemplate.queryForObject("select count(*)" + where, Long.class, args.toArray());
|
||||
List<Object> dataArgs = new ArrayList<>(args);
|
||||
dataArgs.add(size);
|
||||
dataArgs.add((page - 1L) * size);
|
||||
List<ItemResponse> rows = jdbcTemplate.query("""
|
||||
select id, source_type, title, original_url, mime_type, size_bytes, status,
|
||||
error_code, error_message, summary, tags_json, captured_at, parsed_at
|
||||
""" + where + " order by captured_at desc, id desc limit ? offset ?",
|
||||
(rs, rowNum) -> new ItemResponse(rs.getLong("id"), rs.getString("source_type"),
|
||||
rs.getString("title"), rs.getString("original_url"), rs.getString("mime_type"),
|
||||
rs.getLong("size_bytes"), rs.getString("status"), rs.getString("error_code"),
|
||||
rs.getString("error_message"), rs.getString("summary"), tags(rs.getString("tags_json")),
|
||||
rs.getObject("captured_at", LocalDateTime.class),
|
||||
rs.getObject("parsed_at", LocalDateTime.class)), dataArgs.toArray());
|
||||
return new PageResponse<>(List.copyOf(rows), total == null ? 0 : total, page, size);
|
||||
}
|
||||
|
||||
public Map<String, Object> item(PersonalOwner owner, long itemId) {
|
||||
@@ -32,13 +146,80 @@ public class PersonalSpaceService {
|
||||
tags_json, captured_at, parsed_at, create_time
|
||||
from aihr_personal_item
|
||||
where tenant_id = ? and owner_user_id = ? and id = ?
|
||||
and status <> 'DELETED'
|
||||
and status not in ('DELETING','DELETED')
|
||||
""", owner.tenantId(), owner.userId(), itemId);
|
||||
} catch (EmptyResultDataAccessException ex) {
|
||||
throw new ServiceException("PERSONAL_ITEM_NOT_FOUND");
|
||||
throw new ServiceException(ITEM_NOT_FOUND);
|
||||
}
|
||||
}
|
||||
|
||||
public ItemResponse itemResponse(PersonalOwner owner, long itemId) {
|
||||
return itemResponse(item(owner, itemId));
|
||||
}
|
||||
|
||||
public DownloadUrlResponse downloadUrl(PersonalOwner owner, long itemId) {
|
||||
requireOwner(owner);
|
||||
List<Map<String, Object>> rows = jdbcTemplate.queryForList("""
|
||||
select i.oss_id, o.file_name, o.service
|
||||
from aihr_personal_item i
|
||||
join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id
|
||||
and o.create_by = i.owner_user_id
|
||||
where i.tenant_id = ? and i.owner_user_id = ? and i.id = ?
|
||||
and i.status not in ('DELETING','DELETED')
|
||||
limit 1
|
||||
""", owner.tenantId(), owner.userId(), itemId);
|
||||
if (rows.isEmpty()) throw new ServiceException(ITEM_NOT_FOUND);
|
||||
Map<String, Object> row = rows.get(0);
|
||||
int minutes = Math.max(1, Math.min(60, properties.getDownloadUrlMinutes()));
|
||||
String url = downloadSigner.sign(text(row, "service"), text(row, "file_name"), Duration.ofMinutes(minutes));
|
||||
return new DownloadUrlResponse(url, LocalDateTime.now().plusMinutes(minutes));
|
||||
}
|
||||
|
||||
public List<SessionResponse> sessions(PersonalOwner owner) {
|
||||
requireOwner(owner);
|
||||
return List.copyOf(jdbcTemplate.query("""
|
||||
select id, title, default_scope, update_time
|
||||
from aihr_personal_chat_session
|
||||
where tenant_id = ? and owner_user_id = ? and status = 'ACTIVE'
|
||||
order by update_time desc, id desc limit 100
|
||||
""", (rs, rowNum) -> new SessionResponse(rs.getLong("id"), rs.getString("title"),
|
||||
rs.getString("default_scope"), rs.getObject("update_time", LocalDateTime.class)),
|
||||
owner.tenantId(), owner.userId()));
|
||||
}
|
||||
|
||||
public SessionDetailResponse session(PersonalOwner owner, long sessionId) {
|
||||
requireOwner(owner);
|
||||
List<Map<String, Object>> sessions = jdbcTemplate.queryForList("""
|
||||
select id, title from aihr_personal_chat_session
|
||||
where tenant_id = ? and owner_user_id = ? and id = ? and status = 'ACTIVE'
|
||||
limit 1
|
||||
""", owner.tenantId(), owner.userId(), sessionId);
|
||||
if (sessions.isEmpty()) throw new ServiceException(SESSION_NOT_FOUND);
|
||||
List<ChatMessageResponse> messages = jdbcTemplate.query("""
|
||||
select id, role, content, citations_json, create_time
|
||||
from aihr_personal_chat_message
|
||||
where tenant_id = ? and owner_user_id = ? and session_id = ?
|
||||
order by create_time, id
|
||||
""", (rs, rowNum) -> new ChatMessageResponse(rs.getLong("id"), rs.getString("role"),
|
||||
rs.getString("content"), citations(rs.getString("citations_json")),
|
||||
rs.getObject("create_time", LocalDateTime.class)), owner.tenantId(), owner.userId(), sessionId);
|
||||
return new SessionDetailResponse(sessionId, text(sessions.get(0), "title"), List.copyOf(messages));
|
||||
}
|
||||
|
||||
@Transactional
|
||||
public void deleteSession(PersonalOwner owner, long sessionId) {
|
||||
requireOwner(owner);
|
||||
int hidden = jdbcTemplate.update("""
|
||||
update aihr_personal_chat_session set status = 'DELETED', update_time = now()
|
||||
where tenant_id = ? and owner_user_id = ? and id = ? and status = 'ACTIVE'
|
||||
""", owner.tenantId(), owner.userId(), sessionId);
|
||||
if (hidden != 1) throw new ServiceException(SESSION_NOT_FOUND);
|
||||
jdbcTemplate.update("""
|
||||
delete from aihr_personal_chat_message
|
||||
where tenant_id = ? and owner_user_id = ? and session_id = ?
|
||||
""", owner.tenantId(), owner.userId(), sessionId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Ensures and locks the current owner's space, then validates that one future item of the given size can fit.
|
||||
* This check intentionally does not mutate usage counters. The future ingestion application service must call
|
||||
@@ -94,4 +275,94 @@ public class PersonalSpaceService {
|
||||
private long defaultQuotaBytes() {
|
||||
return Math.multiplyExact(properties.getMaxSpaceMb(), 1024L * 1024L);
|
||||
}
|
||||
|
||||
private ItemResponse itemResponse(Map<String, Object> row) {
|
||||
return new ItemResponse(number(row, "id"), text(row, "source_type"), text(row, "title"),
|
||||
nullableText(row, "original_url"), nullableText(row, "mime_type"), number(row, "size_bytes"),
|
||||
text(row, "status"), nullableText(row, "error_code"), nullableText(row, "error_message"),
|
||||
nullableText(row, "summary"), tags(row.get("tags_json")), dateTime(row.get("captured_at")),
|
||||
dateTime(row.get("parsed_at")));
|
||||
}
|
||||
|
||||
private List<String> tags(Object value) {
|
||||
if (value == null || String.valueOf(value).isBlank()) return List.of();
|
||||
try {
|
||||
return objectMapper.readValue(String.valueOf(value), new TypeReference<>() { });
|
||||
} catch (Exception ex) {
|
||||
return List.of();
|
||||
}
|
||||
}
|
||||
|
||||
private List<CitationResponse> citations(String value) {
|
||||
if (value == null || value.isBlank()) return List.of();
|
||||
try {
|
||||
return objectMapper.readValue(value, new TypeReference<>() { });
|
||||
} catch (Exception ex) {
|
||||
throw new ServiceException("PERSONAL_SESSION_DATA_INVALID");
|
||||
}
|
||||
}
|
||||
|
||||
private static String normalized(String value, Set<String> allowed, String error) {
|
||||
if (value == null || value.isBlank()) return null;
|
||||
String normalized = value.trim().toUpperCase(java.util.Locale.ROOT);
|
||||
if (!allowed.contains(normalized)) throw new ServiceException(error);
|
||||
return normalized;
|
||||
}
|
||||
|
||||
private static void validateDates(LocalDate from, LocalDate to) {
|
||||
if (from != null && to != null && from.isAfter(to)) throw new ServiceException("PERSONAL_ITEM_DATE_INVALID");
|
||||
if (to != null) {
|
||||
try {
|
||||
to.plusDays(1);
|
||||
} catch (DateTimeException ex) {
|
||||
throw new ServiceException("PERSONAL_ITEM_DATE_INVALID");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static void requireOwner(PersonalOwner owner) {
|
||||
if (owner == null || owner.tenantId() == null || owner.tenantId().isBlank() || owner.userId() <= 0) {
|
||||
throw new ServiceException("PERSONAL_OWNER_INVALID");
|
||||
}
|
||||
}
|
||||
|
||||
private static long number(Map<String, Object> row, String key) {
|
||||
if (row.get(key) instanceof Number number) return number.longValue();
|
||||
throw new ServiceException("PERSONAL_DATA_INVALID");
|
||||
}
|
||||
|
||||
private static String text(Map<String, Object> row, String key) {
|
||||
String value = nullableText(row, key);
|
||||
if (value == null) throw new ServiceException("PERSONAL_DATA_INVALID");
|
||||
return value;
|
||||
}
|
||||
|
||||
private static String nullableText(Map<String, Object> row, String key) {
|
||||
Object value = row.get(key);
|
||||
return value == null || String.valueOf(value).isBlank() ? null : String.valueOf(value);
|
||||
}
|
||||
|
||||
private static LocalDateTime dateTime(Object value) {
|
||||
if (value == null) return null;
|
||||
if (value instanceof LocalDateTime time) return time;
|
||||
if (value instanceof java.sql.Timestamp time) return time.toLocalDateTime();
|
||||
throw new ServiceException("PERSONAL_DATA_INVALID");
|
||||
}
|
||||
|
||||
private static String escapeLike(String value) {
|
||||
return value.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_");
|
||||
}
|
||||
|
||||
private static String presign(String service, String objectKey, Duration duration) {
|
||||
OssClient client = service == null || service.isBlank() ? OssFactory.instance() : OssFactory.instance(service);
|
||||
if (client == null || client.getAccessPolicy() != AccessPolicyType.PRIVATE) {
|
||||
throw new ServiceException("PERSONAL_OSS_NOT_PRIVATE");
|
||||
}
|
||||
return client.createPresignedGetUrl(objectKey, duration);
|
||||
}
|
||||
|
||||
@FunctionalInterface
|
||||
public interface DownloadSigner {
|
||||
String sign(String service, String objectKey, Duration duration);
|
||||
}
|
||||
}
|
||||
|
||||
+69
@@ -0,0 +1,69 @@
|
||||
package org.dromara.aihr.personal;
|
||||
|
||||
import cn.dev33.satoken.annotation.SaIgnore;
|
||||
import org.dromara.aihr.personal.controller.PersonalAssistantController;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.TextItemRequest;
|
||||
import org.dromara.aihr.personal.domain.PersonalAssistantDto.UrlItemRequest;
|
||||
import org.dromara.aihr.personal.service.PersonalAnswerService;
|
||||
import org.dromara.aihr.personal.service.PersonalCleanupService;
|
||||
import org.dromara.aihr.personal.service.PersonalIngestionService;
|
||||
import org.dromara.aihr.personal.service.PersonalRetrievalService;
|
||||
import org.dromara.aihr.personal.service.PersonalSpaceService;
|
||||
import org.dromara.aihr.personal.service.PersonalUrlFetchService;
|
||||
import org.dromara.aihr.personal.support.PersonalOwner;
|
||||
import org.dromara.aihr.personal.support.PersonalOwnerProvider;
|
||||
import org.junit.jupiter.api.Tag;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.times;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
@Tag("dev")
|
||||
class PersonalAssistantControllerTest {
|
||||
|
||||
@Test
|
||||
void everyP0EntryDerivesOwnerAndControllerHasNoAnonymousBypass() {
|
||||
PersonalOwnerProvider owners = mock(PersonalOwnerProvider.class);
|
||||
PersonalOwner owner = new PersonalOwner("000000", 101L, null);
|
||||
when(owners.current()).thenReturn(owner);
|
||||
PersonalSpaceService spaces = mock(PersonalSpaceService.class);
|
||||
PersonalIngestionService ingestion = mock(PersonalIngestionService.class);
|
||||
PersonalUrlFetchService urls = mock(PersonalUrlFetchService.class);
|
||||
PersonalRetrievalService retrieval = mock(PersonalRetrievalService.class);
|
||||
PersonalAnswerService answers = mock(PersonalAnswerService.class);
|
||||
PersonalCleanupService cleanup = mock(PersonalCleanupService.class);
|
||||
when(urls.fetch("https://example.com/a")).thenReturn(new PersonalUrlFetchService.FetchResult(
|
||||
java.net.URI.create("https://example.com/a"), 200, "text/plain", "a".getBytes(),
|
||||
java.time.Instant.now(), "hash"));
|
||||
|
||||
PersonalAssistantController controller = new PersonalAssistantController(owners, spaces, ingestion, urls,
|
||||
retrieval, answers, cleanup);
|
||||
controller.space();
|
||||
controller.items(1, 20, null, null, null, null, null);
|
||||
controller.createText(new TextItemRequest("note", "body", null, List.of()));
|
||||
controller.createFile(mock(MultipartFile.class), null, null);
|
||||
controller.createUrl(new UrlItemRequest("https://example.com/a", null, null));
|
||||
controller.item(9L);
|
||||
controller.retry(9L);
|
||||
controller.deleteItem(9L);
|
||||
controller.downloadUrl(9L);
|
||||
controller.search(null);
|
||||
controller.ask(null);
|
||||
controller.sessions();
|
||||
controller.session(3L);
|
||||
controller.deleteSession(3L);
|
||||
|
||||
verify(owners, times(14)).current();
|
||||
assertFalse(PersonalAssistantController.class.isAnnotationPresent(SaIgnore.class));
|
||||
for (var method : PersonalAssistantController.class.getDeclaredMethods()) {
|
||||
assertFalse(method.isAnnotationPresent(SaIgnore.class), method.getName());
|
||||
}
|
||||
}
|
||||
}
|
||||
+124
@@ -0,0 +1,124 @@
|
||||
package org.dromara.aihr.personal;
|
||||
|
||||
import org.dromara.aihr.personal.service.PersonalCleanupService;
|
||||
import org.dromara.aihr.personal.service.PersonalVectorStore;
|
||||
import org.dromara.aihr.personal.support.PersonalOwner;
|
||||
import org.dromara.common.core.exception.ServiceException;
|
||||
import org.junit.jupiter.api.Tag;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.InOrder;
|
||||
import org.springframework.dao.EmptyResultDataAccessException;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.contains;
|
||||
import static org.mockito.ArgumentMatchers.eq;
|
||||
import static org.mockito.Mockito.inOrder;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.never;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
import static org.mockito.Mockito.doThrow;
|
||||
|
||||
@Tag("dev")
|
||||
class PersonalCleanupServiceTest {
|
||||
|
||||
private static final PersonalOwner OWNER = new PersonalOwner("000000", 101L, null);
|
||||
|
||||
@Test
|
||||
void deletionHidesItemAndPersistsJobBeforeExternalCleanup() {
|
||||
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||
PersonalVectorStore vectors = mock(PersonalVectorStore.class);
|
||||
PersonalCleanupService.OssCleanup oss = mock(PersonalCleanupService.OssCleanup.class);
|
||||
when(jdbc.queryForMap(contains("for update"), eq("000000"), eq(101L), eq(9L)))
|
||||
.thenReturn(item("READY"));
|
||||
when(jdbc.update(contains("set status = 'DELETING'"), eq("000000"), eq(101L), eq(9L), eq("READY")))
|
||||
.thenReturn(1);
|
||||
when(jdbc.update(contains("insert into aihr_personal_cleanup_job"), any(), eq("000000"), eq(101L), eq(9L)))
|
||||
.thenReturn(1);
|
||||
PersonalCleanupService service = PersonalCleanupService.forTest(jdbc, vectors, oss, () -> 7001L,
|
||||
action -> action.get());
|
||||
|
||||
assertEquals(7001L, service.requestDelete(OWNER, 9L));
|
||||
|
||||
InOrder order = inOrder(jdbc);
|
||||
order.verify(jdbc).queryForMap(contains("for update"), eq("000000"), eq(101L), eq(9L));
|
||||
order.verify(jdbc).update(contains("set status = 'DELETING'"), eq("000000"), eq(101L), eq(9L), eq("READY"));
|
||||
order.verify(jdbc).update(contains("insert into aihr_personal_cleanup_job"), eq(7001L), eq("000000"),
|
||||
eq(101L), eq(9L));
|
||||
verify(vectors, never()).deleteItem(any(), any(Long.class));
|
||||
verify(oss, never()).delete(any(Long.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
void foreignItemUsesStableNotFoundWithoutCreatingJob() {
|
||||
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||
when(jdbc.queryForMap(anyString(), eq("000000"), eq(202L), eq(9L)))
|
||||
.thenThrow(new EmptyResultDataAccessException(1));
|
||||
PersonalCleanupService service = PersonalCleanupService.forTest(jdbc, mock(PersonalVectorStore.class),
|
||||
mock(PersonalCleanupService.OssCleanup.class), () -> 7001L, action -> action.get());
|
||||
|
||||
ServiceException error = assertThrows(ServiceException.class,
|
||||
() -> service.requestDelete(new PersonalOwner("000000", 202L, null), 9L));
|
||||
|
||||
assertEquals("PERSONAL_ITEM_NOT_FOUND", error.getMessage());
|
||||
verify(jdbc, never()).update(contains("insert into aihr_personal_cleanup_job"), any(), any(), any(), any());
|
||||
}
|
||||
|
||||
@Test
|
||||
void cleanupUsesFixedOrderAndIsIdempotent() {
|
||||
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||
PersonalVectorStore vectors = mock(PersonalVectorStore.class);
|
||||
PersonalCleanupService.OssCleanup oss = mock(PersonalCleanupService.OssCleanup.class);
|
||||
when(jdbc.queryForList(contains("from aihr_personal_cleanup_job j"), eq(7001L)))
|
||||
.thenReturn(List.of(cleanupRow()), List.of());
|
||||
when(jdbc.update(contains("delete from aihr_personal_fragment"), eq("000000"), eq(101L), eq(9L)))
|
||||
.thenReturn(3);
|
||||
when(jdbc.update(contains("set status = 'DELETED'"), eq("000000"), eq(101L), eq(9L)))
|
||||
.thenReturn(1);
|
||||
PersonalCleanupService service = PersonalCleanupService.forTest(jdbc, vectors, oss, () -> 7001L,
|
||||
action -> action.get());
|
||||
|
||||
service.cleanup(7001L);
|
||||
service.cleanup(7001L);
|
||||
|
||||
InOrder order = inOrder(vectors, jdbc, oss);
|
||||
order.verify(vectors).deleteItem(OWNER, 9L);
|
||||
order.verify(jdbc).update(contains("delete from aihr_personal_fragment"), eq("000000"), eq(101L), eq(9L));
|
||||
order.verify(oss).delete(55L);
|
||||
order.verify(jdbc).update(contains("set status = 'DELETED'"), eq("000000"), eq(101L), eq(9L));
|
||||
}
|
||||
|
||||
@Test
|
||||
void externalFailureKeepsDeletingAndPersistsRetry() {
|
||||
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||
PersonalVectorStore vectors = mock(PersonalVectorStore.class);
|
||||
when(jdbc.queryForList(contains("from aihr_personal_cleanup_job j"), eq(7001L)))
|
||||
.thenReturn(List.of(cleanupRow()));
|
||||
doThrow(new IllegalStateException("qdrant unavailable")).when(vectors).deleteItem(OWNER, 9L);
|
||||
PersonalCleanupService service = PersonalCleanupService.forTest(jdbc, vectors,
|
||||
mock(PersonalCleanupService.OssCleanup.class), () -> 7001L, action -> action.get());
|
||||
|
||||
ServiceException error = assertThrows(ServiceException.class, () -> service.cleanup(7001L));
|
||||
|
||||
assertEquals("PERSONAL_CLEANUP_RETRY_PENDING", error.getMessage());
|
||||
verify(jdbc).update(contains("set status = 'RETRY'"), eq("IllegalStateException"), eq(7001L),
|
||||
eq("000000"), eq(101L), eq(9L));
|
||||
verify(jdbc, never()).update(contains("set status = 'DELETED'"), any(), any(), any());
|
||||
}
|
||||
|
||||
private static Map<String, Object> item(String status) {
|
||||
return Map.of("id", 9L, "space_id", 3L, "size_bytes", 100L, "status", status);
|
||||
}
|
||||
|
||||
private static Map<String, Object> cleanupRow() {
|
||||
return Map.of("job_id", 7001L, "tenant_id", "000000", "owner_user_id", 101L, "item_id", 9L,
|
||||
"space_id", 3L, "size_bytes", 100L, "oss_id", 55L, "job_status", "PENDING");
|
||||
}
|
||||
}
|
||||
+6
-1
@@ -28,8 +28,9 @@ class PersonalSchemaContractTest {
|
||||
String fragment = tableDefinition(sql, "aihr_personal_fragment");
|
||||
String session = tableDefinition(sql, "aihr_personal_chat_session");
|
||||
String message = tableDefinition(sql, "aihr_personal_chat_message");
|
||||
String cleanup = tableDefinition(sql, "aihr_personal_cleanup_job");
|
||||
|
||||
for (String definition : new String[] {space, item, fragment, session, message}) {
|
||||
for (String definition : new String[] {space, item, fragment, session, message, cleanup}) {
|
||||
assertTrue(definition.contains("`owner_user_id` bigint not null"),
|
||||
"Every personal table must carry a non-null owner_user_id");
|
||||
}
|
||||
@@ -88,6 +89,10 @@ class PersonalSchemaContractTest {
|
||||
assertTrue(message.contains(
|
||||
"key `idx_personal_message_owner` (`tenant_id`, `owner_user_id`, `session_id`, `create_time`)"));
|
||||
|
||||
assertTrue(cleanup.contains("`status` varchar(20) not null default 'pending'"));
|
||||
assertTrue(cleanup.contains("unique key `uk_personal_cleanup_item` (`tenant_id`, `owner_user_id`, `item_id`)"));
|
||||
assertTrue(cleanup.contains("key `idx_personal_cleanup_status` (`status`, `update_time`)"));
|
||||
|
||||
assertFalse(sql.contains("alter table aihr_knowledge_fragment"),
|
||||
"Personal schema must not mutate enterprise knowledge tables");
|
||||
assertFalse(sql.contains("alter table `aihr_knowledge_fragment`"),
|
||||
|
||||
+44
@@ -27,11 +27,13 @@ import org.springframework.transaction.support.TransactionTemplate;
|
||||
import org.mockito.InOrder;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.List;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.contains;
|
||||
import static org.mockito.ArgumentMatchers.eq;
|
||||
import static org.mockito.Mockito.mock;
|
||||
@@ -44,6 +46,48 @@ import static org.mockito.Mockito.when;
|
||||
@Tag("dev")
|
||||
class PersonalSpaceServiceTest {
|
||||
|
||||
@Test
|
||||
void itemListAndSessionsAreAlwaysOwnerScoped() {
|
||||
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||
when(jdbc.queryForObject(contains("from aihr_personal_item"), eq(Long.class), eq("000000"), eq(101L)))
|
||||
.thenReturn(0L);
|
||||
when(jdbc.query(anyString(), any(org.springframework.jdbc.core.RowMapper.class),
|
||||
eq("000000"), eq(101L), eq(20), eq(0L))).thenReturn(List.of());
|
||||
when(jdbc.query(contains("from aihr_personal_chat_session"),
|
||||
any(org.springframework.jdbc.core.RowMapper.class), eq("000000"), eq(101L))).thenReturn(List.of());
|
||||
PersonalSpaceService service = new PersonalSpaceService(jdbc, properties());
|
||||
PersonalOwner owner = new PersonalOwner("000000", 101L, null);
|
||||
|
||||
service.items(owner, 1, 20, null, null, null, null, null);
|
||||
service.sessions(owner);
|
||||
|
||||
verify(jdbc).queryForObject(contains("tenant_id = ? and owner_user_id = ?"), eq(Long.class),
|
||||
eq("000000"), eq(101L));
|
||||
verify(jdbc).query(contains("from aihr_personal_chat_session"),
|
||||
any(org.springframework.jdbc.core.RowMapper.class), eq("000000"), eq(101L));
|
||||
}
|
||||
|
||||
@Test
|
||||
void downloadChecksOwnerBeforeSigningPrivateObject() {
|
||||
JdbcTemplate jdbc = mock(JdbcTemplate.class);
|
||||
when(jdbc.queryForList(contains("join sys_oss"), eq("000000"), eq(101L), eq(9L)))
|
||||
.thenReturn(List.of(Map.of("oss_id", 55L, "file_name", "personal/000000/101/9/a.pdf",
|
||||
"service", "private")));
|
||||
PersonalSpaceService.DownloadSigner signer = mock(PersonalSpaceService.DownloadSigner.class);
|
||||
when(signer.sign(eq("private"), eq("personal/000000/101/9/a.pdf"), any(java.time.Duration.class)))
|
||||
.thenReturn("https://signed.example/a");
|
||||
PersonalSpaceService service = PersonalSpaceService.forTest(jdbc, properties(),
|
||||
new com.fasterxml.jackson.databind.ObjectMapper(), signer);
|
||||
|
||||
assertEquals("https://signed.example/a",
|
||||
service.downloadUrl(new PersonalOwner("000000", 101L, null), 9L).url());
|
||||
|
||||
verify(jdbc).queryForList(contains("i.tenant_id = ? and i.owner_user_id = ? and i.id = ?"),
|
||||
eq("000000"), eq(101L), eq(9L));
|
||||
verify(signer).sign(eq("private"), eq("personal/000000/101/9/a.pdf"),
|
||||
eq(java.time.Duration.ofMinutes(5)));
|
||||
}
|
||||
|
||||
@Test
|
||||
void propertiesHaveExplicitSafeDefaults() {
|
||||
PersonalKnowledgeProperties properties = properties();
|
||||
|
||||
@@ -96,3 +96,19 @@ CREATE TABLE IF NOT EXISTS `aihr_personal_chat_message` (
|
||||
KEY `idx_personal_message_session` (`session_id`, `create_time`),
|
||||
KEY `idx_personal_message_owner` (`tenant_id`, `owner_user_id`, `session_id`, `create_time`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='个人AI助理对话消息';
|
||||
|
||||
CREATE TABLE IF NOT EXISTS `aihr_personal_cleanup_job` (
|
||||
`id` bigint NOT NULL COMMENT '清理任务ID',
|
||||
`tenant_id` varchar(20) NOT NULL COMMENT '租户编号',
|
||||
`owner_user_id` bigint NOT NULL COMMENT '资料所属用户ID',
|
||||
`item_id` bigint NOT NULL COMMENT '资料ID',
|
||||
`status` varchar(20) NOT NULL DEFAULT 'PENDING' COMMENT 'PENDING/RETRY/DONE',
|
||||
`attempt_count` int NOT NULL DEFAULT 0 COMMENT '执行次数',
|
||||
`last_error` varchar(100) DEFAULT NULL COMMENT '脱敏后的错误类型',
|
||||
`completed_at` datetime DEFAULT NULL COMMENT '完成时间',
|
||||
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
|
||||
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
|
||||
PRIMARY KEY (`id`),
|
||||
UNIQUE KEY `uk_personal_cleanup_item` (`tenant_id`, `owner_user_id`, `item_id`),
|
||||
KEY `idx_personal_cleanup_status` (`status`, `update_time`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='个人AI助理资料清理任务';
|
||||
|
||||
Reference in New Issue
Block a user