fix(personal): bound DNS fallback and smoke redirects
This commit is contained in:
+170
-139
@@ -5,13 +5,6 @@ import org.dromara.common.core.exception.ServiceException;
|
|||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
|
|
||||||
import javax.naming.Context;
|
|
||||||
import javax.naming.NamingEnumeration;
|
|
||||||
import javax.naming.NamingException;
|
|
||||||
import javax.naming.directory.Attribute;
|
|
||||||
import javax.naming.directory.Attributes;
|
|
||||||
import javax.naming.directory.DirContext;
|
|
||||||
import javax.naming.directory.InitialDirContext;
|
|
||||||
import javax.net.ssl.SNIHostName;
|
import javax.net.ssl.SNIHostName;
|
||||||
import javax.net.ssl.SSLParameters;
|
import javax.net.ssl.SSLParameters;
|
||||||
import javax.net.ssl.SSLSocket;
|
import javax.net.ssl.SSLSocket;
|
||||||
@@ -22,20 +15,23 @@ import java.io.IOException;
|
|||||||
import java.io.InputStream;
|
import java.io.InputStream;
|
||||||
import java.io.OutputStream;
|
import java.io.OutputStream;
|
||||||
import java.net.IDN;
|
import java.net.IDN;
|
||||||
|
import java.net.DatagramPacket;
|
||||||
|
import java.net.DatagramSocket;
|
||||||
import java.net.Inet4Address;
|
import java.net.Inet4Address;
|
||||||
import java.net.Inet6Address;
|
import java.net.Inet6Address;
|
||||||
import java.net.InetAddress;
|
import java.net.InetAddress;
|
||||||
import java.net.InetSocketAddress;
|
import java.net.InetSocketAddress;
|
||||||
import java.net.Socket;
|
import java.net.Socket;
|
||||||
|
import java.net.SocketTimeoutException;
|
||||||
import java.net.URI;
|
import java.net.URI;
|
||||||
import java.net.URISyntaxException;
|
import java.net.URISyntaxException;
|
||||||
import java.nio.charset.StandardCharsets;
|
import java.nio.charset.StandardCharsets;
|
||||||
|
import java.nio.file.Files;
|
||||||
|
import java.nio.file.Path;
|
||||||
import java.security.MessageDigest;
|
import java.security.MessageDigest;
|
||||||
import java.security.NoSuchAlgorithmException;
|
import java.security.NoSuchAlgorithmException;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.Arrays;
|
|
||||||
import java.util.Hashtable;
|
|
||||||
import java.util.HashSet;
|
import java.util.HashSet;
|
||||||
import java.util.HexFormat;
|
import java.util.HexFormat;
|
||||||
import java.util.LinkedHashMap;
|
import java.util.LinkedHashMap;
|
||||||
@@ -51,7 +47,9 @@ import java.util.concurrent.RejectedExecutionException;
|
|||||||
import java.util.concurrent.ThreadPoolExecutor;
|
import java.util.concurrent.ThreadPoolExecutor;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
import java.util.concurrent.TimeoutException;
|
import java.util.concurrent.TimeoutException;
|
||||||
|
import java.util.concurrent.ThreadLocalRandom;
|
||||||
import java.util.concurrent.atomic.AtomicInteger;
|
import java.util.concurrent.atomic.AtomicInteger;
|
||||||
|
import java.util.function.IntSupplier;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class PersonalUrlFetchService {
|
public class PersonalUrlFetchService {
|
||||||
@@ -69,7 +67,6 @@ public class PersonalUrlFetchService {
|
|||||||
private static final int MAX_LINE_BYTES = 8 * 1024;
|
private static final int MAX_LINE_BYTES = 8 * 1024;
|
||||||
private static final long TOTAL_TIMEOUT_NANOS = 15_000_000_000L;
|
private static final long TOTAL_TIMEOUT_NANOS = 15_000_000_000L;
|
||||||
private static final long HARD_MAX_BODY_BYTES = 10L * 1024 * 1024;
|
private static final long HARD_MAX_BODY_BYTES = 10L * 1024 * 1024;
|
||||||
private static final ExecutorService DNS_EXECUTOR = boundedExecutor("personal-url-dns", 2, 8);
|
|
||||||
private static final ExecutorService WRITE_EXECUTOR = boundedExecutor("personal-url-write", 2, 8);
|
private static final ExecutorService WRITE_EXECUTOR = boundedExecutor("personal-url-write", 2, 8);
|
||||||
private static final ExecutorService HANDSHAKE_EXECUTOR = boundedExecutor("personal-url-tls", 2, 8);
|
private static final ExecutorService HANDSHAKE_EXECUTOR = boundedExecutor("personal-url-tls", 2, 8);
|
||||||
private static final String USER_AGENT = "wygj-personal-url-fetch/1.0";
|
private static final String USER_AGENT = "wygj-personal-url-fetch/1.0";
|
||||||
@@ -94,9 +91,8 @@ public class PersonalUrlFetchService {
|
|||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public PersonalUrlFetchService(PersonalKnowledgeProperties properties) {
|
public PersonalUrlFetchService(PersonalKnowledgeProperties properties) {
|
||||||
this(properties, new FallbackResolver(
|
this(properties, new UdpDnsResolver(configuredDnsServers(), PersonalUrlFetchService::exchangeDns,
|
||||||
new DeadlineDnsResolver(DNS_EXECUTOR, new JndiDnsQuery()),
|
() -> ThreadLocalRandom.current().nextInt(0x10000)), new RawSocketFetcher());
|
||||||
new DeadlineSystemResolver(DNS_EXECUTOR, InetAddress::getAllByName)), new RawSocketFetcher());
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private PersonalUrlFetchService(PersonalKnowledgeProperties properties, Resolver resolver, Fetcher fetcher) {
|
private PersonalUrlFetchService(PersonalKnowledgeProperties properties, Resolver resolver, Fetcher fetcher) {
|
||||||
@@ -362,154 +358,189 @@ public class PersonalUrlFetchService {
|
|||||||
List<InetAddress> resolve(String host, long deadlineNanos) throws IOException;
|
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
|
@FunctionalInterface
|
||||||
interface SystemAddressQuery {
|
interface DnsExchange {
|
||||||
InetAddress[] resolve(String host) throws IOException;
|
byte[] exchange(InetSocketAddress server, byte[] request, int timeoutMillis) throws IOException;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Bounds the JVM/system resolver with the same end-to-end deadline used by the fetch. */
|
/** Direct bounded UDP resolver. Closing the socket terminates every timed-out query without worker threads. */
|
||||||
static final class DeadlineSystemResolver implements Resolver {
|
static final class UdpDnsResolver implements Resolver {
|
||||||
private final ExecutorService executor;
|
private static final int TYPE_A = 1;
|
||||||
private final SystemAddressQuery query;
|
private static final int TYPE_AAAA = 28;
|
||||||
|
private final List<InetSocketAddress> servers;
|
||||||
|
private final DnsExchange exchange;
|
||||||
|
private final IntSupplier transactionIds;
|
||||||
|
|
||||||
DeadlineSystemResolver(ExecutorService executor, SystemAddressQuery query) {
|
UdpDnsResolver(List<InetSocketAddress> servers, DnsExchange exchange, IntSupplier transactionIds) {
|
||||||
this.executor = executor;
|
this.servers = List.copyOf(servers);
|
||||||
this.query = query;
|
this.exchange = exchange;
|
||||||
|
this.transactionIds = transactionIds;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public List<InetAddress> resolve(String host, long deadlineNanos) throws IOException {
|
public List<InetAddress> resolve(String host, long deadlineNanos) throws IOException {
|
||||||
long remaining = deadlineNanos - System.nanoTime();
|
if (servers.isEmpty()) throw new IOException("DNS resolver unavailable");
|
||||||
if (remaining <= 0) throw new IOException("resolution deadline exceeded");
|
IOException last = null;
|
||||||
Future<InetAddress[]> future;
|
for (int serverIndex = 0; serverIndex < servers.size(); serverIndex++) {
|
||||||
try {
|
List<InetAddress> addresses = new ArrayList<>();
|
||||||
future = executor.submit(() -> query.resolve(host));
|
boolean received = false;
|
||||||
} catch (RejectedExecutionException ex) {
|
for (int typeIndex = 0; typeIndex < 2; typeIndex++) {
|
||||||
throw new IOException("resolution unavailable");
|
int type = typeIndex == 0 ? TYPE_A : TYPE_AAAA;
|
||||||
}
|
int operationsLeft = (servers.size() - serverIndex) * 2 - typeIndex;
|
||||||
try {
|
int timeout = dnsTimeout(deadlineNanos, operationsLeft);
|
||||||
InetAddress[] addresses = future.get(remaining, TimeUnit.NANOSECONDS);
|
int transactionId = transactionIds.getAsInt() & 0xffff;
|
||||||
return addresses == null ? List.of() : List.copyOf(Arrays.asList(addresses));
|
byte[] request = dnsQuery(host, type, transactionId);
|
||||||
} catch (TimeoutException ex) {
|
try {
|
||||||
cancelAndPurge(executor, future);
|
byte[] response = exchange.exchange(servers.get(serverIndex), request, timeout);
|
||||||
throw new IOException("resolution deadline exceeded");
|
addresses.addAll(dnsAnswers(response, transactionId, type));
|
||||||
} catch (InterruptedException ex) {
|
received = true;
|
||||||
cancelAndPurge(executor, future);
|
} catch (SocketTimeoutException ex) {
|
||||||
Thread.currentThread().interrupt();
|
last = ex;
|
||||||
throw new IOException("resolution interrupted");
|
} catch (IOException ex) {
|
||||||
} catch (ExecutionException ex) {
|
last = ex;
|
||||||
cancelAndPurge(executor, future);
|
}
|
||||||
throw new IOException("resolution failed");
|
}
|
||||||
|
if (received) return addresses.stream().distinct().toList();
|
||||||
}
|
}
|
||||||
|
throw last == null ? new IOException("DNS resolution failed") : last;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@FunctionalInterface
|
private static List<InetSocketAddress> configuredDnsServers() {
|
||||||
interface DnsQuery {
|
String configured = System.getProperty("aihr.personal.dns-servers");
|
||||||
List<String> resolve(String host, int timeoutMillis, int retries) throws NamingException;
|
if (configured == null || configured.isBlank()) configured = System.getenv("AIHR_PERSONAL_DNS_SERVERS");
|
||||||
|
List<String> literals = new ArrayList<>();
|
||||||
|
if (configured != null && !configured.isBlank()) {
|
||||||
|
for (String value : configured.split("[,\\s]+")) if (!value.isBlank()) literals.add(value.trim());
|
||||||
|
} else {
|
||||||
|
try {
|
||||||
|
for (String line : Files.readAllLines(Path.of("/etc/resolv.conf"), StandardCharsets.US_ASCII)) {
|
||||||
|
String value = line.replaceFirst("#.*$", "").trim();
|
||||||
|
if (!value.startsWith("nameserver")) continue;
|
||||||
|
String[] parts = value.split("\\s+");
|
||||||
|
if (parts.length == 2) literals.add(parts[1]);
|
||||||
|
}
|
||||||
|
} catch (IOException ignored) {
|
||||||
|
return List.of();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
List<InetSocketAddress> servers = new ArrayList<>();
|
||||||
|
for (String literal : literals) {
|
||||||
|
if (servers.size() >= 4) break;
|
||||||
|
try {
|
||||||
|
servers.add(new InetSocketAddress(numericAddress(literal), 53));
|
||||||
|
} catch (IOException ignored) {
|
||||||
|
// Invalid configured resolver entries are not resolved as hostnames.
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return List.copyOf(servers);
|
||||||
}
|
}
|
||||||
|
|
||||||
static final class DeadlineDnsResolver implements Resolver {
|
private static byte[] exchangeDns(InetSocketAddress server, byte[] request, int timeoutMillis) throws IOException {
|
||||||
private final ExecutorService executor;
|
try (DatagramSocket socket = new DatagramSocket()) {
|
||||||
private final DnsQuery query;
|
socket.connect(server);
|
||||||
|
socket.setSoTimeout(timeoutMillis);
|
||||||
DeadlineDnsResolver(DnsQuery query) {
|
socket.send(new DatagramPacket(request, request.length));
|
||||||
this(DNS_EXECUTOR, query);
|
byte[] buffer = new byte[4096];
|
||||||
}
|
DatagramPacket response = new DatagramPacket(buffer, buffer.length);
|
||||||
|
socket.receive(response);
|
||||||
DeadlineDnsResolver(ExecutorService executor, DnsQuery query) {
|
validateDnsSource(server, response);
|
||||||
this.executor = executor;
|
return java.util.Arrays.copyOf(response.getData(), response.getLength());
|
||||||
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<List<String>> future;
|
|
||||||
try {
|
|
||||||
future = executor.submit(() -> {
|
|
||||||
long taskRemainingMillis = TimeUnit.NANOSECONDS.toMillis(deadlineNanos - System.nanoTime());
|
|
||||||
if (taskRemainingMillis <= 0) throw new NamingException("resolution deadline exceeded");
|
|
||||||
return query.resolve(host, (int) Math.min(5_000L, taskRemainingMillis), 0);
|
|
||||||
});
|
|
||||||
} catch (RejectedExecutionException ex) {
|
|
||||||
throw new IOException("resolution unavailable");
|
|
||||||
}
|
|
||||||
List<String> literals;
|
|
||||||
try {
|
|
||||||
literals = future.get(remaining, TimeUnit.NANOSECONDS);
|
|
||||||
} 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");
|
|
||||||
}
|
|
||||||
if (System.nanoTime() >= deadlineNanos) throw new IOException("resolution deadline exceeded");
|
|
||||||
List<InetAddress> addresses = new ArrayList<>();
|
|
||||||
if (literals != null) {
|
|
||||||
for (String literal : literals) addresses.add(numericAddress(literal));
|
|
||||||
}
|
|
||||||
return List.copyOf(addresses);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
static final class JndiDnsQuery implements DnsQuery {
|
static void validateDnsSource(InetSocketAddress server, DatagramPacket response) throws IOException {
|
||||||
@Override
|
if (!server.getAddress().equals(response.getAddress()) || server.getPort() != response.getPort()) {
|
||||||
public List<String> resolve(String host, int timeoutMillis, int retries) throws NamingException {
|
throw new IOException("DNS response source mismatch");
|
||||||
Hashtable<String, String> environment = environment(timeoutMillis);
|
}
|
||||||
environment.put("com.sun.jndi.dns.timeout.retries", Integer.toString(Math.max(0, retries)));
|
}
|
||||||
DirContext context = new InitialDirContext(environment);
|
|
||||||
try {
|
private static int dnsTimeout(long deadlineNanos, int operationsLeft) throws IOException {
|
||||||
Attributes attributes = context.getAttributes(host, new String[] {"A", "AAAA"});
|
long remaining = deadlineNanos - System.nanoTime();
|
||||||
List<String> values = new ArrayList<>();
|
if (remaining <= 0) throw new IOException("DNS resolution deadline exceeded");
|
||||||
collect(attributes.get("A"), values);
|
long millis = Math.max(1, TimeUnit.NANOSECONDS.toMillis(remaining) / Math.max(1, operationsLeft));
|
||||||
collect(attributes.get("AAAA"), values);
|
return (int) Math.min(2_000, millis);
|
||||||
return values;
|
}
|
||||||
} finally {
|
|
||||||
context.close();
|
private static byte[] dnsQuery(String host, int type, int transactionId) throws IOException {
|
||||||
|
ByteArrayOutputStream output = new ByteArrayOutputStream();
|
||||||
|
output.write((transactionId >>> 8) & 0xff);
|
||||||
|
output.write(transactionId & 0xff);
|
||||||
|
output.write(new byte[]{1, 0, 0, 1, 0, 0, 0, 0, 0, 0});
|
||||||
|
for (String label : host.split("\\.")) {
|
||||||
|
byte[] bytes = label.getBytes(StandardCharsets.US_ASCII);
|
||||||
|
if (bytes.length == 0 || bytes.length > 63) throw new IOException("invalid DNS name");
|
||||||
|
output.write(bytes.length);
|
||||||
|
output.write(bytes);
|
||||||
|
}
|
||||||
|
output.write(0);
|
||||||
|
output.write((type >>> 8) & 0xff);
|
||||||
|
output.write(type & 0xff);
|
||||||
|
output.write(new byte[]{0, 1});
|
||||||
|
return output.toByteArray();
|
||||||
|
}
|
||||||
|
|
||||||
|
private static List<InetAddress> dnsAnswers(byte[] response, int transactionId, int expectedType)
|
||||||
|
throws IOException {
|
||||||
|
if (response == null || response.length < 12 || response.length > 4096
|
||||||
|
|| unsigned16(response, 0) != transactionId) throw new IOException("invalid DNS response");
|
||||||
|
int flags = unsigned16(response, 2);
|
||||||
|
if ((flags & 0x8000) == 0 || (flags & 0x0200) != 0 || (flags & 0x000f) != 0
|
||||||
|
|| unsigned16(response, 4) != 1) throw new IOException("invalid DNS response");
|
||||||
|
int answerCount = unsigned16(response, 6);
|
||||||
|
int totalRecords = answerCount + unsigned16(response, 8) + unsigned16(response, 10);
|
||||||
|
if (answerCount > 64 || totalRecords > 128) throw new IOException("invalid DNS response");
|
||||||
|
int position = skipDnsName(response, 12);
|
||||||
|
requireDnsBytes(response, position, 4);
|
||||||
|
int questionType = unsigned16(response, position);
|
||||||
|
int questionClass = unsigned16(response, position + 2);
|
||||||
|
if (questionType != expectedType || questionClass != 1) throw new IOException("invalid DNS response");
|
||||||
|
position += 4;
|
||||||
|
List<InetAddress> addresses = new ArrayList<>();
|
||||||
|
for (int index = 0; index < answerCount; index++) {
|
||||||
|
position = skipDnsName(response, position);
|
||||||
|
requireDnsBytes(response, position, 10);
|
||||||
|
int type = unsigned16(response, position);
|
||||||
|
int recordClass = unsigned16(response, position + 2);
|
||||||
|
int length = unsigned16(response, position + 8);
|
||||||
|
position += 10;
|
||||||
|
requireDnsBytes(response, position, length);
|
||||||
|
if (recordClass == 1 && type == expectedType
|
||||||
|
&& ((type == UdpDnsResolver.TYPE_A && length == 4)
|
||||||
|
|| (type == UdpDnsResolver.TYPE_AAAA && length == 16))) {
|
||||||
|
addresses.add(InetAddress.getByAddress(java.util.Arrays.copyOfRange(response, position, position + length)));
|
||||||
}
|
}
|
||||||
|
position += length;
|
||||||
}
|
}
|
||||||
|
return List.copyOf(addresses);
|
||||||
|
}
|
||||||
|
|
||||||
static Hashtable<String, String> environment(int timeoutMillis) {
|
private static int skipDnsName(byte[] message, int position) throws IOException {
|
||||||
Hashtable<String, String> environment = new Hashtable<>();
|
for (int labels = 0; labels < 128; labels++) {
|
||||||
environment.put(Context.INITIAL_CONTEXT_FACTORY, "com.sun.jndi.dns.DnsContextFactory");
|
requireDnsBytes(message, position, 1);
|
||||||
environment.put("com.sun.jndi.dns.timeout.initial", Integer.toString(Math.max(1, timeoutMillis)));
|
int length = message[position] & 0xff;
|
||||||
environment.put("com.sun.jndi.dns.timeout.retries", "0");
|
if (length == 0) return position + 1;
|
||||||
return environment;
|
if ((length & 0xc0) == 0xc0) {
|
||||||
|
requireDnsBytes(message, position, 2);
|
||||||
|
int pointer = ((length & 0x3f) << 8) | (message[position + 1] & 0xff);
|
||||||
|
if (pointer >= message.length) throw new IOException("invalid DNS compression pointer");
|
||||||
|
return position + 2;
|
||||||
|
}
|
||||||
|
if ((length & 0xc0) != 0 || length > 63) throw new IOException("invalid DNS label");
|
||||||
|
position++;
|
||||||
|
requireDnsBytes(message, position, length);
|
||||||
|
position += length;
|
||||||
}
|
}
|
||||||
|
throw new IOException("DNS name too deep");
|
||||||
|
}
|
||||||
|
|
||||||
private static void collect(Attribute attribute, List<String> values) throws NamingException {
|
private static int unsigned16(byte[] value, int offset) throws IOException {
|
||||||
if (attribute == null) return;
|
requireDnsBytes(value, offset, 2);
|
||||||
NamingEnumeration<?> all = attribute.getAll();
|
return ((value[offset] & 0xff) << 8) | (value[offset + 1] & 0xff);
|
||||||
while (all.hasMore()) values.add(String.valueOf(all.next()).trim());
|
}
|
||||||
}
|
|
||||||
|
private static void requireDnsBytes(byte[] value, int offset, int length) throws IOException {
|
||||||
|
if (offset < 0 || length < 0 || offset > value.length - length) throw new IOException("truncated DNS response");
|
||||||
}
|
}
|
||||||
|
|
||||||
private static InetAddress numericAddress(String literal) throws IOException {
|
private static InetAddress numericAddress(String literal) throws IOException {
|
||||||
|
|||||||
+108
-97
@@ -14,9 +14,12 @@ import java.io.ByteArrayOutputStream;
|
|||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
import java.io.InputStream;
|
import java.io.InputStream;
|
||||||
import java.io.OutputStream;
|
import java.io.OutputStream;
|
||||||
|
import java.net.DatagramPacket;
|
||||||
import java.net.InetAddress;
|
import java.net.InetAddress;
|
||||||
|
import java.net.InetSocketAddress;
|
||||||
import java.net.ServerSocket;
|
import java.net.ServerSocket;
|
||||||
import java.net.Socket;
|
import java.net.Socket;
|
||||||
|
import java.net.SocketTimeoutException;
|
||||||
import java.net.URI;
|
import java.net.URI;
|
||||||
import java.nio.charset.StandardCharsets;
|
import java.nio.charset.StandardCharsets;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
@@ -223,112 +226,86 @@ class PersonalUrlFetchServiceTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void nativeDnsQueryReceivesRemainingTimeoutAndNumericAnswersOnly() {
|
void udpDnsMovesPastTwoSilentResolversWithoutWorkerPoolExhaustion() throws Exception {
|
||||||
|
List<InetSocketAddress> servers = List.of(
|
||||||
|
new InetSocketAddress("127.0.0.1", 5301),
|
||||||
|
new InetSocketAddress("127.0.0.1", 5302),
|
||||||
|
new InetSocketAddress("127.0.0.1", 5303));
|
||||||
|
AtomicInteger exchanges = new AtomicInteger();
|
||||||
|
var resolver = new PersonalUrlFetchService.UdpDnsResolver(servers, (server, request, timeoutMillis) -> {
|
||||||
|
exchanges.incrementAndGet();
|
||||||
|
if (server.getPort() != 5303) throw new SocketTimeoutException("silent resolver");
|
||||||
|
return dnsResponse(request, request[request.length - 3] == 1 ? PUBLIC : null);
|
||||||
|
}, () -> 0x1234);
|
||||||
|
|
||||||
|
List<InetAddress> result = resolver.resolve("example.com", System.nanoTime() + TimeUnit.SECONDS.toNanos(1));
|
||||||
|
|
||||||
|
assertEquals(List.of(PUBLIC), result);
|
||||||
|
assertEquals(6, exchanges.get());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void udpDnsFallbackAddressesStillUsePublicPolicyAndHonorDeadline() {
|
||||||
|
InetSocketAddress server = new InetSocketAddress("127.0.0.1", 5301);
|
||||||
|
var privateResolver = new PersonalUrlFetchService.UdpDnsResolver(List.of(server),
|
||||||
|
(ignored, request, timeoutMillis) -> dnsResponse(request, address("127.0.0.1")), () -> 7);
|
||||||
|
assertCode("PERSONAL_URL_BLOCKED", () -> fixture(privateResolver, request -> ok("text/plain", "ok"))
|
||||||
|
.validate("https://example.com/"), "private UDP answer");
|
||||||
|
|
||||||
AtomicInteger timeoutSeen = new AtomicInteger();
|
AtomicInteger timeoutSeen = new AtomicInteger();
|
||||||
AtomicInteger retriesSeen = new AtomicInteger(-1);
|
var silent = new PersonalUrlFetchService.UdpDnsResolver(List.of(server),
|
||||||
var resolver = new PersonalUrlFetchService.DeadlineDnsResolver((host, timeoutMillis, retries) -> {
|
(ignored, request, timeoutMillis) -> {
|
||||||
timeoutSeen.set(timeoutMillis); retriesSeen.set(retries);
|
timeoutSeen.set(timeoutMillis);
|
||||||
return List.of("93.184.216.34", "2606:2800:220:1:248:1893:25c8:1946");
|
throw new SocketTimeoutException("silent resolver");
|
||||||
});
|
}, () -> 8);
|
||||||
long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(200);
|
long started = System.nanoTime();
|
||||||
assertEquals(2, assertDoesNotThrow(() -> resolver.resolve("example.com", deadline)).size());
|
assertThrows(IOException.class, () -> silent.resolve("example.com",
|
||||||
assertTrue(timeoutSeen.get() > 0 && timeoutSeen.get() <= 200);
|
started + TimeUnit.MILLISECONDS.toNanos(40)));
|
||||||
assertEquals(0, retriesSeen.get());
|
assertTrue(timeoutSeen.get() > 0 && timeoutSeen.get() <= 40);
|
||||||
assertThrows(IOException.class, () -> resolver.resolve("example.com", System.nanoTime() - 1));
|
assertTrue(TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - started) < 500);
|
||||||
var nonNumeric = new PersonalUrlFetchService.DeadlineDnsResolver(
|
}
|
||||||
(host, timeoutMillis, retries) -> List.of("internal.example", "fe80::1%en0"));
|
|
||||||
assertThrows(IOException.class, () -> nonNumeric.resolve("example.com",
|
@Test
|
||||||
|
void udpDnsParsesAAndAaaaAndAcceptsValidEmptyAnswer() throws Exception {
|
||||||
|
InetAddress ipv6 = address("2606:4700:4700::1111");
|
||||||
|
InetSocketAddress server = new InetSocketAddress("127.0.0.1", 5301);
|
||||||
|
var resolver = new PersonalUrlFetchService.UdpDnsResolver(List.of(server),
|
||||||
|
(ignored, request, timeoutMillis) -> dnsResponse(request,
|
||||||
|
request[request.length - 3] == 1 ? PUBLIC : ipv6), () -> 0x2211);
|
||||||
|
assertEquals(List.of(PUBLIC, ipv6), resolver.resolve("example.com",
|
||||||
System.nanoTime() + TimeUnit.SECONDS.toNanos(1)));
|
System.nanoTime() + TimeUnit.SECONDS.toNanos(1)));
|
||||||
|
|
||||||
Hashtable<String, String> environment = PersonalUrlFetchService.JndiDnsQuery.environment(123);
|
var empty = new PersonalUrlFetchService.UdpDnsResolver(List.of(server),
|
||||||
assertEquals("123", environment.get("com.sun.jndi.dns.timeout.initial"));
|
(ignored, request, timeoutMillis) -> dnsResponse(request, null), () -> 0x2212);
|
||||||
assertEquals("0", environment.get("com.sun.jndi.dns.timeout.retries"));
|
assertEquals(List.of(), empty.resolve("example.com",
|
||||||
|
System.nanoTime() + TimeUnit.SECONDS.toNanos(1)));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void fallsBackToBoundedSystemDnsOnlyWhenJndiResolutionFails() throws Exception {
|
void udpDnsRejectsTransactionMismatchTruncationAndInvalidCompressionPointer() {
|
||||||
AtomicInteger fallbackCalls = new AtomicInteger();
|
InetSocketAddress server = new InetSocketAddress("127.0.0.1", 5301);
|
||||||
var resolver = new PersonalUrlFetchService.FallbackResolver(
|
assertMalformedDns(server, response -> response[1] ^= 1, "transaction mismatch");
|
||||||
(host, deadline) -> { throw new IOException("JNDI unavailable"); },
|
assertMalformedDns(server, response -> response[2] |= 0x02, "truncated response flag");
|
||||||
(host, deadline) -> { fallbackCalls.incrementAndGet(); return List.of(PUBLIC); });
|
assertMalformedDns(server, response -> {
|
||||||
var service = fixture(resolver, request -> ok("text/plain", "ok"));
|
int answerOffset = dnsQuestionEnd(response);
|
||||||
|
response[answerOffset] = (byte) 0xff;
|
||||||
assertEquals("https://example.com/", service.validate("https://example.com/").toString());
|
response[answerOffset + 1] = (byte) 0xff;
|
||||||
assertEquals(1, fallbackCalls.get());
|
}, "compression pointer out of bounds");
|
||||||
|
var emptyPacket = new PersonalUrlFetchService.UdpDnsResolver(List.of(server),
|
||||||
fallbackCalls.set(0);
|
(ignored, request, timeoutMillis) -> new byte[0], () -> 0x3311);
|
||||||
var emptyPrimary = new PersonalUrlFetchService.FallbackResolver(
|
assertThrows(IOException.class, () -> emptyPacket.resolve("example.com",
|
||||||
(host, deadline) -> List.of(),
|
System.nanoTime() + TimeUnit.SECONDS.toNanos(1)), "empty packet");
|
||||||
(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
|
@Test
|
||||||
void validatesEveryFallbackAddressAndFailsClosedWhenFallbackFails() {
|
void udpDnsRejectsUnexpectedResponseSource() throws Exception {
|
||||||
var privateFallback = new PersonalUrlFetchService.FallbackResolver(
|
InetSocketAddress expected = new InetSocketAddress(address("127.0.0.1"), 5301);
|
||||||
(host, deadline) -> { throw new IOException("JNDI unavailable"); },
|
DatagramPacket wrongAddress = new DatagramPacket(new byte[1], 1,
|
||||||
(host, deadline) -> List.of(PUBLIC, address("127.0.0.1")));
|
address("127.0.0.2"), 5301);
|
||||||
assertCode("PERSONAL_URL_BLOCKED", () -> fixture(privateFallback, request -> ok("text/plain", "ok"))
|
DatagramPacket wrongPort = new DatagramPacket(new byte[1], 1,
|
||||||
.validate("https://example.com/"), "mixed fallback addresses");
|
address("127.0.0.1"), 5302);
|
||||||
|
assertThrows(IOException.class, () -> PersonalUrlFetchService.validateDnsSource(expected, wrongAddress));
|
||||||
var failedFallback = new PersonalUrlFetchService.FallbackResolver(
|
assertThrows(IOException.class, () -> PersonalUrlFetchService.validateDnsSource(expected, wrongPort));
|
||||||
(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");
|
|
||||||
CountDownLatch entered = new CountDownLatch(1);
|
|
||||||
CountDownLatch release = new CountDownLatch(1);
|
|
||||||
try {
|
|
||||||
var resolver = new PersonalUrlFetchService.DeadlineDnsResolver(executor, (host, timeoutMillis, retries) -> {
|
|
||||||
entered.countDown();
|
|
||||||
boolean done = false;
|
|
||||||
while (!done) {
|
|
||||||
try { release.await(); done = true; }
|
|
||||||
catch (InterruptedException ignored) { }
|
|
||||||
}
|
|
||||||
return List.of("93.184.216.34");
|
|
||||||
});
|
|
||||||
long started = System.nanoTime();
|
|
||||||
assertThrows(IOException.class, () -> resolver.resolve("example.com", started + TimeUnit.MILLISECONDS.toNanos(40)));
|
|
||||||
assertTrue(entered.await(200, TimeUnit.MILLISECONDS));
|
|
||||||
assertTrue(TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - started) < 500);
|
|
||||||
} finally {
|
|
||||||
release.countDown();
|
|
||||||
executor.shutdownNow();
|
|
||||||
assertTrue(executor.awaitTermination(1, TimeUnit.SECONDS));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
@@ -634,6 +611,40 @@ class PersonalUrlFetchServiceTest {
|
|||||||
catch (Exception ex) { throw new AssertionError(ex); }
|
catch (Exception ex) { throw new AssertionError(ex); }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static byte[] dnsResponse(byte[] request, InetAddress answer) throws IOException {
|
||||||
|
ByteArrayOutputStream output = new ByteArrayOutputStream();
|
||||||
|
output.write(request, 0, 2);
|
||||||
|
output.write(new byte[]{(byte) 0x81, (byte) 0x80, 0, 1, 0, (byte) (answer == null ? 0 : 1), 0, 0, 0, 0});
|
||||||
|
output.write(request, 12, request.length - 12);
|
||||||
|
if (answer != null) {
|
||||||
|
byte[] address = answer.getAddress();
|
||||||
|
output.write(new byte[]{(byte) 0xc0, 0x0c});
|
||||||
|
output.write(request, request.length - 4, 2);
|
||||||
|
output.write(new byte[]{0, 1, 0, 0, 0, 30, 0, (byte) address.length});
|
||||||
|
output.write(address);
|
||||||
|
}
|
||||||
|
return output.toByteArray();
|
||||||
|
}
|
||||||
|
|
||||||
|
private static void assertMalformedDns(InetSocketAddress server,
|
||||||
|
java.util.function.Consumer<byte[]> mutation,
|
||||||
|
String context) {
|
||||||
|
var resolver = new PersonalUrlFetchService.UdpDnsResolver(List.of(server),
|
||||||
|
(ignored, request, timeoutMillis) -> {
|
||||||
|
byte[] response = dnsResponse(request, PUBLIC);
|
||||||
|
mutation.accept(response);
|
||||||
|
return response;
|
||||||
|
}, () -> 0x3311);
|
||||||
|
assertThrows(IOException.class, () -> resolver.resolve("example.com",
|
||||||
|
System.nanoTime() + TimeUnit.SECONDS.toNanos(1)), context);
|
||||||
|
}
|
||||||
|
|
||||||
|
private static int dnsQuestionEnd(byte[] response) {
|
||||||
|
int position = 12;
|
||||||
|
while ((response[position] & 0xff) != 0) position += 1 + (response[position] & 0xff);
|
||||||
|
return position + 5;
|
||||||
|
}
|
||||||
|
|
||||||
private static ByteArrayInputStream stream(String value) {
|
private static ByteArrayInputStream stream(String value) {
|
||||||
return new ByteArrayInputStream(value.getBytes(StandardCharsets.US_ASCII));
|
return new ByteArrayInputStream(value.getBytes(StandardCharsets.US_ASCII));
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -52,6 +52,6 @@
|
|||||||
- 大模型不作为演示硬依赖:模型管理已启用 chat 模型时,三角色对练为真实 LLM 生成与评分(asr/tts 配置后语音输入/播报可用);未配置或现场调用失败时全链路自动回退 seed,演示不中断。
|
- 大模型不作为演示硬依赖:模型管理已启用 chat 模型时,三角色对练为真实 LLM 生成与评分(asr/tts 配置后语音输入/播报可用);未配置或现场调用失败时全链路自动回退 seed,演示不中断。
|
||||||
# 个人 AI 助理 P0 验收
|
# 个人 AI 助理 P0 验收
|
||||||
|
|
||||||
先运行 `./scripts/personal-assistant-smoke.sh`,必须输出 `PASS`。脚本会真实创建并解析 TEXT、`cupsfilter` PDF 与公开网页(默认 `https://example.com/`),验证三类资料 READY、按 itemIds 可检索、回答引用来自实际命中资料,并在删除后确认 MySQL/MinIO/Qdrant 零残留;公网不可达必须失败,不允许改用 localhost 绕过 SSRF。浏览器使用手机号 A 登录后,依次收藏文字、PDF 与公开网页,等待资料状态变为 READY;按采集日期检索,并分别验证个人、企业与 mixed 问答的引用域。删除个人资料后,详情和搜索应立即不可见。
|
先运行 `./scripts/personal-assistant-smoke.sh`,必须输出 `PASS`。脚本会真实创建并解析 TEXT、`cupsfilter` PDF 与公开网页(默认 `https://example.com/`),验证三类资料 READY、按 itemIds 可检索、回答引用来自实际命中资料,并在删除后确认 MySQL/MinIO/Qdrant 及临时用户全部个人会话零残留;公开网页的 DNS、重定向逐跳校验与最终 READY 状态以后端为唯一判定,不做客户端 `curl --location` 预检,公网不可达必须失败,不允许改用 localhost 绕过 SSRF。浏览器使用手机号 A 登录后,依次收藏文字、PDF 与公开网页,等待资料状态变为 READY;按采集日期检索,并分别验证个人、企业与 mixed 问答的引用域。删除个人资料后,详情和搜索应立即不可见。
|
||||||
|
|
||||||
再使用手机号 B 登录,确认看不到 A 的资料标题、会话与引用,且不能访问 A 的详情、下载、重试或删除接口。私网 URL 与云元数据 URL 必须显示明确的 `PERSONAL_URL_BLOCKED`,回答不得出现无引用内容。企业知识未配置明确授权 allowlist 时,ENTERPRISE/mixed 必须 fail-closed。
|
再使用手机号 B 登录,确认看不到 A 的资料标题、会话与引用,且不能访问 A 的详情、下载、重试或删除接口。私网 URL 与云元数据 URL 必须显示明确的 `PERSONAL_URL_BLOCKED`,回答不得出现无引用内容。企业知识未配置明确授权 allowlist 时,ENTERPRISE/mixed 必须 fail-closed。
|
||||||
|
|||||||
+3
-1
@@ -209,4 +209,6 @@ curl -k -s https://peilian.njzhmj.top/h5/ | sed -n '1,20p'
|
|||||||
./scripts/personal-assistant-smoke.sh
|
./scripts/personal-assistant-smoke.sh
|
||||||
```
|
```
|
||||||
|
|
||||||
脚本每次生成唯一 smoke 手机号与 run marker,通过开发短信登录创建 A/B,并真实采集 TEXT、由 macOS `cupsfilter` 生成的可检索 PDF、公开网页 `https://example.com/`。它会验证三类资料 READY/检索/引用、owner 隔离、私有 OSS 匿名 403、SSRF、幂等删除,以及 MySQL/MinIO/Qdrant 零残留;退出时只按本次 user/item/session/OSS/job ID 与 run marker 回查清理。公开网页可用 `AIHR_PERSONAL_SMOKE_PUBLIC_URL` 覆盖,页面检索词可用 `AIHR_PERSONAL_SMOKE_PUBLIC_QUERY` 覆盖;公网不可达会明确失败,不会退回 localhost 或假数据。脚本不会输出 token。`./scripts/personal-assistant-smoke.sh --signal-self-test` 可单独验证 INT/TERM 分别返回 130/143。
|
脚本每次生成唯一 smoke 手机号与 run marker,通过开发短信登录创建 A/B,并真实采集 TEXT、由 macOS `cupsfilter` 生成的可检索 PDF、公开网页 `https://example.com/`。它会验证三类资料 READY/检索/引用、owner 隔离、私有 OSS 匿名 403、SSRF、幂等删除,以及 MySQL/MinIO/Qdrant 零残留;退出时按本次临时用户清理其全部个人会话/消息,并按 user/item/OSS/job ID 与 run marker 回查清理。公开网页可用 `AIHR_PERSONAL_SMOKE_PUBLIC_URL` 覆盖,页面检索词可用 `AIHR_PERSONAL_SMOKE_PUBLIC_QUERY` 覆盖;客户端只检查 URL 语法,不预先跟随重定向,DNS、逐跳 SSRF 校验和最终 READY 状态以后端为准,公网不可达会明确失败。需要为重定向目标做精确断言时可设置 `AIHR_PERSONAL_SMOKE_EXPECTED_PUBLIC_URL`。脚本不会输出 token。`./scripts/personal-assistant-smoke.sh --signal-self-test` 可单独验证 INT/TERM 分别返回 130/143。
|
||||||
|
|
||||||
|
个人网页采集默认从 `/etc/resolv.conf` 读取最多 4 个 DNS resolver,并使用有 socket deadline 的原生 UDP 查询;如运行环境的 resolver 配置不可用,可通过 `AIHR_PERSONAL_DNS_SERVERS=223.5.5.5,1.1.1.1` 显式覆盖。配置项只接受数字 IP,不会递归解析 DNS 服务器名称。
|
||||||
|
|||||||
@@ -20,6 +20,8 @@ PDF_TITLE="$TITLE-pdf"
|
|||||||
URL_TITLE="$TITLE-url"
|
URL_TITLE="$TITLE-url"
|
||||||
PDF_QUERY="Personal PDF verification evidence"
|
PDF_QUERY="Personal PDF verification evidence"
|
||||||
PUBLIC_URL="${AIHR_PERSONAL_SMOKE_PUBLIC_URL:-https://example.com/}"
|
PUBLIC_URL="${AIHR_PERSONAL_SMOKE_PUBLIC_URL:-https://example.com/}"
|
||||||
|
PUBLIC_EXPECTED_URL="${AIHR_PERSONAL_SMOKE_EXPECTED_PUBLIC_URL:-}"
|
||||||
|
[[ -n "${AIHR_PERSONAL_SMOKE_PUBLIC_URL:-}" ]] || PUBLIC_EXPECTED_URL="https://example.com/"
|
||||||
PUBLIC_QUERY="${AIHR_PERSONAL_SMOKE_PUBLIC_QUERY:-Example Domain}"
|
PUBLIC_QUERY="${AIHR_PERSONAL_SMOKE_PUBLIC_QUERY:-Example Domain}"
|
||||||
TMP_ROOT="$(mktemp -d "${TMPDIR:-/tmp}/wygj-personal-smoke.XXXXXX")"
|
TMP_ROOT="$(mktemp -d "${TMPDIR:-/tmp}/wygj-personal-smoke.XXXXXX")"
|
||||||
chmod 700 "$TMP_ROOT"
|
chmod 700 "$TMP_ROOT"
|
||||||
@@ -96,9 +98,6 @@ cleanup_once() {
|
|||||||
USER_B="$(mysql "select user_id from sys_user where phonenumber='$PHONE_B' and remark='移动端短信自动注册' order by user_id desc limit 1" | head -1)"
|
USER_B="$(mysql "select user_id from sys_user where phonenumber='$PHONE_B' and remark='移动端短信自动注册' order by user_id desc limit 1" | head -1)"
|
||||||
fi
|
fi
|
||||||
discover_run_items
|
discover_run_items
|
||||||
if [[ "$USER_A" =~ ^[0-9]+$ && ! "$SESSION_ID" =~ ^[0-9]+$ ]]; then
|
|
||||||
SESSION_ID="$(mysql "select s.id from aihr_personal_chat_session s join aihr_personal_chat_message m on m.session_id=s.id and m.owner_user_id=s.owner_user_id where s.tenant_id='000000' and s.owner_user_id=$USER_A and m.content like '%$RUN_ID%' order by s.id desc limit 1" | head -1)"
|
|
||||||
fi
|
|
||||||
local index item_id oss_id object_key
|
local index item_id oss_id object_key
|
||||||
for index in "${!ITEM_IDS[@]}"; do
|
for index in "${!ITEM_IDS[@]}"; do
|
||||||
item_id="${ITEM_IDS[$index]}"
|
item_id="${ITEM_IDS[$index]}"
|
||||||
@@ -119,10 +118,11 @@ cleanup_once() {
|
|||||||
delete from aihr_personal_cleanup_job where tenant_id='000000' and owner_user_id=$USER_A and item_id=$item_id;
|
delete from aihr_personal_cleanup_job where tenant_id='000000' and owner_user_id=$USER_A and item_id=$item_id;
|
||||||
delete from aihr_personal_item where tenant_id='000000' and owner_user_id=$USER_A and id=$item_id and title like '$TITLE-%';" >/dev/null 2>&1 || true
|
delete from aihr_personal_item where tenant_id='000000' and owner_user_id=$USER_A and id=$item_id and title like '$TITLE-%';" >/dev/null 2>&1 || true
|
||||||
done
|
done
|
||||||
if [[ "$SESSION_ID" =~ ^[0-9]+$ && "$USER_A" =~ ^[0-9]+$ ]]; then
|
for owner in "$USER_A" "$USER_B"; do
|
||||||
mysql "delete from aihr_personal_chat_message where tenant_id='000000' and owner_user_id=$USER_A and session_id=$SESSION_ID;
|
[[ "$owner" =~ ^[0-9]+$ ]] || continue
|
||||||
delete from aihr_personal_chat_session where tenant_id='000000' and owner_user_id=$USER_A and id=$SESSION_ID;" >/dev/null 2>&1 || true
|
mysql "delete from aihr_personal_chat_message where tenant_id='000000' and owner_user_id=$owner;
|
||||||
fi
|
delete from aihr_personal_chat_session where tenant_id='000000' and owner_user_id=$owner;" >/dev/null 2>&1 || true
|
||||||
|
done
|
||||||
for owner in "$USER_A" "$USER_B"; do
|
for owner in "$USER_A" "$USER_B"; do
|
||||||
[[ "$owner" =~ ^[0-9]+$ ]] || continue
|
[[ "$owner" =~ ^[0-9]+$ ]] || continue
|
||||||
mysql "delete from aihr_personal_space where tenant_id='000000' and owner_user_id=$owner and not exists
|
mysql "delete from aihr_personal_space where tenant_id='000000' and owner_user_id=$owner and not exists
|
||||||
@@ -145,7 +145,37 @@ cleanup_once() {
|
|||||||
done < <(redis-cli -h 127.0.0.1 -p 16379 -a ruoyi123 --scan --pattern "*resource/sms/code:$phone*" 2>/dev/null)
|
done < <(redis-cli -h 127.0.0.1 -p 16379 -a ruoyi123 --scan --pattern "*resource/sms/code:$phone*" 2>/dev/null)
|
||||||
done
|
done
|
||||||
}
|
}
|
||||||
on_exit() { local code=$?; trap - EXIT INT TERM; cleanup_once; exit "$code"; }
|
|
||||||
|
assert_identity_cleanup() {
|
||||||
|
local owner session_count message_count
|
||||||
|
for owner in "$USER_A" "$USER_B"; do
|
||||||
|
[[ "$owner" =~ ^[0-9]+$ ]] || continue
|
||||||
|
session_count="$(mysql "select count(*) from aihr_personal_chat_session where tenant_id='000000' and owner_user_id=$owner")"
|
||||||
|
message_count="$(mysql "select count(*) from aihr_personal_chat_message where tenant_id='000000' and owner_user_id=$owner")"
|
||||||
|
[[ "$session_count" == 0 && "$message_count" == 0 ]] || {
|
||||||
|
echo "FAIL: personal chat residue owner=$owner sessions=$session_count messages=$message_count" >&2
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
done
|
||||||
|
[[ "$(mysql "select count(*) from sys_user where phonenumber in ('$PHONE_A','$PHONE_B')")" == 0 ]] || {
|
||||||
|
echo "FAIL: smoke users were not removed" >&2
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
on_exit() {
|
||||||
|
local code=$?
|
||||||
|
trap - EXIT INT TERM
|
||||||
|
cleanup_once
|
||||||
|
if [[ "$code" != 0 && "${AIHR_SMOKE_CLEANUP_DRY_RUN:-0}" != 1 ]]; then
|
||||||
|
if assert_identity_cleanup; then
|
||||||
|
echo "PASS: failure-window owner session and user cleanup"
|
||||||
|
else
|
||||||
|
code=1
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
exit "$code"
|
||||||
|
}
|
||||||
on_int() { trap - EXIT INT TERM; cleanup_once; exit 130; }
|
on_int() { trap - EXIT INT TERM; cleanup_once; exit 130; }
|
||||||
on_term() { trap - EXIT INT TERM; cleanup_once; exit 143; }
|
on_term() { trap - EXIT INT TERM; cleanup_once; exit 143; }
|
||||||
trap on_exit EXIT
|
trap on_exit EXIT
|
||||||
@@ -253,6 +283,8 @@ assert_single_item_ask() {
|
|||||||
expect_success "A $label single-item answer"
|
expect_success "A $label single-item answer"
|
||||||
response_session="$(jq -er '.data.sessionId | tostring' <<<"$HTTP_BODY")"
|
response_session="$(jq -er '.data.sessionId | tostring' <<<"$HTTP_BODY")"
|
||||||
[[ "$response_session" =~ ^[0-9]+$ ]] || fail "A $label answer did not persist session"
|
[[ "$response_session" =~ ^[0-9]+$ ]] || fail "A $label answer did not persist session"
|
||||||
|
[[ "${AIHR_SMOKE_FAIL_AFTER_SESSION_CREATED:-0}" != 1 ]] \
|
||||||
|
|| fail "injected failure after backend session creation"
|
||||||
if [[ "$SESSION_ID" =~ ^[0-9]+$ ]]; then
|
if [[ "$SESSION_ID" =~ ^[0-9]+$ ]]; then
|
||||||
expect_code "$response_session" "$SESSION_ID" "A $label answer session continuity"
|
expect_code "$response_session" "$SESSION_ID" "A $label answer session continuity"
|
||||||
else
|
else
|
||||||
@@ -326,14 +358,12 @@ case "$PUBLIC_URL" in
|
|||||||
http://localhost*|https://localhost*|http://127.*|https://127.*|http://\[*|https://\[*|http://169.254.*|https://169.254.*)
|
http://localhost*|https://localhost*|http://127.*|https://127.*|http://\[*|https://\[*|http://169.254.*|https://169.254.*)
|
||||||
fail "AIHR_PERSONAL_SMOKE_PUBLIC_URL must be a public URL, not localhost/private metadata"
|
fail "AIHR_PERSONAL_SMOKE_PUBLIC_URL must be a public URL, not localhost/private metadata"
|
||||||
;;
|
;;
|
||||||
http://*|https://*) ;;
|
http://*|https://*)
|
||||||
|
[[ "$PUBLIC_URL" =~ ^https?://[^/?#]+([/?#].*)?$ ]] \
|
||||||
|
|| fail "AIHR_PERSONAL_SMOKE_PUBLIC_URL has no valid HTTP(S) authority"
|
||||||
|
;;
|
||||||
*) fail "AIHR_PERSONAL_SMOKE_PUBLIC_URL must use http or https" ;;
|
*) fail "AIHR_PERSONAL_SMOKE_PUBLIC_URL must use http or https" ;;
|
||||||
esac
|
esac
|
||||||
PUBLIC_EFFECTIVE_URL="$(curl --fail --location --silent --show-error --max-time 20 \
|
|
||||||
--output /dev/null --write-out '%{url_effective}' "$PUBLIC_URL")" \
|
|
||||||
|| fail "public URL unreachable: $PUBLIC_URL"
|
|
||||||
[[ "$PUBLIC_EFFECTIVE_URL" == http://* || "$PUBLIC_EFFECTIVE_URL" == https://* ]] \
|
|
||||||
|| fail "public URL did not resolve to HTTP(S): $PUBLIC_URL"
|
|
||||||
request GET /auth/tenant/list
|
request GET /auth/tenant/list
|
||||||
expect_success "backend health"
|
expect_success "backend health"
|
||||||
docker ps --format '{{.Names}}' | grep -qx "$DB_CONTAINER" || fail "database container not running: $DB_CONTAINER"
|
docker ps --format '{{.Names}}' | grep -qx "$DB_CONTAINER" || fail "database container not running: $DB_CONTAINER"
|
||||||
@@ -375,7 +405,9 @@ append_item_metadata "$PDF_ITEM_ID" "$PDF_TITLE"
|
|||||||
|
|
||||||
request POST /api/aihr/personal-assistant/items/url "$TOKEN_A" "$CLIENT_A" \
|
request POST /api/aihr/personal-assistant/items/url "$TOKEN_A" "$CLIENT_A" \
|
||||||
"$(jq -cn --arg url "$PUBLIC_URL" --arg title "$URL_TITLE" '{url:$url,title:$title}')"
|
"$(jq -cn --arg url "$PUBLIC_URL" --arg title "$URL_TITLE" '{url:$url,title:$title}')"
|
||||||
expect_success "A create public URL item"
|
if [[ "$HTTP_STATUS" != 200 || "$(jq -r '.code // empty' <<<"$HTTP_BODY")" != 200 ]]; then
|
||||||
|
fail "public URL capture failed HTTP=$HTTP_STATUS code=$(jq -r '.code // empty' <<<"$HTTP_BODY") msg=$(jq -r '.msg // empty' <<<"$HTTP_BODY")"
|
||||||
|
fi
|
||||||
URL_ITEM_ID="$(jq -er '.data.itemId | tostring' <<<"$HTTP_BODY")"
|
URL_ITEM_ID="$(jq -er '.data.itemId | tostring' <<<"$HTTP_BODY")"
|
||||||
append_item_metadata "$URL_ITEM_ID" "$URL_TITLE"
|
append_item_metadata "$URL_ITEM_ID" "$URL_TITLE"
|
||||||
expect_code "${#ITEM_IDS[@]}" 3 "three personal items captured"
|
expect_code "${#ITEM_IDS[@]}" 3 "three personal items captured"
|
||||||
@@ -387,7 +419,9 @@ expect_code "$(jq -r '.data.sourceType' <<<"$HTTP_BODY")" FILE "A PDF source typ
|
|||||||
expect_code "$(jq -r '.data.mimeType' <<<"$HTTP_BODY")" application/pdf "A PDF mime type"
|
expect_code "$(jq -r '.data.mimeType' <<<"$HTTP_BODY")" application/pdf "A PDF mime type"
|
||||||
wait_ready "$URL_ITEM_ID" "A public URL item"
|
wait_ready "$URL_ITEM_ID" "A public URL item"
|
||||||
expect_code "$(jq -r '.data.sourceType' <<<"$HTTP_BODY")" URL "A URL source type"
|
expect_code "$(jq -r '.data.sourceType' <<<"$HTTP_BODY")" URL "A URL source type"
|
||||||
expect_code "$(jq -r '.data.originalUrl' <<<"$HTTP_BODY")" "$PUBLIC_EFFECTIVE_URL" "A URL originalUrl"
|
URL_ORIGINAL="$(jq -r '.data.originalUrl // empty' <<<"$HTTP_BODY")"
|
||||||
|
[[ "$URL_ORIGINAL" == http://* || "$URL_ORIGINAL" == https://* ]] || fail "A URL originalUrl is not HTTP(S)"
|
||||||
|
[[ -z "$PUBLIC_EXPECTED_URL" ]] || expect_code "$URL_ORIGINAL" "$PUBLIC_EXPECTED_URL" "A URL originalUrl"
|
||||||
|
|
||||||
for index in "${!OSS_URLS[@]}"; do
|
for index in "${!OSS_URLS[@]}"; do
|
||||||
assert_anonymous_private "${OSS_URLS[$index]}" "personal object ${ITEM_IDS[$index]}"
|
assert_anonymous_private "${OSS_URLS[$index]}" "personal object ${ITEM_IDS[$index]}"
|
||||||
@@ -434,6 +468,7 @@ expect_code "$cleanup_done" 3 "three cleanup jobs completed"
|
|||||||
|
|
||||||
assert_business_cleanup
|
assert_business_cleanup
|
||||||
cleanup_once
|
cleanup_once
|
||||||
|
assert_identity_cleanup
|
||||||
expect_code "$(mysql "select count(*) from aihr_personal_item where tenant_id='000000' and owner_user_id=$USER_A and title like '$TITLE-%'")" 0 "run items residual"
|
expect_code "$(mysql "select count(*) from aihr_personal_item where tenant_id='000000' and owner_user_id=$USER_A and title like '$TITLE-%'")" 0 "run items residual"
|
||||||
expect_code "$(mysql "select count(*) from aihr_personal_fragment where tenant_id='000000' and owner_user_id=$USER_A and item_id in ($TEXT_ITEM_ID,$PDF_ITEM_ID,$URL_ITEM_ID)")" 0 "run fragments residual"
|
expect_code "$(mysql "select count(*) from aihr_personal_fragment where tenant_id='000000' and owner_user_id=$USER_A and item_id in ($TEXT_ITEM_ID,$PDF_ITEM_ID,$URL_ITEM_ID)")" 0 "run fragments residual"
|
||||||
expect_code "$(mysql "select count(*) from aihr_personal_cleanup_job where tenant_id='000000' and owner_user_id=$USER_A and item_id in ($TEXT_ITEM_ID,$PDF_ITEM_ID,$URL_ITEM_ID)")" 0 "run cleanup jobs residual"
|
expect_code "$(mysql "select count(*) from aihr_personal_cleanup_job where tenant_id='000000' and owner_user_id=$USER_A and item_id in ($TEXT_ITEM_ID,$PDF_ITEM_ID,$URL_ITEM_ID)")" 0 "run cleanup jobs residual"
|
||||||
|
|||||||
Reference in New Issue
Block a user