test(personal): cover PDF and public URL capture

This commit is contained in:
2026-07-12 16:26:59 +08:00
parent c7756f86e2
commit 50e2b9e5ff
7 changed files with 421 additions and 72 deletions
@@ -8,6 +8,7 @@ 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.dromara.common.tenant.helper.TenantHelper;
import org.dromara.system.domain.vo.SysOssVo;
import org.dromara.system.service.ISysOssService;
import org.springframework.beans.factory.annotation.Autowired;
@@ -33,6 +34,7 @@ public class PersonalIngestionWorker {
private final KnowledgeDocumentParser parser;
private final TransactionTemplate transactionTemplate;
private final StoredObjectReader objectReader;
private final TenantRunner tenantRunner;
private final long maxInputBytes;
private final int chunkSize;
private final int chunkOverlap;
@@ -46,17 +48,18 @@ public class PersonalIngestionWorker {
this(jdbcTemplate, parser, new TransactionTemplate(transactionManager),
defaultReader(ossService, PersonalIngestionWorker::ossClient), configuredMaxBytes(properties),
properties.getChunkSize(), properties.getChunkOverlap(), properties.getParsingLeaseMinutes(),
properties.getMaxParseAttempts());
properties.getMaxParseAttempts(), PersonalIngestionWorker::runInTenant);
}
private PersonalIngestionWorker(JdbcTemplate jdbcTemplate, KnowledgeDocumentParser parser,
TransactionTemplate transactionTemplate, StoredObjectReader objectReader,
long maxInputBytes, int chunkSize, int chunkOverlap,
int parsingLeaseMinutes, int maxParseAttempts) {
int parsingLeaseMinutes, int maxParseAttempts, TenantRunner tenantRunner) {
this.jdbcTemplate = jdbcTemplate;
this.parser = parser;
this.transactionTemplate = transactionTemplate;
this.objectReader = objectReader;
this.tenantRunner = tenantRunner;
this.maxInputBytes = maxInputBytes;
if (chunkSize <= 0 || chunkOverlap < 0 || chunkOverlap >= chunkSize
|| parsingLeaseMinutes <= 0 || maxParseAttempts <= 0) {
@@ -73,7 +76,15 @@ public class PersonalIngestionWorker {
TransactionTemplate transactionTemplate,
StoredObjectReader objectReader) {
return new PersonalIngestionWorker(jdbcTemplate, parser, transactionTemplate, objectReader,
20L * 1024 * 1024, 800, 120, 15, 3);
20L * 1024 * 1024, 800, 120, 15, 3, (tenantId, operation) -> operation.execute());
}
public static PersonalIngestionWorker forTest(JdbcTemplate jdbcTemplate, ISysOssService ossService,
KnowledgeDocumentParser parser,
TransactionTemplate transactionTemplate,
StoredObjectReader objectReader, TenantRunner tenantRunner) {
return new PersonalIngestionWorker(jdbcTemplate, parser, transactionTemplate, objectReader,
20L * 1024 * 1024, 800, 120, 15, 3, tenantRunner);
}
public static StoredObjectReader objectReaderForTest(ISysOssService ossService,
@@ -157,8 +168,8 @@ public class PersonalIngestionWorker {
int attemptVersion = Math.addExact(item.attemptCount(), 1);
try {
StoredObject stored = objectReader.read(
item.ossId(), ownerObjectPrefix(item), item.ownerUserId(), maxInputBytes);
StoredObject stored = tenantRunner.execute(item.tenantId(), () -> objectReader.read(
item.ossId(), ownerObjectPrefix(item), item.ownerUserId(), maxInputBytes));
ParsedDocument document = parser.parse(stored.fileName(), item.mimeType(), stored.bytes());
List<String> chunks = document.chunks(chunkSize, chunkOverlap);
if (chunks.isEmpty()) {
@@ -345,6 +356,16 @@ public class PersonalIngestionWorker {
StoredObject read(long ossId, String expectedPrefix, long ownerUserId, long maxBytes) throws Exception;
}
@FunctionalInterface
public interface TenantOperation {
StoredObject execute() throws Exception;
}
@FunctionalInterface
public interface TenantRunner {
StoredObject execute(String tenantId, TenantOperation operation) throws Exception;
}
@FunctionalInterface
public interface OssClientProvider {
OssClient get(String configKey);
@@ -373,4 +394,15 @@ public class PersonalIngestionWorker {
? OssFactory.instance()
: OssFactory.instance(configKey);
}
private static StoredObject runInTenant(String tenantId, TenantOperation operation) throws Exception {
String previous = TenantHelper.getDynamic();
TenantHelper.setDynamic(tenantId);
try {
return operation.execute();
} finally {
TenantHelper.clearDynamic();
if (previous != null && !previous.isBlank()) TenantHelper.setDynamic(previous);
}
}
}
@@ -34,6 +34,7 @@ import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Hashtable;
import java.util.HashSet;
import java.util.HexFormat;
@@ -93,7 +94,9 @@ public class PersonalUrlFetchService {
@Autowired
public PersonalUrlFetchService(PersonalKnowledgeProperties properties) {
this(properties, new DeadlineDnsResolver(DNS_EXECUTOR, new JndiDnsQuery()), new RawSocketFetcher());
this(properties, new FallbackResolver(
new DeadlineDnsResolver(DNS_EXECUTOR, new JndiDnsQuery()),
new DeadlineSystemResolver(DNS_EXECUTOR, InetAddress::getAllByName)), new RawSocketFetcher());
}
private PersonalUrlFetchService(PersonalKnowledgeProperties properties, Resolver resolver, Fetcher fetcher) {
@@ -359,6 +362,69 @@ public class PersonalUrlFetchService {
List<InetAddress> resolve(String host, long deadlineNanos) throws IOException;
}
/** Falls back only when the primary resolver is unavailable; empty or unsafe answers remain fail-closed. */
static final class FallbackResolver implements Resolver {
private final Resolver primary;
private final Resolver fallback;
FallbackResolver(Resolver primary, Resolver fallback) {
this.primary = primary;
this.fallback = fallback;
}
@Override
public List<InetAddress> resolve(String host, long deadlineNanos) throws IOException {
try {
return primary.resolve(host, deadlineNanos);
} catch (IOException primaryFailure) {
if (deadlineNanos - System.nanoTime() <= 0) throw primaryFailure;
return fallback.resolve(host, deadlineNanos);
}
}
}
@FunctionalInterface
interface SystemAddressQuery {
InetAddress[] resolve(String host) throws IOException;
}
/** Bounds the JVM/system resolver with the same end-to-end deadline used by the fetch. */
static final class DeadlineSystemResolver implements Resolver {
private final ExecutorService executor;
private final SystemAddressQuery query;
DeadlineSystemResolver(ExecutorService executor, SystemAddressQuery query) {
this.executor = executor;
this.query = query;
}
@Override
public List<InetAddress> resolve(String host, long deadlineNanos) throws IOException {
long remaining = deadlineNanos - System.nanoTime();
if (remaining <= 0) throw new IOException("resolution deadline exceeded");
Future<InetAddress[]> future;
try {
future = executor.submit(() -> query.resolve(host));
} catch (RejectedExecutionException ex) {
throw new IOException("resolution unavailable");
}
try {
InetAddress[] addresses = future.get(remaining, TimeUnit.NANOSECONDS);
return addresses == null ? List.of() : List.copyOf(Arrays.asList(addresses));
} catch (TimeoutException ex) {
cancelAndPurge(executor, future);
throw new IOException("resolution deadline exceeded");
} catch (InterruptedException ex) {
cancelAndPurge(executor, future);
Thread.currentThread().interrupt();
throw new IOException("resolution interrupted");
} catch (ExecutionException ex) {
cancelAndPurge(executor, future);
throw new IOException("resolution failed");
}
}
}
@FunctionalInterface
interface DnsQuery {
List<String> resolve(String host, int timeoutMillis, int retries) throws NamingException;
@@ -20,6 +20,7 @@ import java.sql.Timestamp;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -36,6 +37,30 @@ import static org.mockito.Mockito.when;
@Tag("dev")
class PersonalIngestionWorkerTest {
@Test
void workerReadsPrivateObjectInsideItemTenantScope() throws Exception {
JdbcTemplate jdbc = mock(JdbcTemplate.class);
KnowledgeDocumentParser parser = mock(KnowledgeDocumentParser.class);
AtomicReference<String> tenantSeen = new AtomicReference<>();
when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item()));
when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L), eq(0))).thenReturn(1);
when(parser.parse(any(), any(), any(byte[].class)))
.thenThrow(new KnowledgeDocumentParser.ParseException(
KnowledgeDocumentParser.Failure.INVALID, "stop after tenant-scoped read"));
PersonalIngestionWorker worker = PersonalIngestionWorker.forTest(
jdbc, mock(ISysOssService.class), parser, immediateTransactions(),
(ossId, prefix, ownerUserId, maxBytes) ->
new PersonalIngestionWorker.StoredObject("sample.pdf", new byte[]{1}),
(tenantId, operation) -> {
tenantSeen.set(tenantId);
return operation.execute();
});
assertTrue(worker.processNext());
assertEquals("000000", tenantSeen.get());
}
@Test
void workerQueueOnlySelectsReadyUploadIntents() {
JdbcTemplate jdbc = mock(JdbcTemplate.class);
@@ -245,6 +245,66 @@ class PersonalUrlFetchServiceTest {
assertEquals("0", environment.get("com.sun.jndi.dns.timeout.retries"));
}
@Test
void fallsBackToBoundedSystemDnsOnlyWhenJndiResolutionFails() throws Exception {
AtomicInteger fallbackCalls = new AtomicInteger();
var resolver = new PersonalUrlFetchService.FallbackResolver(
(host, deadline) -> { throw new IOException("JNDI unavailable"); },
(host, deadline) -> { fallbackCalls.incrementAndGet(); return List.of(PUBLIC); });
var service = fixture(resolver, request -> ok("text/plain", "ok"));
assertEquals("https://example.com/", service.validate("https://example.com/").toString());
assertEquals(1, fallbackCalls.get());
fallbackCalls.set(0);
var emptyPrimary = new PersonalUrlFetchService.FallbackResolver(
(host, deadline) -> List.of(),
(host, deadline) -> { fallbackCalls.incrementAndGet(); return List.of(PUBLIC); });
assertCode("PERSONAL_URL_BLOCKED", () -> fixture(emptyPrimary, request -> ok("text/plain", "ok"))
.validate("https://example.com/"), "empty primary result must fail closed");
assertEquals(0, fallbackCalls.get());
}
@Test
void validatesEveryFallbackAddressAndFailsClosedWhenFallbackFails() {
var privateFallback = new PersonalUrlFetchService.FallbackResolver(
(host, deadline) -> { throw new IOException("JNDI unavailable"); },
(host, deadline) -> List.of(PUBLIC, address("127.0.0.1")));
assertCode("PERSONAL_URL_BLOCKED", () -> fixture(privateFallback, request -> ok("text/plain", "ok"))
.validate("https://example.com/"), "mixed fallback addresses");
var failedFallback = new PersonalUrlFetchService.FallbackResolver(
(host, deadline) -> { throw new IOException("JNDI unavailable"); },
(host, deadline) -> { throw new IOException("system DNS unavailable"); });
assertCode("PERSONAL_URL_BLOCKED", () -> fixture(failedFallback, request -> ok("text/plain", "ok"))
.validate("https://example.com/"), "fallback failure");
}
@Test
void boundsSystemDnsFallbackByTheSharedDeadline() throws Exception {
ExecutorService executor = boundedExecutor("system-dns-wall-test");
CountDownLatch entered = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
try {
var resolver = new PersonalUrlFetchService.DeadlineSystemResolver(executor, host -> {
entered.countDown();
boolean done = false;
while (!done) {
try { release.await(); done = true; }
catch (InterruptedException ignored) { }
}
return new InetAddress[]{PUBLIC};
});
long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(30);
assertThrows(IOException.class, () -> resolver.resolve("example.com", deadline));
assertTrue(entered.await(1, TimeUnit.SECONDS));
} finally {
release.countDown();
executor.shutdownNow();
assertTrue(executor.awaitTermination(1, TimeUnit.SECONDS));
}
}
@Test
void outerDnsDeadlineReturnsWhenQueryIgnoresInterrupt() throws Exception {
ExecutorService executor = boundedExecutor("dns-wall-test");