From 3e6c8babbef77e7abb36d2beaa1c527c9da26790 Mon Sep 17 00:00:00 2001 From: damingishere-coder Date: Mon, 24 Aug 2026 18:58:37 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=EF=BC=9A=E5=A2=9E=E5=BC=BA?= =?UTF-8?q?=20AI=20Provider=20=E7=A8=B3=E5=AE=9A=E6=80=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 1 + ...77\347\224\250\346\214\207\345\215\227.md" | 1 + front/app/env-config/page.tsx | 15 + .../controller/AiConfigController.java | 11 +- .../controller/BossController.java | 7 + .../service/AiProviderException.java | 57 +++ .../application/service/AiService.java | 473 +++++++++++++----- .../ChromeJobAnalysisQueueService.java | 44 +- .../application/service/CodexCliService.java | 53 +- .../application/service/ConfigService.java | 28 +- .../service/JobAiAnalysisService.java | 354 ++++++++----- .../service/JobAnalysisTaskStore.java | 14 +- .../AiConfigControllerJobTaskTest.java | 24 + .../BossControllerAiKeywordTest.java | 62 +++ .../service/AiServiceRemoteHttpTest.java | 282 +++++++++++ .../ChromeJobAnalysisQueueServiceTest.java | 59 +++ .../service/CodexCliServiceTest.java | 35 ++ .../service/ConfigServiceTest.java | 32 ++ .../JobAiAnalysisServiceStatusTest.java | 164 +++++- .../service/JobAnalysisTaskStoreTest.java | 21 + tasks/2026-08-24-p1-2-provider-resilience.md | 70 +++ 21 files changed, 1546 insertions(+), 261 deletions(-) create mode 100644 src/main/java/com/getjobs/application/service/AiProviderException.java create mode 100644 src/test/java/com/getjobs/application/controller/BossControllerAiKeywordTest.java create mode 100644 src/test/java/com/getjobs/application/service/AiServiceRemoteHttpTest.java create mode 100644 tasks/2026-08-24-p1-2-provider-resilience.md diff --git a/.env.example b/.env.example index 1d57429..f86359f 100644 --- a/.env.example +++ b/.env.example @@ -67,6 +67,7 @@ CODEX_PATH=codex CODEX_HOME= CODEX_MODEL=gpt-5.6-sol CODEX_TIMEOUT_SECONDS=300 +AI_REQUEST_TIMEOUT_SECONDS=120 BASE_URL= API_KEY= MODEL= diff --git "a/doc/\344\275\277\347\224\250\346\214\207\345\215\227.md" "b/doc/\344\275\277\347\224\250\346\214\207\345\215\227.md" index 1bd8955..2064b6c 100644 --- "a/doc/\344\275\277\347\224\250\346\214\207\345\215\227.md" +++ "b/doc/\344\275\277\347\224\250\346\214\207\345\215\227.md" @@ -132,6 +132,7 @@ AI 配置主要用于 Boss 直聘岗位匹配和打招呼语生成。常见字 | `HOOK_URL` | 企业微信机器人 webhook,用于推送运行消息 | | `AI_PROVIDER` | 默认 `codex`,复用本机 Codex 登录;填写 `api` 时改用远程接口 | | `CODEX_MODEL` | Codex CLI 模型,默认 `gpt-5.6-sol` | +| `AI_REQUEST_TIMEOUT_SECONDS` | 一次远程 AI 调用的总超时(含限流重试或兼容切换),默认 120 秒 | | `BASE_URL` | 仅远程 API 模式使用的模型服务地址 | | `API_KEY` | 仅远程 API 模式使用的密钥 | | `MODEL` | 仅远程 API 模式使用的模型名称 | diff --git a/front/app/env-config/page.tsx b/front/app/env-config/page.tsx index 5b93528..17790c7 100644 --- a/front/app/env-config/page.tsx +++ b/front/app/env-config/page.tsx @@ -16,6 +16,7 @@ export default function EnvConfig() { codexPath: 'codex', codexModel: 'gpt-5.6-sol', codexTimeoutSeconds: '300', + apiTimeoutSeconds: '120', baseUrl: '', apiKey: '', model: '', @@ -54,6 +55,7 @@ export default function EnvConfig() { codexPath: result.data.CODEX_PATH || 'codex', codexModel: result.data.CODEX_MODEL || 'gpt-5.6-sol', codexTimeoutSeconds: result.data.CODEX_TIMEOUT_SECONDS || '300', + apiTimeoutSeconds: result.data.AI_REQUEST_TIMEOUT_SECONDS || '120', baseUrl: result.data.BASE_URL || '', apiKey: '', model: result.data.MODEL || '', @@ -90,6 +92,7 @@ export default function EnvConfig() { CODEX_PATH: envConfig.codexPath, CODEX_MODEL: envConfig.codexModel, CODEX_TIMEOUT_SECONDS: envConfig.codexTimeoutSeconds, + AI_REQUEST_TIMEOUT_SECONDS: envConfig.apiTimeoutSeconds, BASE_URL: envConfig.baseUrl, MODEL: envConfig.model, BOT_IS_SEND: String(envConfig.botIsSend ?? 0), @@ -312,6 +315,18 @@ export default function EnvConfig() { />

DeepSeek 推荐模型 deepseek-chat;也可以填写其他 OpenAI-compatible 模型名

+
+ + setEnvConfig({ ...envConfig, apiTimeoutSeconds: e.target.value })} + /> +

默认 120 秒;超时或网络中断不会自动重发,以避免重复计费。

+
diff --git a/src/main/java/com/getjobs/application/controller/AiConfigController.java b/src/main/java/com/getjobs/application/controller/AiConfigController.java index 8c884ef..4b4743e 100644 --- a/src/main/java/com/getjobs/application/controller/AiConfigController.java +++ b/src/main/java/com/getjobs/application/controller/AiConfigController.java @@ -250,10 +250,13 @@ public ResponseEntity> savePriorityCompanies( public ResponseEntity> analyzeJob(@RequestBody JobAiAnalysisService.JobAnalysisRequest request) { Map response = new HashMap<>(); try { - response.put("success", true); - response.put("data", jobAiAnalysisService.analyzeJob(request)); - response.put("message", "AI岗位分析完成"); - return ResponseEntity.ok(response); + JobAiAnalysisService.AnalysisResult result = jobAiAnalysisService.analyzeJob(request); + response.put("success", !result.isFailure()); + response.put("data", result); + response.put("message", result.isFailure() ? result.getSummary() : "AI岗位分析完成"); + return result.isFailure() + ? ResponseEntity.status(502).body(response) + : ResponseEntity.ok(response); } catch (Exception e) { log.error("AI岗位分析失败", e); response.put("success", false); diff --git a/src/main/java/com/getjobs/application/controller/BossController.java b/src/main/java/com/getjobs/application/controller/BossController.java index add1be8..6b46558 100644 --- a/src/main/java/com/getjobs/application/controller/BossController.java +++ b/src/main/java/com/getjobs/application/controller/BossController.java @@ -379,6 +379,13 @@ public ResponseEntity> generateBossAiKeywords(@RequestBody(r "message", e.getMessage(), "keywords", List.of() )); + } catch (RuntimeException e) { + log.warn("生成 Boss AI 关键词失败: {}", e.getMessage()); + return ResponseEntity.status(502).body(Map.of( + "success", false, + "message", e.getMessage() == null ? "AI 关键词生成失败" : e.getMessage(), + "keywords", List.of() + )); } } diff --git a/src/main/java/com/getjobs/application/service/AiProviderException.java b/src/main/java/com/getjobs/application/service/AiProviderException.java new file mode 100644 index 0000000..d96c496 --- /dev/null +++ b/src/main/java/com/getjobs/application/service/AiProviderException.java @@ -0,0 +1,57 @@ +package com.getjobs.application.service; + +/** + * 不携带 Provider 原始响应的安全异常;可直接用于任务错误信息和本机 API 提示。 + */ +public class AiProviderException extends RuntimeException { + private final Code code; + private final Integer httpStatus; + private final String clientRequestId; + private final String providerRequestId; + private final boolean outcomeUnknown; + + public AiProviderException(Code code, + String message, + Integer httpStatus, + String clientRequestId, + String providerRequestId, + boolean outcomeUnknown, + Throwable cause) { + super(message, cause); + this.code = code; + this.httpStatus = httpStatus; + this.clientRequestId = clientRequestId; + this.providerRequestId = providerRequestId; + this.outcomeUnknown = outcomeUnknown; + } + + public Code getCode() { + return code; + } + + public Integer getHttpStatus() { + return httpStatus; + } + + public String getClientRequestId() { + return clientRequestId; + } + + public String getProviderRequestId() { + return providerRequestId; + } + + public boolean isOutcomeUnknown() { + return outcomeUnknown; + } + + public enum Code { + RATE_LIMITED, + HTTP_4XX, + HTTP_5XX, + TIMEOUT, + NETWORK, + EMPTY_RESPONSE, + INVALID_RESPONSE + } +} diff --git a/src/main/java/com/getjobs/application/service/AiService.java b/src/main/java/com/getjobs/application/service/AiService.java index cf16ce6..bb206cd 100644 --- a/src/main/java/com/getjobs/application/service/AiService.java +++ b/src/main/java/com/getjobs/application/service/AiService.java @@ -11,18 +11,26 @@ import org.springframework.transaction.annotation.Transactional; import org.springframework.context.annotation.DependsOn; +import java.io.IOException; import java.net.URI; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; +import java.net.http.HttpTimeoutException; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; import java.util.Base64; +import java.util.HexFormat; import java.util.LinkedHashMap; import java.util.Map; +import java.util.UUID; import java.time.Duration; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneId; +import java.time.ZonedDateTime; import java.time.format.DateTimeFormatter; +import java.time.format.DateTimeParseException; /** * AI 服务(Spring 管理) @@ -33,6 +41,10 @@ @RequiredArgsConstructor @DependsOn("profileService") public class AiService { + private static final int DEFAULT_API_TIMEOUT_SECONDS = 120; + private static final int MAX_REMOTE_REQUESTS = 2; + private static final long DEFAULT_RATE_LIMIT_DELAY_MILLIS = 500; + private static final long MAX_RATE_LIMIT_DELAY_MILLIS = 10_000; private final ConfigService configService; private final AiMapper aiMapper; private final ProfileService profileService; @@ -48,7 +60,6 @@ public class AiService { * @return AI 回复文本 */ public String sendRequest(String content) { - // 读取并校验配置 var cfg = configService.getAiConfigs(); if ("codex".equalsIgnoreCase(cfg.get("AI_PROVIDER"))) { return codexCliService.generateText(content, cfg); @@ -60,12 +71,11 @@ public String sendRequest(String content) { String endpoint = isResponsesModel(model) ? buildResponsesEndpoint(baseUrl) : buildChatCompletionsEndpoint(baseUrl); - - int timeoutInSeconds = 60; - - HttpClient client = HttpClient.newBuilder() - .connectTimeout(Duration.ofSeconds(timeoutInSeconds)) - .build(); + Duration timeout = requestTimeout(cfg); + HttpClient client = buildHttpClient(timeout); + String clientRequestId = UUID.randomUUID().toString(); + RequestBudget budget = new RequestBudget(MAX_REMOTE_REQUESTS); + long deadlineNanos = deadlineAfter(timeout); // 构建 JSON 请求体 JSONObject requestData = new JSONObject(); @@ -88,58 +98,28 @@ public String sendRequest(String content) { requestData.put("messages", messages); } - HttpRequest request = HttpRequest.newBuilder() - .uri(URI.create(endpoint)) - .header("Content-Type", "application/json") - .header("Accept", "application/json") - .header("Authorization", "Bearer " + apiKey) - // 某些服务(例如 Azure OpenAI)需要 api-key 头,额外加一层兼容 - .header("api-key", apiKey) - .POST(HttpRequest.BodyPublishers.ofString(requestData.toString())) - .build(); - - try { - HttpResponse response = client.send(request, HttpResponse.BodyHandlers.ofString()); - if (response.statusCode() == 200) { - JSONObject responseObject = new JSONObject(response.body()); - - String requestId = responseObject.optString("id"); - long created = responseObject.optLong("created", 0); - String usedModel = responseObject.optString("model"); - - String responseContent = endpoint.endsWith("/responses") - ? extractResponsesContent(responseObject, response.body()) - : extractChatContent(responseObject, response.body()); - - JSONObject usageObject = responseObject.optJSONObject("usage"); - int promptTokens = usageObject != null ? usageObject.optInt("prompt_tokens", -1) : -1; - int completionTokens = usageObject != null ? usageObject.optInt("completion_tokens", -1) : -1; - int totalTokens = usageObject != null ? usageObject.optInt("total_tokens", -1) : -1; - - LocalDateTime createdTime = created > 0 - ? Instant.ofEpochSecond(created).atZone(ZoneId.systemDefault()).toLocalDateTime() - : LocalDateTime.now(); - DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); - - log.info("AI响应: id={}, time={}, model={}, promptTokens={}, completionTokens={}, totalTokens={}", - requestId, createdTime.format(formatter), usedModel, promptTokens, completionTokens, totalTokens); - - return responseContent; - } else { - // 更详细的错误日志,便于定位 400 问题 - log.error("AI请求失败: status={}, endpoint={}, body={}", response.statusCode(), endpoint, response.body()); - // 针对 Responses-only 模型误用 Chat Completions 的常见错误做一次自动重试 - if (!endpoint.endsWith("/responses") && containsReasoningParamError(response.body())) { - String fallbackEndpoint = buildResponsesEndpoint(baseUrl); - log.warn("检测到 reasoning 相关参数错误,自动切换到 Responses API 重试: {}", fallbackEndpoint); - return sendRequestViaResponses(content, apiKey, model, fallbackEndpoint); - } - throw new RuntimeException("AI请求失败,状态码: " + response.statusCode() + ", 详情: " + response.body()); - } - } catch (Exception e) { - log.error("调用AI服务异常", e); - throw e instanceof RuntimeException ? (RuntimeException) e : new RuntimeException(e); - } + ProviderHttpResponse response = sendRemoteJson( + client, endpoint, apiKey, requestData, deadlineNanos, clientRequestId, budget, true); + if (isReasoningCompatibilityStatus(response.statusCode()) + && !endpoint.endsWith("/responses") + && containsReasoningParamError(response.body()) + && budget.hasRemaining()) { + String fallbackEndpoint = buildResponsesEndpoint(baseUrl); + log.warn("AI endpoint 兼容切换: clientRequestId={}, fromHost={}, toHost={}", + clientRequestId, endpointHost(endpoint), endpointHost(fallbackEndpoint)); + JSONObject responsesData = new JSONObject(); + responsesData.put("model", model); + responsesData.put("temperature", 0.5); + responsesData.put("input", content); + response = sendRemoteJson( + client, fallbackEndpoint, apiKey, responsesData, deadlineNanos, + clientRequestId, budget, true); + endpoint = fallbackEndpoint; + } + if (response.statusCode() != 200) { + throw providerHttpError(response, endpoint, clientRequestId); + } + return parseTextResponse(response, endpoint, clientRequestId); } /** @@ -157,6 +137,9 @@ public String extractResumeFromImage(byte[] imageBytes, String mimeType) { String apiKey = cfg.get("API_KEY"); String model = cfg.get("MODEL"); String endpoint = buildChatCompletionsEndpoint(baseUrl); + Duration timeout = requestTimeout(cfg); + String clientRequestId = UUID.randomUUID().toString(); + long deadlineNanos = deadlineAfter(timeout); String dataUrl = "data:" + (mimeType == null || mimeType.isBlank() ? "image/jpeg" : mimeType) + ";base64," + Base64.getEncoder().encodeToString(imageBytes); @@ -181,31 +164,23 @@ public String extractResumeFromImage(byte[] imageBytes, String mimeType) { messages.put(message); requestData.put("messages", messages); - HttpRequest request = HttpRequest.newBuilder() - .uri(URI.create(endpoint)) - .header("Content-Type", "application/json") - .header("Accept", "application/json") - .header("Authorization", "Bearer " + apiKey) - .header("api-key", apiKey) - .POST(HttpRequest.BodyPublishers.ofString(requestData.toString())) - .build(); - - try { - HttpClient client = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(60)).build(); - HttpResponse response = client.send(request, HttpResponse.BodyHandlers.ofString()); - if (response.statusCode() != 200) { - log.error("图片简历解析失败: status={}, body={}", response.statusCode(), response.body()); - throw new RuntimeException("图片简历解析失败,状态码: " + response.statusCode()); - } - JSONObject responseObject = new JSONObject(response.body()); - JSONObject messageObject = responseObject.getJSONArray("choices") - .getJSONObject(0) - .getJSONObject("message"); - return messageObject.getString("content"); - } catch (Exception e) { - log.error("图片简历解析异常", e); - throw e instanceof RuntimeException ? (RuntimeException) e : new RuntimeException(e); + ProviderHttpResponse response = sendRemoteJson( + buildHttpClient(timeout), endpoint, apiKey, requestData, deadlineNanos, + clientRequestId, new RequestBudget(MAX_REMOTE_REQUESTS), true); + if (response.statusCode() != 200) { + throw providerHttpError(response, endpoint, clientRequestId); + } + String result = parseChatContent(response, endpoint, clientRequestId); + if (result.isBlank()) { + throw providerResponseError( + AiProviderException.Code.EMPTY_RESPONSE, + "AI Provider 返回空内容", + response, + endpoint, + clientRequestId + ); } + return result; } /** @@ -358,55 +333,290 @@ private boolean containsReasoningParamError(String body) { || s.contains("reasoning.summary"); } - /** - * 使用 Responses API 发送一次请求(用于自动降级/重试) - */ - private String sendRequestViaResponses(String content, String apiKey, String model, String endpoint) { - int timeoutInSeconds = 60; - HttpClient client = HttpClient.newBuilder() - .connectTimeout(Duration.ofSeconds(timeoutInSeconds)) - .build(); + private HttpClient buildHttpClient(Duration timeout) { + Duration connectTimeout = timeout.compareTo(Duration.ofSeconds(30)) < 0 + ? timeout + : Duration.ofSeconds(30); + return HttpClient.newBuilder().connectTimeout(connectTimeout).build(); + } - JSONObject requestData = new JSONObject(); - requestData.put("model", model); - requestData.put("temperature", 0.5); - requestData.put("input", content); + private Duration requestTimeout(Map config) { + String raw = config == null ? null : config.get("AI_REQUEST_TIMEOUT_SECONDS"); + try { + int seconds = raw == null || raw.isBlank() + ? DEFAULT_API_TIMEOUT_SECONDS + : Integer.parseInt(raw.trim()); + return Duration.ofSeconds(Math.max(1, Math.min(1800, seconds))); + } catch (NumberFormatException ignored) { + return Duration.ofSeconds(DEFAULT_API_TIMEOUT_SECONDS); + } + } - HttpRequest request = HttpRequest.newBuilder() - .uri(URI.create(endpoint)) - .header("Content-Type", "application/json") - .header("Accept", "application/json") - .header("Authorization", "Bearer " + apiKey) - .header("api-key", apiKey) - .POST(HttpRequest.BodyPublishers.ofString(requestData.toString())) - .build(); + private ProviderHttpResponse sendRemoteJson(HttpClient client, + String endpoint, + String apiKey, + JSONObject requestData, + long deadlineNanos, + String clientRequestId, + RequestBudget budget, + boolean allowRateLimitRetry) { + boolean rateLimitRetried = false; + while (budget.consume()) { + Duration remaining = remainingDuration(deadlineNanos, clientRequestId); + HttpRequest request = HttpRequest.newBuilder() + .uri(URI.create(endpoint)) + .timeout(remaining) + .header("Content-Type", "application/json") + .header("Accept", "application/json") + .header("Authorization", "Bearer " + apiKey) + .header("api-key", apiKey) + .header("X-Request-ID", clientRequestId) + .POST(HttpRequest.BodyPublishers.ofString(requestData.toString())) + .build(); + try { + HttpResponse response = client.send(request, HttpResponse.BodyHandlers.ofString()); + ProviderHttpResponse captured = new ProviderHttpResponse( + response.statusCode(), + response.body() == null ? "" : response.body(), + providerRequestId(response) + ); + if (captured.statusCode() == 429 + && allowRateLimitRetry + && !rateLimitRetried + && budget.hasRemaining()) { + long delayMillis = retryAfterMillis(response); + if (delayMillis >= 0 + && delayMillis <= MAX_RATE_LIMIT_DELAY_MILLIS + && delayMillis < remainingMillis(deadlineNanos)) { + rateLimitRetried = true; + log.warn("AI Provider 限流,执行唯一一次有界重试: clientRequestId={}, providerRequestId={}, " + + "endpointHost={}, retryAfterMs={}", + clientRequestId, captured.providerRequestId(), endpointHost(endpoint), delayMillis); + sleepBeforeRetry(delayMillis, clientRequestId); + continue; + } + } + return captured; + } catch (HttpTimeoutException e) { + throw new AiProviderException( + AiProviderException.Code.TIMEOUT, + "AI Provider 请求超时,请先确认任务状态再重试(requestId=" + clientRequestId + ")", + null, clientRequestId, "", true, e); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AiProviderException( + AiProviderException.Code.NETWORK, + "AI Provider 请求被中断,请先确认任务状态再重试(requestId=" + clientRequestId + ")", + null, clientRequestId, "", true, e); + } catch (IOException e) { + throw new AiProviderException( + AiProviderException.Code.NETWORK, + "AI Provider 网络异常,请先确认任务状态再重试(requestId=" + clientRequestId + ")", + null, clientRequestId, "", true, e); + } + } + throw new AiProviderException( + AiProviderException.Code.INVALID_RESPONSE, + "AI Provider 请求预算已耗尽(requestId=" + clientRequestId + ")", + null, clientRequestId, "", false, null); + } + private long retryAfterMillis(HttpResponse response) { + String raw = response.headers().firstValue("Retry-After").orElse("").trim(); + if (raw.isEmpty()) return DEFAULT_RATE_LIMIT_DELAY_MILLIS; try { - HttpResponse response = client.send(request, HttpResponse.BodyHandlers.ofString()); - if (response.statusCode() == 200) { - JSONObject resp = new JSONObject(response.body()); - return extractResponsesContent(resp, response.body()); + long seconds = Long.parseLong(raw); + if (seconds < 0) return -1; + return Math.min(Long.MAX_VALUE / 1000, seconds) * 1000; + } catch (NumberFormatException ignored) { + try { + long millis = Duration.between( + Instant.now(), + ZonedDateTime.parse(raw, DateTimeFormatter.RFC_1123_DATE_TIME).toInstant() + ).toMillis(); + return Math.max(0, millis); + } catch (DateTimeParseException invalidDate) { + return -1; } - log.error("Responses API 调用失败: status={}, endpoint={}, body={}", response.statusCode(), endpoint, response.body()); - throw new RuntimeException("AI请求失败,状态码: " + response.statusCode() + ", 详情: " + response.body()); - } catch (Exception e) { - log.error("Responses API 调用异常", e); - throw e instanceof RuntimeException ? (RuntimeException) e : new RuntimeException(e); } } - private String extractChatContent(JSONObject responseObject, String rawBody) { + private boolean isReasoningCompatibilityStatus(int statusCode) { + return statusCode == 400 || statusCode == 422; + } + + private long deadlineAfter(Duration timeout) { + return System.nanoTime() + timeout.toNanos(); + } + + private Duration remainingDuration(long deadlineNanos, String clientRequestId) { + long remainingMillis = remainingMillis(deadlineNanos); + if (remainingMillis <= 0) { + throw new AiProviderException( + AiProviderException.Code.TIMEOUT, + "AI Provider 总请求时间已耗尽,请先确认任务状态再重试(requestId=" + clientRequestId + ")", + null, clientRequestId, "", true, null); + } + return Duration.ofMillis(remainingMillis); + } + + private long remainingMillis(long deadlineNanos) { + long remainingNanos = deadlineNanos - System.nanoTime(); + if (remainingNanos <= 0) return 0; + return Math.max(1, Duration.ofNanos(remainingNanos).toMillis()); + } + + private void sleepBeforeRetry(long delayMillis, String clientRequestId) { + if (delayMillis <= 0) return; + try { + Thread.sleep(delayMillis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AiProviderException( + AiProviderException.Code.NETWORK, + "AI Provider 限流等待被中断(requestId=" + clientRequestId + ")", + 429, clientRequestId, "", false, e); + } + } + + private String parseTextResponse(ProviderHttpResponse response, + String endpoint, + String clientRequestId) { + JSONObject responseObject = parseResponseEnvelope(response, endpoint, clientRequestId); + String responseContent = endpoint.endsWith("/responses") + ? extractResponsesContent(responseObject) + : extractChatContent(responseObject); + if (responseContent == null || responseContent.isBlank()) { + throw providerResponseError( + AiProviderException.Code.EMPTY_RESPONSE, + "AI Provider 返回空内容", + response, endpoint, clientRequestId); + } + + String responseId = responseObject.optString("id", response.providerRequestId()); + long created = responseObject.optLong("created", 0); + String usedModel = responseObject.optString("model"); + JSONObject usageObject = responseObject.optJSONObject("usage"); + int promptTokens = usageObject != null ? usageObject.optInt("prompt_tokens", -1) : -1; + int completionTokens = usageObject != null ? usageObject.optInt("completion_tokens", -1) : -1; + int totalTokens = usageObject != null ? usageObject.optInt("total_tokens", -1) : -1; + LocalDateTime createdTime = created > 0 + ? Instant.ofEpochSecond(created).atZone(ZoneId.systemDefault()).toLocalDateTime() + : LocalDateTime.now(); + DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); + log.info("AI响应: clientRequestId={}, providerRequestId={}, time={}, model={}, " + + "promptTokens={}, completionTokens={}, totalTokens={}", + clientRequestId, responseId, createdTime.format(formatter), usedModel, + promptTokens, completionTokens, totalTokens); + return responseContent; + } + + private String parseChatContent(ProviderHttpResponse response, + String endpoint, + String clientRequestId) { + return extractChatContent(parseResponseEnvelope(response, endpoint, clientRequestId)); + } + + private JSONObject parseResponseEnvelope(ProviderHttpResponse response, + String endpoint, + String clientRequestId) { + if (response.body() == null || response.body().isBlank()) { + throw providerResponseError( + AiProviderException.Code.EMPTY_RESPONSE, + "AI Provider 返回空响应", + response, endpoint, clientRequestId); + } + try { + return new JSONObject(response.body()); + } catch (RuntimeException e) { + throw providerResponseError( + AiProviderException.Code.INVALID_RESPONSE, + "AI Provider 响应不是有效 JSON", + response, endpoint, clientRequestId, e); + } + } + + private AiProviderException providerHttpError(ProviderHttpResponse response, + String endpoint, + String clientRequestId) { + AiProviderException.Code code = response.statusCode() == 429 + ? AiProviderException.Code.RATE_LIMITED + : response.statusCode() >= 500 + ? AiProviderException.Code.HTTP_5XX + : AiProviderException.Code.HTTP_4XX; + String message = switch (code) { + case RATE_LIMITED -> "AI Provider 请求过于频繁,请稍后显式重试"; + case HTTP_5XX -> "AI Provider 服务异常,未自动重试"; + default -> "AI Provider 拒绝了请求"; + }; + return providerResponseError(code, message, response, endpoint, clientRequestId); + } + + private AiProviderException providerResponseError(AiProviderException.Code code, + String message, + ProviderHttpResponse response, + String endpoint, + String clientRequestId) { + return providerResponseError(code, message, response, endpoint, clientRequestId, null); + } + + private AiProviderException providerResponseError(AiProviderException.Code code, + String message, + ProviderHttpResponse response, + String endpoint, + String clientRequestId, + Throwable cause) { + String body = response.body() == null ? "" : response.body(); + log.error("AI Provider 调用失败: code={}, status={}, endpointHost={}, clientRequestId={}, " + + "providerRequestId={}, bodyLength={}, bodyHash={}", + code, response.statusCode(), endpointHost(endpoint), clientRequestId, + response.providerRequestId(), body.length(), shortHash(body)); + boolean outcomeUnknown = code == AiProviderException.Code.HTTP_5XX + || code == AiProviderException.Code.EMPTY_RESPONSE + || code == AiProviderException.Code.INVALID_RESPONSE; + return new AiProviderException( + code, + message + "(requestId=" + clientRequestId + ")", + response.statusCode(), clientRequestId, response.providerRequestId(), outcomeUnknown, cause); + } + + private String providerRequestId(HttpResponse response) { + return response.headers().firstValue("x-request-id") + .or(() -> response.headers().firstValue("request-id")) + .orElse(""); + } + + private String endpointHost(String endpoint) { + try { + URI uri = URI.create(endpoint); + return uri.getHost() == null ? "unknown" : uri.getHost(); + } catch (RuntimeException ignored) { + return "invalid"; + } + } + + private String shortHash(String value) { + try { + byte[] hash = MessageDigest.getInstance("SHA-256") + .digest((value == null ? "" : value).getBytes(StandardCharsets.UTF_8)); + return HexFormat.of().formatHex(hash).substring(0, 12); + } catch (Exception ignored) { + return "unavailable"; + } + } + + private String extractChatContent(JSONObject responseObject) { try { JSONObject messageObject = responseObject.getJSONArray("choices") .getJSONObject(0) .getJSONObject("message"); - return flattenContent(messageObject.opt("content"), rawBody); + return flattenContent(messageObject.opt("content")); } catch (Exception ignore) { - return rawBody; + return ""; } } - private String extractResponsesContent(JSONObject responseObject, String rawBody) { + private String extractResponsesContent(JSONObject responseObject) { String outputText = responseObject.optString("output_text", null); if (outputText != null && !outputText.isEmpty()) { return outputText; @@ -420,13 +630,13 @@ private String extractResponsesContent(JSONObject responseObject, String rawBody JSONObject messageObject = responseObject.getJSONArray("choices") .getJSONObject(0) .getJSONObject("message"); - String content = flattenContent(messageObject.opt("content"), rawBody); + String content = flattenContent(messageObject.opt("content")); if (content != null && !content.isBlank()) { return content; } } catch (Exception ignore) { } - return rawBody; + return ""; } private void collectResponseOutputText(Object value, StringBuilder out) { @@ -450,8 +660,8 @@ private void collectResponseOutputText(Object value, StringBuilder out) { } } - private String flattenContent(Object content, String fallback) { - if (content == null) return fallback; + private String flattenContent(Object content) { + if (content == null) return ""; if (content instanceof String text) return text; if (content instanceof JSONArray arr) { StringBuilder out = new StringBuilder(); @@ -471,11 +681,32 @@ private String flattenContent(Object content, String fallback) { } } } - return out.isEmpty() ? fallback : out.toString(); + return out.toString(); } return String.valueOf(content); } + private record ProviderHttpResponse(int statusCode, String body, String providerRequestId) { + } + + private static final class RequestBudget { + private int remaining; + + private RequestBudget(int remaining) { + this.remaining = Math.max(1, remaining); + } + + private boolean consume() { + if (remaining <= 0) return false; + remaining--; + return true; + } + + private boolean hasRemaining() { + return remaining > 0; + } + } + // ================= 合并的 AI 配置管理方法 ================= /** diff --git a/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java b/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java index 93a4cf2..8404a5d 100644 --- a/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java +++ b/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java @@ -113,6 +113,20 @@ public java.util.List listTasks(long profileId, i } public JobAnalysisTaskStore.RetryResult retry(long taskId, long profileId, boolean confirmUnknown) { + JobAnalysisTaskStore.TaskRecord current = taskStore.findByIdAndProfile(taskId, profileId); + if (current != null + && current.statusEnum() == JobAnalysisTaskStore.Status.UNKNOWN + && confirmUnknown) { + JobAiAnalysisService.JobAnalysisRequest request = taskStore.deserialize(current); + JobAiAnalysisService.PlatformAnalysisState platformState = + jobAiAnalysisService.inspectPlatformAnalysis(request); + if (DeliveryStatus.AI_ANALYZING.equals(platformState.status()) + && !jobAiAnalysisService.markAnalysisInterrupted( + request, "用户已确认 UNKNOWN 任务,重试前安全复位 AI 分析状态")) { + return JobAnalysisTaskStore.RetryResult.rejected( + "岗位仍处于 AI 分析中且安全复位失败,未重新调用 AI Provider;请稍后重试"); + } + } JobAnalysisTaskStore.RetryResult result = taskStore.retry(taskId, profileId, confirmUnknown); if (result.accepted() && result.task() != null) { schedule(result.task().id()); @@ -205,22 +219,32 @@ private void runPersistedTask(long taskId) { } boolean failed = result.isFailure(); String summary = Objects.toString(result.getSummary(), ""); - boolean completed = taskStore.complete(taskId, leaseToken, failed, summary); - if (!completed) { + boolean completed = result.isProviderOutcomeUnknown() + ? taskStore.completeUnknown(taskId, leaseToken, summary) + : taskStore.complete(taskId, leaseToken, failed, summary); + if (!completed && !result.isProviderOutcomeUnknown()) { reconcileLateWorkerResult(claimed, leaseToken, failed, summary); } + if (!completed && result.isProviderOutcomeUnknown()) { + log.warn("AI 任务 {} 的 UNKNOWN 终态写入未命中当前租约,将等待过期对账", taskId); + } - if (failed) { + if (result.isProviderOutcomeUnknown()) { + emit(progress, JobProgressMessage.warning( + platform, "AI分析结果未知:" + jobName + "," + summary)); + } else if (failed) { emit(progress, JobProgressMessage.warning(platform, "AI分析失败:" + jobName + "," + summary)); } + String completionLabel = "跳过:"; + if (result.shouldApply()) { + completionLabel = DeliveryStatus.WAITING_CONFIRM + ":"; + } else if (result.isProviderOutcomeUnknown()) { + completionLabel = "结果未知:"; + } else if (failed) { + completionLabel = DeliveryStatus.AI_ANALYSIS_FAILED + ":"; + } emit(progress, JobProgressMessage.progress( - platform, - (result.shouldApply() - ? DeliveryStatus.WAITING_CONFIRM + ":" - : failed ? DeliveryStatus.AI_ANALYSIS_FAILED + ":" : "跳过:") + jobName, - current, - total - )); + platform, completionLabel + jobName, current, total)); } catch (Exception e) { log.warn("Chrome 后台 AI 分析任务 {} 失败: {}", taskId, e.getMessage(), e); if (claimed != null) { diff --git a/src/main/java/com/getjobs/application/service/CodexCliService.java b/src/main/java/com/getjobs/application/service/CodexCliService.java index 5d4f54d..ef31319 100644 --- a/src/main/java/com/getjobs/application/service/CodexCliService.java +++ b/src/main/java/com/getjobs/application/service/CodexCliService.java @@ -48,6 +48,7 @@ String run(String content, Path imagePath, Map config) { String executable = resolveExecutable(value(config, "CODEX_PATH", "codex")); String model = value(config, "CODEX_MODEL", "gpt-5.6-sol"); int timeoutSeconds = parseTimeout(value(config, "CODEX_TIMEOUT_SECONDS", "300")); + long deadlineNanos = System.nanoTime() + TimeUnit.SECONDS.toNanos(timeoutSeconds); Path tempDirectory = null; Path outputPath = null; @@ -66,19 +67,20 @@ String run(String content, Path imagePath, Map config) { builder.environment().put("CODEX_HOME", Path.of(codexHome).toAbsolutePath().normalize().toString()); } - slotAcquired = CODEX_SLOTS.tryAcquire(timeoutSeconds, TimeUnit.SECONDS); + slotAcquired = CODEX_SLOTS.tryAcquire(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); if (!slotAcquired) { throw new IllegalStateException("Codex CLI 当前任务过多,等待执行超时"); } + long processBudgetMillis = remainingMillis(deadlineNanos); + if (processBudgetMillis <= 0) { + throw new IllegalStateException("Codex CLI 总执行时间已耗尽"); + } process = builder.start(); try (var writer = process.outputWriter(StandardCharsets.UTF_8)) { writer.write(buildPrompt(content, imagePath != null)); } - if (!process.waitFor(timeoutSeconds, TimeUnit.SECONDS)) { - process.destroy(); - if (!process.waitFor(2, TimeUnit.SECONDS)) { - process.destroyForcibly(); - } + if (!process.waitFor(processBudgetMillis, TimeUnit.MILLISECONDS)) { + terminateProcessTree(process); throw new IllegalStateException("Codex CLI 执行超时(>" + timeoutSeconds + " 秒)"); } if (process.exitValue() != 0) { @@ -98,8 +100,8 @@ String run(String content, Path imagePath, Map config) { } catch (IOException e) { throw new IllegalStateException("Codex CLI 无法启动,请检查 CODEX_PATH", e); } finally { - if (process != null && process.isAlive()) { - process.destroyForcibly(); + if (process != null) { + terminateProcessTree(process); } if (slotAcquired) { CODEX_SLOTS.release(); @@ -226,6 +228,41 @@ private int parseTimeout(String raw) { } } + long remainingMillis(long deadlineNanos) { + long remainingNanos = deadlineNanos - System.nanoTime(); + if (remainingNanos <= 0) return 0; + return Math.max(1, TimeUnit.NANOSECONDS.toMillis(remainingNanos)); + } + + void terminateProcessTree(Process process) { + if (process == null) return; + List descendants; + try { + descendants = process.toHandle().descendants().toList(); + } catch (RuntimeException ignored) { + descendants = List.of(); + } + for (int i = descendants.size() - 1; i >= 0; i--) { + descendants.get(i).destroy(); + } + process.destroy(); + try { + if (!process.waitFor(2, TimeUnit.SECONDS)) { + for (int i = descendants.size() - 1; i >= 0; i--) { + ProcessHandle child = descendants.get(i); + if (child.isAlive()) child.destroyForcibly(); + } + if (process.isAlive()) process.destroyForcibly(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + for (ProcessHandle child : descendants) { + if (child.isAlive()) child.destroyForcibly(); + } + if (process.isAlive()) process.destroyForcibly(); + } + } + private String extensionForMimeType(String mimeType) { String normalized = mimeType == null ? "" : mimeType.toLowerCase(Locale.ROOT); if (normalized.contains("png")) return ".png"; diff --git a/src/main/java/com/getjobs/application/service/ConfigService.java b/src/main/java/com/getjobs/application/service/ConfigService.java index be9f876..f65a5c5 100644 --- a/src/main/java/com/getjobs/application/service/ConfigService.java +++ b/src/main/java/com/getjobs/application/service/ConfigService.java @@ -37,6 +37,7 @@ public class ConfigService { "BASE_URL", "API_KEY", "MODEL", + "AI_REQUEST_TIMEOUT_SECONDS", "CODEX_PATH", "CODEX_MODEL", "CODEX_TIMEOUT_SECONDS", @@ -83,6 +84,13 @@ public Map getUiConfigsAsMap() { configMap.put(key, config.getConfigValue()); } } + for (String key : UI_CONFIG_KEYS) { + if (SENSITIVE_UI_CONFIG_KEYS.contains(key) || configMap.containsKey(key)) continue; + String environmentValue = environment.getProperty(key); + if (environmentValue != null && !environmentValue.isBlank()) { + configMap.put(key, environmentValue.trim()); + } + } return configMap; } @@ -185,6 +193,7 @@ public Map getAiConfigs() { result.put("CODEX_HOME", optionalAiConfigValue("CODEX_HOME", "")); result.put("CODEX_MODEL", optionalAiConfigValue("CODEX_MODEL", "gpt-5.6-sol")); result.put("CODEX_TIMEOUT_SECONDS", optionalAiConfigValue("CODEX_TIMEOUT_SECONDS", "300")); + result.put("AI_REQUEST_TIMEOUT_SECONDS", optionalAiConfigValue("AI_REQUEST_TIMEOUT_SECONDS", "120")); if ("codex".equals(provider)) { result.put("BASE_URL", optionalAiConfigValue("BASE_URL", "")); result.put("API_KEY", optionalAiConfigValue("API_KEY", "")); @@ -357,6 +366,21 @@ private void validateUiConfigValue(String configKey, String configValue) { if ("CODEX_PATH".equals(configKey)) { codexCliService.validateExecutableName(configValue); } + if ("CODEX_TIMEOUT_SECONDS".equals(configKey) || "AI_REQUEST_TIMEOUT_SECONDS".equals(configKey)) { + validateTimeoutSeconds(configKey, configValue); + } + } + + private void validateTimeoutSeconds(String configKey, String configValue) { + try { + int value = Integer.parseInt(configValue == null ? "" : configValue.trim()); + int minimum = "CODEX_TIMEOUT_SECONDS".equals(configKey) ? 10 : 1; + if (value < minimum || value > 1800) { + throw new IllegalArgumentException(configKey + " 必须在 " + minimum + " 到 1800 秒之间"); + } + } catch (NumberFormatException e) { + throw new IllegalArgumentException(configKey + " 必须是整数秒数", e); + } } private String normalizeUiConfigKey(String configKey) { @@ -368,7 +392,8 @@ private String resolveConfigCategory(String configKey) { return "general"; } return switch (configKey) { - case "AI_PROVIDER", "BASE_URL", "API_KEY", "MODEL", "CODEX_PATH", "CODEX_HOME", "CODEX_MODEL", "CODEX_TIMEOUT_SECONDS" -> "ai"; + case "AI_PROVIDER", "BASE_URL", "API_KEY", "MODEL", "AI_REQUEST_TIMEOUT_SECONDS", + "CODEX_PATH", "CODEX_HOME", "CODEX_MODEL", "CODEX_TIMEOUT_SECONDS" -> "ai"; case "HOOK_URL", "BOT_IS_SEND" -> "notification"; default -> "general"; }; @@ -383,6 +408,7 @@ private String resolveConfigDescription(String configKey) { case "BASE_URL" -> "AI 服务地址"; case "API_KEY" -> "AI 服务密钥"; case "MODEL" -> "AI 模型名称"; + case "AI_REQUEST_TIMEOUT_SECONDS" -> "一次远程 AI 调用的总超时秒数"; case "CODEX_PATH" -> "Codex CLI 可执行文件"; case "CODEX_HOME" -> "Codex 登录配置目录"; case "CODEX_MODEL" -> "Codex 模型名称"; diff --git a/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java b/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java index aa619cd..286140e 100644 --- a/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java +++ b/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java @@ -27,16 +27,20 @@ import org.springframework.context.annotation.DependsOn; import java.io.ByteArrayInputStream; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; import java.time.LocalDateTime; import java.util.ArrayList; import java.util.Base64; import java.util.HashMap; +import java.util.HexFormat; import java.util.List; import java.util.Locale; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.BooleanSupplier; import java.util.stream.Collectors; @@ -239,14 +243,16 @@ public AnalysisResult analyzeJob(JobAnalysisRequest request, if (resumeText == null || resumeText.trim().isEmpty()) { if (!isLeaseCurrent(leaseIsCurrent)) return AnalysisResult.staleLease(); AnalysisResult result = AnalysisResult.failed(DeliveryStatus.AI_ANALYSIS_FAILED, "请先在AI配置页保存简历内容"); + result.setErrorCode("AI_RESUME_MISSING"); result.setPriorityCompany(priority); + AtomicReference storedResult = new AtomicReference<>(result); if (!executeLeaseWrite(leaseWriteGuard, () -> { - persistAnalysis(request, result, "{\"error\":\"missing resume\"}"); - updatePlatformCache(request, result); + storedResult.set(persistAndUpdate( + request, result, "{\"errorCode\":\"AI_RESUME_MISSING\"}", false)); })) { return AnalysisResult.staleLease(); } - return result; + return storedResult.get(); } String prompt = buildPrompt(resumeText, request, priority, threshold); @@ -267,26 +273,31 @@ public AnalysisResult analyzeJob(JobAnalysisRequest request, result.setDecision("SKIP"); } if (!isLeaseCurrent(leaseIsCurrent)) return AnalysisResult.staleLease(); + AtomicReference storedResult = new AtomicReference<>(result); if (!executeLeaseWrite(leaseWriteGuard, () -> { - persistAnalysis(request, result, raw); - updatePlatformCache(request, result); + storedResult.set(persistAndUpdate( + request, result, responseDiagnostic(raw), true)); })) { return AnalysisResult.staleLease(); } - return result; + return storedResult.get(); } catch (Exception e) { log.warn("AI岗位分析失败: {}", e.getMessage()); if (!isLeaseCurrent(leaseIsCurrent)) return AnalysisResult.staleLease(); AnalysisResult result = AnalysisResult.failed(DeliveryStatus.AI_ANALYSIS_FAILED, e.getMessage()); + result.setErrorCode(errorCode(e)); + result.setProviderOutcomeUnknown( + e instanceof AiProviderException providerError && providerError.isOutcomeUnknown()); result.setPriorityCompany(priority); result.setThreshold(threshold); + AtomicReference storedResult = new AtomicReference<>(result); if (!executeLeaseWrite(leaseWriteGuard, () -> { - persistAnalysis(request, result, "{\"error\":\"" + escape(e.getMessage()) + "\"}"); - updatePlatformCache(request, result); + storedResult.set(persistAndUpdate( + request, result, errorDiagnostic(e), true)); })) { return AnalysisResult.staleLease(); } - return result; + return storedResult.get(); } } @@ -326,16 +337,11 @@ public List generateBossSearchKeywords(List existingKeywords, in "只返回JSON数组,不要使用Markdown代码块,不要解释。关键词要适合直接填入 Boss 搜索框,优先2到8个字,避免过宽泛。\n" + "已配置关键词(不要重复):\n" + new JSONArray(existing).toString() + "\n\n" + "简历:\n" + limit(resumeText, 5000); - try { - String raw = aiService.sendRequest(prompt); - return parseKeywordArray(raw).stream() - .filter(s -> existing.stream().noneMatch(e -> e.equalsIgnoreCase(s))) - .limit(max) - .collect(Collectors.toList()); - } catch (Exception e) { - log.warn("AI生成Boss关键词失败: {}", e.getMessage()); - return List.of(); - } + String raw = aiService.sendRequest(prompt); + return parseKeywordArray(raw).stream() + .filter(s -> existing.stream().noneMatch(e -> e.equalsIgnoreCase(s))) + .limit(max) + .collect(Collectors.toList()); } private String buildPrompt(String resumeText, JobAnalysisRequest request, boolean priority, int threshold) { @@ -359,10 +365,37 @@ private String buildPrompt(String resumeText, JobAnalysisRequest request, boolea private AnalysisResult parseResult(String raw) { JSONObject obj = new JSONObject(repairJsonObject(extractJson(raw))); + for (String field : List.of("score", "decision", "summary", "strengths", "risks", "greeting")) { + if (!obj.has(field) || obj.isNull(field)) { + throw outputError("AI_OUTPUT_MISSING_FIELD", "AI 返回缺少字段: " + field, raw); + } + } + Object scoreValue = obj.opt("score"); + if (!(scoreValue instanceof Number number) + || number.doubleValue() != Math.rint(number.doubleValue()) + || number.intValue() < 0 + || number.intValue() > 100) { + throw outputError("AI_OUTPUT_INVALID_SCORE", "AI 返回 score 必须是 0 到 100 的整数", raw); + } + String decision = obj.optString("decision", "").trim().toUpperCase(Locale.ROOT); + if (!"APPLY".equals(decision) && !"SKIP".equals(decision)) { + throw outputError("AI_OUTPUT_INVALID_DECISION", "AI 返回 decision 必须是 APPLY 或 SKIP", raw); + } + if (!(obj.opt("summary") instanceof String summary) || summary.isBlank()) { + throw outputError("AI_OUTPUT_INVALID_SCHEMA", "AI 返回 summary 不能为空", raw); + } + if (!(obj.opt("strengths") instanceof JSONArray strengths) + || !(obj.opt("risks") instanceof JSONArray risks) + || !(obj.opt("greeting") instanceof String)) { + throw outputError("AI_OUTPUT_INVALID_SCHEMA", "AI 返回字段类型不符合约定", raw); + } + if (!containsOnlyStrings(strengths) || !containsOnlyStrings(risks)) { + throw outputError("AI_OUTPUT_INVALID_SCHEMA", "AI 返回 strengths/risks 必须是字符串数组", raw); + } AnalysisResult result = new AnalysisResult(); - result.setScore(clampScore(obj.has("score") ? obj.optInt("score") : 0)); - result.setDecision(obj.optString("decision", "SKIP")); - result.setSummary(obj.optString("summary", "")); + result.setScore(number.intValue()); + result.setDecision(decision); + result.setSummary(summary); result.setStrengths(toStringList(obj.opt("strengths"))); result.setRisks(toStringList(obj.opt("risks"))); result.setGreeting(obj.optString("greeting", "")); @@ -374,7 +407,11 @@ private List parseKeywordArray(String raw) { JSONArray arr = new JSONArray(json); List out = new ArrayList<>(); for (int i = 0; i < arr.length(); i++) { - String keyword = arr.optString(i, "").trim(); + Object value = arr.opt(i); + if (!(value instanceof String)) { + throw outputError("AI_OUTPUT_INVALID_SCHEMA", "AI 返回的搜索关键词必须是字符串数组", raw); + } + String keyword = ((String) value).trim(); if (!keyword.isEmpty() && out.stream().noneMatch(existing -> existing.equalsIgnoreCase(keyword))) { out.add(keyword); } @@ -383,7 +420,9 @@ private List parseKeywordArray(String raw) { } private String extractJson(String raw) { - if (raw == null || raw.trim().isEmpty()) return "{}"; + if (raw == null || raw.trim().isEmpty()) { + throw outputError("AI_OUTPUT_EMPTY", "AI 返回空内容", raw); + } String s = raw.trim(); if (s.startsWith("```")) { s = s.replaceFirst("^```[a-zA-Z]*\\s*", "").replaceFirst("\\s*```$", "").trim(); @@ -395,7 +434,9 @@ private String extractJson(String raw) { } private String repairJsonObject(String raw) { - if (raw == null || raw.trim().isEmpty()) return "{}"; + if (raw == null || raw.trim().isEmpty()) { + throw outputError("AI_OUTPUT_EMPTY", "AI 返回空内容", raw); + } String s = raw.trim() .replace('\u201c', '"') .replace('\u201d', '"') @@ -423,35 +464,7 @@ private String repairJsonObject(String raw) { } catch (Exception ignored) { } - return fallbackJsonFromText(raw); - } - - private String fallbackJsonFromText(String raw) { - String text = raw == null ? "" : raw.trim(); - JSONObject obj = new JSONObject(); - obj.put("score", extractScore(text)); - obj.put("decision", extractDecision(text)); - obj.put("summary", limit(text.isEmpty() ? "AI返回格式异常,已按跳过处理" : text, 500)); - obj.put("strengths", new JSONArray()); - obj.put("risks", new JSONArray(List.of("AI返回不是标准JSON,建议检查模型输出或重试分析"))); - obj.put("greeting", ""); - return obj.toString(); - } - - private int extractScore(String text) { - if (text == null) return 0; - java.util.regex.Matcher matcher = java.util.regex.Pattern.compile("(?i)(score|分数|得分)\\D{0,12}(\\d{1,3})").matcher(text); - if (matcher.find()) { - try { - return Math.max(0, Math.min(100, Integer.parseInt(matcher.group(2)))); - } catch (Exception ignored) { - } - } - return 0; - } - - private int clampScore(int score) { - return Math.max(0, Math.min(100, score)); + throw outputError("AI_OUTPUT_INVALID_JSON", "AI 返回无法修复为有效 JSON", raw); } private int resolveApplyThreshold(Long profileId, boolean priority) { @@ -464,16 +477,10 @@ private int resolveApplyThreshold(Long profileId, boolean priority) { return configured == null ? fallback : Math.max(0, Math.min(100, configured)); } - private String extractDecision(String text) { - if (text == null) return "SKIP"; - java.util.regex.Matcher matcher = java.util.regex.Pattern - .compile("(?i)(decision|决策)\\D{0,20}(APPLY|SKIP)") - .matcher(text); - return matcher.find() ? matcher.group(2).toUpperCase(Locale.ROOT) : "SKIP"; - } - private String extractJsonArray(String raw) { - if (raw == null || raw.trim().isEmpty()) return "[]"; + if (raw == null || raw.trim().isEmpty()) { + throw outputError("AI_OUTPUT_EMPTY", "AI 返回空内容", raw); + } String s = raw.trim(); if (s.startsWith("```")) { s = s.replaceFirst("^```[a-zA-Z]*\\s*", "").replaceFirst("\\s*```$", "").trim(); @@ -484,6 +491,13 @@ private String extractJsonArray(String raw) { return s; } + private boolean containsOnlyStrings(JSONArray values) { + for (int i = 0; i < values.length(); i++) { + if (!(values.opt(i) instanceof String)) return false; + } + return true; + } + private List toStringList(Object value) { List out = new ArrayList<>(); if (value instanceof JSONArray arr) { @@ -494,7 +508,49 @@ private List toStringList(Object value) { return out.stream().filter(s -> s != null && !s.isBlank()).collect(Collectors.toList()); } - private void persistAnalysis(JobAnalysisRequest request, AnalysisResult result, String raw) { + private AnalysisResult persistAndUpdate(JobAnalysisRequest request, + AnalysisResult result, + String diagnostic, + boolean providerWasCalled) { + if (!persistAnalysis(request, result, diagnostic)) { + AnalysisResult failure = AnalysisResult.failed( + DeliveryStatus.AI_ANALYSIS_FAILED, + "AI 分析结果持久化失败,任务未标记成功" + ); + failure.setErrorCode("AI_PERSISTENCE_FAILED"); + failure.setProviderOutcomeUnknown(providerWasCalled); + failure.setPriorityCompany(result.getPriorityCompany()); + failure.setThreshold(result.getThreshold()); + if (!updatePlatformCache(request, failure)) { + safelyResetAnalyzingStatus(request, failure.getSummary()); + } + return failure; + } + if (!updatePlatformCache(request, result)) { + AnalysisResult failure = AnalysisResult.failed( + DeliveryStatus.AI_ANALYSIS_FAILED, + "AI 结果已生成,但岗位状态写回失败,需要人工对账" + ); + failure.setErrorCode("AI_PLATFORM_WRITE_FAILED"); + failure.setProviderOutcomeUnknown(providerWasCalled); + failure.setPriorityCompany(result.getPriorityCompany()); + failure.setThreshold(result.getThreshold()); + safelyResetAnalyzingStatus(request, failure.getSummary()); + return failure; + } + return result; + } + + private void safelyResetAnalyzingStatus(JobAnalysisRequest request, String reason) { + try { + markAnalysisInterrupted(request, reason); + } catch (RuntimeException recoveryError) { + log.warn("AI 岗位状态写回失败后的安全复位也失败: platform={}, rowId={}, error={}", + request.getPlatform(), request.getJobRowId(), recoveryError.getMessage()); + } + } + + private boolean persistAnalysis(JobAnalysisRequest request, AnalysisResult result, String diagnostic) { try { JobAiAnalysisEntity entity = new JobAiAnalysisEntity(); entity.setProfileId(request.getProfileId()); @@ -510,63 +566,72 @@ private void persistAnalysis(JobAnalysisRequest request, AnalysisResult result, entity.setRisks(toJsonArray(result.getRisks())); entity.setGreeting(result.getGreeting()); entity.setPriorityCompany(Boolean.TRUE.equals(result.getPriorityCompany()) ? 1 : 0); - entity.setRawResponse(raw); + entity.setRawResponse(diagnostic); entity.setCreatedAt(LocalDateTime.now()); entity.setUpdatedAt(LocalDateTime.now()); - jobAiAnalysisMapper.insert(entity); + return jobAiAnalysisMapper.insert(entity) == 1; } catch (Exception e) { log.warn("保存AI分析结果失败: {}", e.getMessage()); + return false; } } - public void updatePlatformCache(JobAnalysisRequest request, AnalysisResult result) { - if (request == null || result == null) return; - String reason = result.toReasonText(); - if ("boss".equalsIgnoreCase(request.getPlatform())) { - BossJobDataEntity existing = findBossJobForAnalysis(request); - String nextStatus = DeliveryStatus.protectDelivered( - existing == null ? null : existing.getDeliveryStatus(), - DeliveryStatus.fromAiResult(result) - ); - BossJobDataEntity update = new BossJobDataEntity(); - update.setAiScore(result.getScore()); - update.setAiDecision(result.getDecision()); - update.setAiReason(reason); - update.setPriorityCompany(Boolean.TRUE.equals(result.getPriorityCompany()) ? 1 : 0); - if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { - update.setScanRunId(request.getScanRunId()); - } - if (existing == null || !DeliveryStatus.isFinalStatus(existing.getDeliveryStatus())) { - update.setDeliveryStatus(nextStatus); - } - update.setUpdatedAt(LocalDateTime.now()); - UpdateWrapper wrapper = bossUpdateWrapper(request); - applyExpectedBossStatus(wrapper, existing); - bossJobDataMapper.update(update, wrapper); - } else if ("zhilian".equalsIgnoreCase(request.getPlatform())) { - ZhilianJobDataEntity existing = findZhilianJobForAnalysis(request); - String nextStatus = DeliveryStatus.protectDelivered( - existing == null ? null : existing.getDeliveryStatus(), - DeliveryStatus.fromAiResult(result) - ); - ZhilianJobDataEntity update = new ZhilianJobDataEntity(); - update.setAiScore(result.getScore()); - update.setAiDecision(result.getDecision()); - update.setAiReason(reason); - update.setPriorityCompany(Boolean.TRUE.equals(result.getPriorityCompany()) ? 1 : 0); - if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { - update.setScanRunId(request.getScanRunId()); - } - if (request.getJobDescription() != null && !request.getJobDescription().isBlank()) { - update.setJobDescription(request.getJobDescription()); + public boolean updatePlatformCache(JobAnalysisRequest request, AnalysisResult result) { + if (request == null || result == null) return false; + try { + String reason = result.toReasonText(); + if ("boss".equalsIgnoreCase(request.getPlatform())) { + BossJobDataEntity existing = findBossJobForAnalysis(request); + String nextStatus = DeliveryStatus.protectDelivered( + existing == null ? null : existing.getDeliveryStatus(), + DeliveryStatus.fromAiResult(result) + ); + BossJobDataEntity update = new BossJobDataEntity(); + update.setAiScore(result.getScore()); + update.setAiDecision(result.getDecision()); + update.setAiReason(reason); + update.setPriorityCompany(Boolean.TRUE.equals(result.getPriorityCompany()) ? 1 : 0); + if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { + update.setScanRunId(request.getScanRunId()); + } + if (existing == null || !DeliveryStatus.isFinalStatus(existing.getDeliveryStatus())) { + update.setDeliveryStatus(nextStatus); + } + update.setUpdatedAt(LocalDateTime.now()); + UpdateWrapper wrapper = bossUpdateWrapper(request); + applyExpectedBossStatus(wrapper, existing); + return bossJobDataMapper.update(update, wrapper) == 1; } - if (existing == null || !DeliveryStatus.isFinalStatus(existing.getDeliveryStatus())) { - update.setDeliveryStatus(nextStatus); + if ("zhilian".equalsIgnoreCase(request.getPlatform())) { + ZhilianJobDataEntity existing = findZhilianJobForAnalysis(request); + String nextStatus = DeliveryStatus.protectDelivered( + existing == null ? null : existing.getDeliveryStatus(), + DeliveryStatus.fromAiResult(result) + ); + ZhilianJobDataEntity update = new ZhilianJobDataEntity(); + update.setAiScore(result.getScore()); + update.setAiDecision(result.getDecision()); + update.setAiReason(reason); + update.setPriorityCompany(Boolean.TRUE.equals(result.getPriorityCompany()) ? 1 : 0); + if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { + update.setScanRunId(request.getScanRunId()); + } + if (request.getJobDescription() != null && !request.getJobDescription().isBlank()) { + update.setJobDescription(request.getJobDescription()); + } + if (existing == null || !DeliveryStatus.isFinalStatus(existing.getDeliveryStatus())) { + update.setDeliveryStatus(nextStatus); + } + update.setUpdateTime(LocalDateTime.now()); + UpdateWrapper wrapper = zhilianUpdateWrapper(request); + applyExpectedZhilianStatus(wrapper, existing); + return zhilianJobDataMapper.update(update, wrapper) == 1; } - update.setUpdateTime(LocalDateTime.now()); - UpdateWrapper wrapper = zhilianUpdateWrapper(request); - applyExpectedZhilianStatus(wrapper, existing); - zhilianJobDataMapper.update(update, wrapper); + return false; + } catch (RuntimeException e) { + log.warn("写回 AI 平台状态失败: platform={}, rowId={}, error={}", + request.getPlatform(), request.getJobRowId(), e.getMessage()); + return false; } } @@ -791,8 +856,53 @@ private String safe(String s) { return s == null ? "" : s; } - private String escape(String s) { - return s == null ? "" : s.replace("\"", "\\\""); + private AiOutputException outputError(String code, String message, String raw) { + return new AiOutputException(code, message + "(" + responseFingerprint(raw) + ")"); + } + + private String responseDiagnostic(String raw) { + JSONObject diagnostic = new JSONObject(); + diagnostic.put("kind", "provider_response_fingerprint"); + diagnostic.put("length", raw == null ? 0 : raw.length()); + diagnostic.put("sha256", sha256(raw)); + return diagnostic.toString(); + } + + private String errorDiagnostic(Exception error) { + JSONObject diagnostic = new JSONObject(); + diagnostic.put("kind", "ai_error"); + diagnostic.put("errorCode", errorCode(error)); + diagnostic.put("message", limit(error == null ? "AI 分析失败" : safe(error.getMessage()), 300)); + if (error instanceof AiProviderException providerError) { + diagnostic.put("clientRequestId", safe(providerError.getClientRequestId())); + diagnostic.put("providerRequestId", safe(providerError.getProviderRequestId())); + diagnostic.put("httpStatus", providerError.getHttpStatus() == null + ? JSONObject.NULL + : providerError.getHttpStatus()); + } + return diagnostic.toString(); + } + + private String errorCode(Exception error) { + if (error instanceof AiOutputException outputError) return outputError.code(); + if (error instanceof AiProviderException providerError) { + return "AI_PROVIDER_" + providerError.getCode().name(); + } + return "AI_ANALYSIS_FAILED"; + } + + private String responseFingerprint(String raw) { + return "length=" + (raw == null ? 0 : raw.length()) + ", sha256=" + sha256(raw); + } + + private String sha256(String raw) { + try { + byte[] digest = MessageDigest.getInstance("SHA-256") + .digest((raw == null ? "" : raw).getBytes(StandardCharsets.UTF_8)); + return HexFormat.of().formatHex(digest); + } catch (Exception ignored) { + return "unavailable"; + } } private Long resolveAnalysisProfileId(JobAnalysisRequest request) { @@ -833,6 +943,19 @@ public static PlatformAnalysisState incomplete(String status) { } } + private static final class AiOutputException extends RuntimeException { + private final String code; + + private AiOutputException(String code, String message) { + super(message); + this.code = code; + } + + private String code() { + return code; + } + } + @FunctionalInterface public interface LeaseWriteGuard { boolean execute(Runnable action); @@ -849,6 +972,8 @@ public static class AnalysisResult { private Boolean priorityCompany; private Integer threshold; private boolean staleLease; + private String errorCode; + private boolean providerOutcomeUnknown; public boolean shouldApply() { return "APPLY".equalsIgnoreCase(decision); @@ -864,6 +989,7 @@ public String toReasonText() { map.put("strengths", strengths); map.put("risks", risks); map.put("threshold", threshold); + map.put("errorCode", errorCode); return new JSONObject(map).toString(); } diff --git a/src/main/java/com/getjobs/application/service/JobAnalysisTaskStore.java b/src/main/java/com/getjobs/application/service/JobAnalysisTaskStore.java index eecf1b4..9b63653 100644 --- a/src/main/java/com/getjobs/application/service/JobAnalysisTaskStore.java +++ b/src/main/java/com/getjobs/application/service/JobAnalysisTaskStore.java @@ -315,6 +315,18 @@ public boolean complete(long taskId, String leaseToken, boolean failed, String m return changed == 1; } + public boolean completeUnknown(long taskId, String leaseToken, String message) { + String now = dbTime(LocalDateTime.now()); + int changed = jdbcTemplate.update("UPDATE job_analysis_task SET status='UNKNOWN', processed_count=0, " + + "success_count=0, failed_count=0, message=?, last_error=?, completed_at=NULL, updated_at=?, " + + "lease_expires_at=NULL WHERE id=? AND status='LEASED' AND lease_owner=? " + + "AND lease_expires_at IS NOT NULL AND lease_expires_at>?", + firstNonBlank(message, "AI Provider 结果未知,需要人工确认"), + firstNonBlank(message, "AI Provider 结果未知,需要人工确认"), + now, taskId, leaseToken, now); + return changed == 1; + } + public boolean reconcileExpired(long taskId, String leaseToken, Status target, @@ -450,7 +462,7 @@ TaskRecord findById(long id) { return rows.isEmpty() ? null : rows.get(0); } - private TaskRecord findByIdAndProfile(long id, long profileId) { + TaskRecord findByIdAndProfile(long id, long profileId) { List rows = jdbcTemplate.query( "SELECT " + selectColumns() + " FROM job_analysis_task WHERE id=? AND profile_id=?", TASK_MAPPER, diff --git a/src/test/java/com/getjobs/application/controller/AiConfigControllerJobTaskTest.java b/src/test/java/com/getjobs/application/controller/AiConfigControllerJobTaskTest.java index fb3a958..59de9d2 100644 --- a/src/test/java/com/getjobs/application/controller/AiConfigControllerJobTaskTest.java +++ b/src/test/java/com/getjobs/application/controller/AiConfigControllerJobTaskTest.java @@ -2,6 +2,8 @@ import com.getjobs.application.service.ChromeJobAnalysisQueueService; import com.getjobs.application.service.JobAnalysisTaskStore; +import com.getjobs.application.service.JobAiAnalysisService; +import com.getjobs.application.service.DeliveryStatus; import com.getjobs.application.service.ProfileService; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -20,14 +22,17 @@ class AiConfigControllerJobTaskTest { private AiConfigController controller; private ProfileService profileService; private ChromeJobAnalysisQueueService queueService; + private JobAiAnalysisService jobAiAnalysisService; @BeforeEach void setUp() { controller = new AiConfigController(); profileService = mock(ProfileService.class); queueService = mock(ChromeJobAnalysisQueueService.class); + jobAiAnalysisService = mock(JobAiAnalysisService.class); ReflectionTestUtils.setField(controller, "profileService", profileService); ReflectionTestUtils.setField(controller, "chromeJobAnalysisQueueService", queueService); + ReflectionTestUtils.setField(controller, "jobAiAnalysisService", jobAiAnalysisService); when(profileService.getCurrentProfileId()).thenReturn(7L); } @@ -58,4 +63,23 @@ void retryUsesCurrentProfileAndReturnsRejectedTransition() { .containsEntry("message", "仅 FAILED 或 UNKNOWN 任务允许显式重试"); verify(queueService).retry(12L, 7L, false); } + + @Test + void directAnalyzeReturnsBadGatewayForExplicitAnalysisFailure() { + JobAiAnalysisService.JobAnalysisRequest request = new JobAiAnalysisService.JobAnalysisRequest(); + JobAiAnalysisService.AnalysisResult failure = JobAiAnalysisService.AnalysisResult.failed( + DeliveryStatus.AI_ANALYSIS_FAILED, + "AI Provider 返回空内容" + ); + failure.setErrorCode("AI_PROVIDER_EMPTY_RESPONSE"); + when(jobAiAnalysisService.analyzeJob(request)).thenReturn(failure); + + ResponseEntity> response = controller.analyzeJob(request); + + assertThat(response.getStatusCode().value()).isEqualTo(502); + assertThat(response.getBody()) + .containsEntry("success", false) + .containsEntry("message", "AI Provider 返回空内容") + .containsEntry("data", failure); + } } diff --git a/src/test/java/com/getjobs/application/controller/BossControllerAiKeywordTest.java b/src/test/java/com/getjobs/application/controller/BossControllerAiKeywordTest.java new file mode 100644 index 0000000..d4f739f --- /dev/null +++ b/src/test/java/com/getjobs/application/controller/BossControllerAiKeywordTest.java @@ -0,0 +1,62 @@ +package com.getjobs.application.controller; + +import com.getjobs.application.service.BossService; +import com.getjobs.application.service.ChromeJobAnalysisQueueService; +import com.getjobs.application.service.ConfigService; +import com.getjobs.application.service.CookieService; +import com.getjobs.application.service.JobAiAnalysisService; +import com.getjobs.application.service.ProfileService; +import com.getjobs.worker.boss.Boss; +import com.getjobs.worker.manager.PlaywrightManager; +import com.getjobs.worker.service.BossJobService; +import com.getjobs.worker.service.JobRunCoordinator; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.core.env.Environment; +import org.springframework.http.ResponseEntity; + +import java.util.List; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.anyList; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class BossControllerAiKeywordTest { + @Test + void providerFailureIsNotReportedAsEmptySuccess() { + JobAiAnalysisService analysisService = mock(JobAiAnalysisService.class); + when(analysisService.generateBossSearchKeywords(anyList(), anyInt())) + .thenThrow(new IllegalStateException("AI Provider 服务异常")); + BossController controller = new BossController( + mock(BossJobService.class), + mock(PlaywrightManager.class), + mock(CookieService.class), + mock(JobRunCoordinator.class), + mock(ConfigService.class), + mockBossProvider(), + mock(BossService.class), + mock(ProfileService.class), + analysisService, + mock(ChromeJobAnalysisQueueService.class), + mock(Environment.class) + ); + + ResponseEntity> response = controller.generateBossAiKeywords( + Map.of("existingKeywords", List.of("Java"), "limit", 3) + ); + + assertThat(response.getStatusCode().value()).isEqualTo(502); + assertThat(response.getBody()) + .containsEntry("success", false) + .containsEntry("message", "AI Provider 服务异常") + .containsEntry("keywords", List.of()); + } + + @SuppressWarnings("unchecked") + private ObjectProvider mockBossProvider() { + return mock(ObjectProvider.class); + } +} diff --git a/src/test/java/com/getjobs/application/service/AiServiceRemoteHttpTest.java b/src/test/java/com/getjobs/application/service/AiServiceRemoteHttpTest.java new file mode 100644 index 0000000..c90b5f8 --- /dev/null +++ b/src/test/java/com/getjobs/application/service/AiServiceRemoteHttpTest.java @@ -0,0 +1,282 @@ +package com.getjobs.application.service; + +import com.getjobs.application.mapper.AiMapper; +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.boot.test.system.CapturedOutput; +import org.springframework.boot.test.system.OutputCaptureExtension; + +import java.io.IOException; +import java.net.InetSocketAddress; +import java.net.ServerSocket; +import java.nio.charset.StandardCharsets; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.when; + +@ExtendWith({MockitoExtension.class, OutputCaptureExtension.class}) +class AiServiceRemoteHttpTest { + @Mock + private ConfigService configService; + @Mock + private AiMapper aiMapper; + @Mock + private ProfileService profileService; + @Mock + private CodexCliService codexCliService; + + private HttpServer server; + private ExecutorService serverExecutor; + private AiService service; + private String baseUrl; + + @BeforeEach + void setUp() throws IOException { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + serverExecutor = Executors.newCachedThreadPool(); + server.setExecutor(serverExecutor); + server.start(); + baseUrl = "http://127.0.0.1:" + server.getAddress().getPort(); + service = new AiService(configService, aiMapper, profileService, codexCliService); + } + + @AfterEach + void tearDown() { + if (server != null) server.stop(0); + if (serverExecutor != null) serverExecutor.shutdownNow(); + } + + @Test + void returnsStrictChatContentAndForwardsClientRequestId() { + List requestIds = new CopyOnWriteArrayList<>(); + server.createContext("/v1/chat/completions", exchange -> { + requestIds.add(exchange.getRequestHeaders().getFirst("X-Request-ID")); + respond(exchange, 200, + "{\"id\":\"provider-1\",\"model\":\"test-model\",\"choices\":[{\"message\":{\"content\":\"provider-ok\"}}]}"); + }); + configure("deepseek-chat", "2"); + + assertThat(service.sendRequest("hello")).isEqualTo("provider-ok"); + assertThat(requestIds).hasSize(1); + assertThat(requestIds.getFirst()).isNotBlank(); + } + + @Test + void retriesRateLimitOnceWithSameRequestId() { + AtomicInteger calls = new AtomicInteger(); + List requestIds = new CopyOnWriteArrayList<>(); + server.createContext("/v1/chat/completions", exchange -> { + requestIds.add(exchange.getRequestHeaders().getFirst("X-Request-ID")); + if (calls.incrementAndGet() == 1) { + exchange.getResponseHeaders().add("Retry-After", "0"); + respond(exchange, 429, "{\"error\":\"rate limited\"}"); + return; + } + respond(exchange, 200, + "{\"id\":\"provider-2\",\"choices\":[{\"message\":{\"content\":\"retried-ok\"}}]}"); + }); + configure("deepseek-chat", "2"); + + assertThat(service.sendRequest("hello")).isEqualTo("retried-ok"); + assertThat(calls).hasValue(2); + assertThat(requestIds).hasSize(2).doesNotContainNull(); + assertThat(requestIds.get(0)).isEqualTo(requestIds.get(1)); + } + + @Test + void doesNotRetryServerErrorOrLeakProviderBody(CapturedOutput output) { + AtomicInteger calls = new AtomicInteger(); + server.createContext("/v1/chat/completions", exchange -> { + calls.incrementAndGet(); + exchange.getResponseHeaders().add("x-request-id", "provider-safe-id"); + respond(exchange, 500, "{\"error\":\"secret-body-marker\"}"); + }); + configure("deepseek-chat", "2"); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, error -> { + assertThat(error.getCode()).isEqualTo(AiProviderException.Code.HTTP_5XX); + assertThat(error.getProviderRequestId()).isEqualTo("provider-safe-id"); + assertThat(error.isOutcomeUnknown()).isTrue(); + assertThat(error.getMessage()).doesNotContain("secret-body-marker"); + }); + assertThat(calls).hasValue(1); + assertThat(output).contains("bodyHash=").doesNotContain("secret-body-marker"); + } + + @Test + void totalDeadlineStopsSlowRequestWithoutRetrying() { + AtomicInteger calls = new AtomicInteger(); + server.createContext("/v1/chat/completions", exchange -> { + calls.incrementAndGet(); + try { + Thread.sleep(1500); + respond(exchange, 200, + "{\"choices\":[{\"message\":{\"content\":\"too-late\"}}]}"); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } catch (IOException ignored) { + // 客户端总超时后关闭连接是本测试的预期行为。 + } + }); + configure("deepseek-chat", "1"); + long started = System.nanoTime(); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, error -> { + assertThat(error.getCode()).isEqualTo(AiProviderException.Code.TIMEOUT); + assertThat(error.isOutcomeUnknown()).isTrue(); + }); + long elapsedMillis = (System.nanoTime() - started) / 1_000_000; + assertThat(elapsedMillis).isLessThan(2500); + assertThat(calls).hasValue(1); + } + + @Test + void networkFailureIsUnknownAndIsNotRetried() throws IOException { + final int unavailablePort; + try (ServerSocket socket = new ServerSocket(0)) { + unavailablePort = socket.getLocalPort(); + } + baseUrl = "http://127.0.0.1:" + unavailablePort; + configure("deepseek-chat", "1"); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, error -> { + assertThat(error.getCode()).isEqualTo(AiProviderException.Code.NETWORK); + assertThat(error.isOutcomeUnknown()).isTrue(); + }); + } + + @Test + void rejectsEmptySuccessEnvelope() { + server.createContext("/v1/chat/completions", exchange -> respond(exchange, 200, + "{\"choices\":[{\"message\":{\"content\":\" \"}}]}")); + configure("deepseek-chat", "2"); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, error -> { + assertThat(error.getCode()).isEqualTo(AiProviderException.Code.EMPTY_RESPONSE); + assertThat(error.isOutcomeUnknown()).isTrue(); + }); + } + + @Test + void rejectsEmptyHttpBody() { + server.createContext("/v1/chat/completions", exchange -> respond(exchange, 200, "")); + configure("deepseek-chat", "2"); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, error -> { + assertThat(error.getCode()).isEqualTo(AiProviderException.Code.EMPTY_RESPONSE); + assertThat(error.isOutcomeUnknown()).isTrue(); + }); + } + + @Test + void malformedRetryAfterDoesNotRetry() { + AtomicInteger calls = new AtomicInteger(); + server.createContext("/v1/chat/completions", exchange -> { + calls.incrementAndGet(); + exchange.getResponseHeaders().add("Retry-After", "not-a-delay"); + respond(exchange, 429, "{\"error\":\"rate limited\"}"); + }); + configure("deepseek-chat", "2"); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, error -> { + assertThat(error.getCode()).isEqualTo(AiProviderException.Code.RATE_LIMITED); + assertThat(error.isOutcomeUnknown()).isFalse(); + }); + assertThat(calls).hasValue(1); + } + + @Test + void serverErrorReasoningMarkerDoesNotTriggerEndpointFallback() { + AtomicInteger chatCalls = new AtomicInteger(); + AtomicInteger responsesCalls = new AtomicInteger(); + server.createContext("/v1/chat/completions", exchange -> { + chatCalls.incrementAndGet(); + respond(exchange, 500, + "{\"error\":{\"message\":\"reasoning.summary unsupported_value\"}}"); + }); + server.createContext("/v1/responses", exchange -> { + responsesCalls.incrementAndGet(); + respond(exchange, 200, "{\"output_text\":\"must-not-be-used\"}"); + }); + configure("deepseek-chat", "2"); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, + error -> assertThat(error.getCode()).isEqualTo(AiProviderException.Code.HTTP_5XX)); + assertThat(chatCalls).hasValue(1); + assertThat(responsesCalls).hasValue(0); + } + + @Test + void invalidSuccessEnvelopeRequiresConfirmedRetry() { + server.createContext("/v1/chat/completions", exchange -> respond(exchange, 200, "not-json")); + configure("deepseek-chat", "2"); + + assertThatThrownBy(() -> service.sendRequest("hello")) + .isInstanceOfSatisfying(AiProviderException.class, error -> { + assertThat(error.getCode()).isEqualTo(AiProviderException.Code.INVALID_RESPONSE); + assertThat(error.isOutcomeUnknown()).isTrue(); + }); + } + + @Test + void reasoningCompatibilityFallbackSharesTwoRequestBudget() { + AtomicInteger calls = new AtomicInteger(); + List requestIds = new CopyOnWriteArrayList<>(); + server.createContext("/v1/chat/completions", exchange -> { + calls.incrementAndGet(); + requestIds.add(exchange.getRequestHeaders().getFirst("X-Request-ID")); + respond(exchange, 400, + "{\"error\":{\"message\":\"reasoning.summary unsupported_value\"}}"); + }); + server.createContext("/v1/responses", exchange -> { + calls.incrementAndGet(); + requestIds.add(exchange.getRequestHeaders().getFirst("X-Request-ID")); + respond(exchange, 200, "{\"id\":\"provider-fallback\",\"output_text\":\"fallback-ok\"}"); + }); + configure("deepseek-chat", "2"); + + assertThat(service.sendRequest("hello")).isEqualTo("fallback-ok"); + assertThat(calls).hasValue(2); + assertThat(requestIds).hasSize(2); + assertThat(requestIds.get(0)).isEqualTo(requestIds.get(1)); + } + + private void configure(String model, String timeoutSeconds) { + when(configService.getAiConfigs()).thenReturn(Map.of( + "AI_PROVIDER", "api", + "BASE_URL", baseUrl, + "API_KEY", "fixture-key", + "MODEL", model, + "AI_REQUEST_TIMEOUT_SECONDS", timeoutSeconds + )); + } + + private void respond(HttpExchange exchange, int status, String body) throws IOException { + exchange.getRequestBody().readAllBytes(); + byte[] bytes = body.getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().set("Content-Type", "application/json; charset=utf-8"); + exchange.sendResponseHeaders(status, bytes.length); + exchange.getResponseBody().write(bytes); + exchange.close(); + } +} diff --git a/src/test/java/com/getjobs/application/service/ChromeJobAnalysisQueueServiceTest.java b/src/test/java/com/getjobs/application/service/ChromeJobAnalysisQueueServiceTest.java index 69ef022..17c3640 100644 --- a/src/test/java/com/getjobs/application/service/ChromeJobAnalysisQueueServiceTest.java +++ b/src/test/java/com/getjobs/application/service/ChromeJobAnalysisQueueServiceTest.java @@ -173,6 +173,65 @@ void completionWriteExceptionReconcilesPersistedPlatformSuccess() { verify(analysisService, timeout(3000).times(1)).analyzeJob(any(), any(), any()); } + @Test + void providerUnknownOutcomeStopsWithoutAutomaticDuplicateCall() { + JobAiAnalysisService.AnalysisResult unknown = JobAiAnalysisService.AnalysisResult.failed( + DeliveryStatus.AI_ANALYSIS_FAILED, + "provider timeout" + ); + unknown.setErrorCode("AI_PROVIDER_TIMEOUT"); + unknown.setProviderOutcomeUnknown(true); + when(analysisService.analyzeJob(any(), any(), any())).thenReturn(unknown); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + ChromeJobAnalysisQueueService.EnqueueResult submitted = queue.enqueue( + job(request("boss", "job-provider-unknown", "run-a"))); + + long taskId = submittedTaskId("job-provider-unknown"); + awaitStatus(taskId, "UNKNOWN"); + assertThat(store.findById(taskId).lastError()).contains("provider timeout"); + verify(analysisService, timeout(3000).times(1)).analyzeJob(any(), any(), any()); + queue.initialize(); + verify(analysisService, times(1)).analyzeJob(any(), any(), any()); + } + + @Test + void confirmedUnknownRetryResetsAnalyzingStatusBeforeCallingProviderAgain() { + long taskId = store.submit(request("boss", "job-confirmed-retry", "run-a")).task().id(); + assertThat(store.claim(taskId, "lease-unknown", Duration.ofMinutes(1))).isNotNull(); + assertThat(store.completeUnknown(taskId, "lease-unknown", "provider result unknown")).isTrue(); + when(analysisService.inspectPlatformAnalysis(any())) + .thenReturn(JobAiAnalysisService.PlatformAnalysisState.incomplete(DeliveryStatus.AI_ANALYZING)); + when(analysisService.markAnalysisInterrupted(any(), any())).thenReturn(true); + when(analysisService.analyzeJob(any(), any(), any())).thenReturn(successResult()); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + JobAnalysisTaskStore.RetryResult retried = queue.retry(taskId, 1L, true); + + assertThat(retried.accepted()).isTrue(); + verify(analysisService).markAnalysisInterrupted(any(), any()); + verify(analysisService, timeout(3000).times(1)).analyzeJob(any(), any(), any()); + awaitStatus(taskId, "SUCCEEDED"); + } + + @Test + void confirmedUnknownRetryFailsClosedWhenAnalyzingStatusCannotBeReset() { + long taskId = store.submit(request("boss", "job-reset-rejected", "run-a")).task().id(); + assertThat(store.claim(taskId, "lease-unknown", Duration.ofMinutes(1))).isNotNull(); + assertThat(store.completeUnknown(taskId, "lease-unknown", "provider result unknown")).isTrue(); + when(analysisService.inspectPlatformAnalysis(any())) + .thenReturn(JobAiAnalysisService.PlatformAnalysisState.incomplete(DeliveryStatus.AI_ANALYZING)); + when(analysisService.markAnalysisInterrupted(any(), any())).thenReturn(false); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + JobAnalysisTaskStore.RetryResult retried = queue.retry(taskId, 1L, true); + + assertThat(retried.accepted()).isFalse(); + assertThat(retried.message()).contains("未重新调用 AI Provider"); + assertThat(store.findById(taskId).status()).isEqualTo("UNKNOWN"); + verify(analysisService, never()).analyzeJob(any(), any(), any()); + } + private long leaseExpiredTask(String platform, String jobKey) { long taskId = store.submit(request(platform, jobKey, "run-a")).task().id(); assertThat(store.claim(taskId, "expired-lease", Duration.ofMinutes(1))).isNotNull(); diff --git a/src/test/java/com/getjobs/application/service/CodexCliServiceTest.java b/src/test/java/com/getjobs/application/service/CodexCliServiceTest.java index 119b936..28f5c87 100644 --- a/src/test/java/com/getjobs/application/service/CodexCliServiceTest.java +++ b/src/test/java/com/getjobs/application/service/CodexCliServiceTest.java @@ -4,9 +4,15 @@ import java.nio.file.Path; import java.util.List; +import java.util.concurrent.TimeUnit; +import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; class CodexCliServiceTest { @Test @@ -67,4 +73,33 @@ void executableValidationRejectsOtherProgramsAndInjectedArguments() { .isInstanceOf(IllegalArgumentException.class) .hasMessageContaining("不允许其他程序或命令参数"); } + + @Test + void remainingMillisUsesOneSharedDeadline() { + CodexCliService service = new CodexCliService(); + + assertThat(service.remainingMillis(System.nanoTime() - 1)).isZero(); + long futureDeadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2); + assertThat(service.remainingMillis(futureDeadline)).isBetween(1L, 2000L); + } + + @Test + void terminateProcessTreeStopsDescendantsBeforeParent() throws Exception { + CodexCliService service = new CodexCliService(); + Process process = mock(Process.class); + ProcessHandle parent = mock(ProcessHandle.class); + ProcessHandle firstChild = mock(ProcessHandle.class); + ProcessHandle secondChild = mock(ProcessHandle.class); + when(process.toHandle()).thenReturn(parent); + when(parent.descendants()).thenReturn(Stream.of(firstChild, secondChild)); + when(process.waitFor(2, TimeUnit.SECONDS)).thenReturn(true); + + service.terminateProcessTree(process); + + var order = inOrder(secondChild, firstChild, process); + order.verify(secondChild).destroy(); + order.verify(firstChild).destroy(); + order.verify(process).destroy(); + verify(process).waitFor(2, TimeUnit.SECONDS); + } } diff --git a/src/test/java/com/getjobs/application/service/ConfigServiceTest.java b/src/test/java/com/getjobs/application/service/ConfigServiceTest.java index d19e48a..7a5c938 100644 --- a/src/test/java/com/getjobs/application/service/ConfigServiceTest.java +++ b/src/test/java/com/getjobs/application/service/ConfigServiceTest.java @@ -18,6 +18,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -115,9 +116,29 @@ void codexIsDefaultAndDoesNotRequireApiKey() { .containsEntry("AI_PROVIDER", "codex") .containsEntry("CODEX_PATH", "codex") .containsEntry("CODEX_MODEL", "gpt-5.6-sol") + .containsEntry("AI_REQUEST_TIMEOUT_SECONDS", "120") .containsEntry("API_KEY", ""); } + @Test + void validatesRemoteApiTimeoutBeforeWriting() { + assertThatThrownBy(() -> configService.batchUpdateConfigs( + Map.of("AI_REQUEST_TIMEOUT_SECONDS", "0"))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("1 到 1800"); + assertThatThrownBy(() -> configService.batchUpdateConfigs( + Map.of("AI_REQUEST_TIMEOUT_SECONDS", "1801"))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("1 到 1800"); + assertThatThrownBy(() -> configService.batchUpdateConfigs( + Map.of("AI_REQUEST_TIMEOUT_SECONDS", "abc"))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("整数"); + + verify(configMapper, never()).insert(any(ConfigEntity.class)); + verify(configMapper, never()).updateById(any(ConfigEntity.class)); + } + @Test void apiKeyValueIsHiddenInLogs(CapturedOutput output) { when(configMapper.selectOne(any())).thenReturn(null); @@ -151,6 +172,17 @@ void uiConfigSnapshotDoesNotExposeSecretsOrCodexHome() { .doesNotContain("private/codex-home"); } + @Test + void uiConfigSnapshotShowsNonSensitiveEnvironmentFallback() { + when(configMapper.selectList(null)).thenReturn(List.of()); + when(environment.getProperty(anyString())).thenAnswer(invocation -> + "AI_REQUEST_TIMEOUT_SECONDS".equals(invocation.getArgument(0)) ? "10" : null); + + Map configs = configService.getUiConfigsAsMap(); + + assertThat(configs).containsEntry("AI_REQUEST_TIMEOUT_SECONDS", "10"); + } + @Test void blankSensitiveValuePreservesExistingConfig() { ConfigEntity model = config("MODEL", "old-model"); diff --git a/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java b/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java index f527a68..2f321f6 100644 --- a/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java +++ b/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java @@ -3,6 +3,7 @@ import com.baomidou.mybatisplus.core.conditions.update.UpdateWrapper; import com.getjobs.application.entity.AiEntity; import com.getjobs.application.entity.BossJobDataEntity; +import com.getjobs.application.entity.JobAiAnalysisEntity; import com.getjobs.application.entity.PriorityCompanyEntity; import com.getjobs.application.entity.ResumeProfileEntity; import com.getjobs.application.entity.ZhilianJobDataEntity; @@ -65,6 +66,10 @@ void setUp() { zhilianJobDataMapper ); lenient().when(priorityCompanyMapper.selectList(any())).thenReturn(List.of()); + lenient().when(jobAiAnalysisMapper.insert( + any(com.getjobs.application.entity.JobAiAnalysisEntity.class))).thenReturn(1); + lenient().when(bossJobDataMapper.update(any(), any(UpdateWrapper.class))).thenReturn(1); + lenient().when(zhilianJobDataMapper.update(any(), any(UpdateWrapper.class))).thenReturn(1); } @Test @@ -312,6 +317,157 @@ void repairsMarkdownWrappedAiJsonAndKeepsWaitingConfirmFlow() { assertThat(update.getDeliveryStatus()).isEqualTo(DeliveryStatus.WAITING_CONFIRM); } + @Test + void emptyProviderOutputBecomesExplicitAiFailureInsteadOfSkip() { + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.NOT_DELIVERED)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(aiService.sendRequest(any())).thenReturn(" "); + + JobAiAnalysisService.AnalysisResult result = service.analyzeJob(bossRequest()); + + assertThat(result.isFailure()).isTrue(); + assertThat(result.getErrorCode()).isEqualTo("AI_OUTPUT_EMPTY"); + assertThat(result.isProviderOutcomeUnknown()).isFalse(); + assertThat(lastBossUpdate().getDeliveryStatus()).isEqualTo(DeliveryStatus.AI_ANALYSIS_FAILED); + } + + @Test + void missingRequiredOutputFieldBecomesExplicitAiFailure() { + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.NOT_DELIVERED)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(aiService.sendRequest(any())).thenReturn(""" + {"score":80,"decision":"APPLY","summary":"匹配","strengths":[],"risks":[]} + """); + + JobAiAnalysisService.AnalysisResult result = service.analyzeJob(bossRequest()); + + assertThat(result.isFailure()).isTrue(); + assertThat(result.getErrorCode()).isEqualTo("AI_OUTPUT_MISSING_FIELD"); + assertThat(lastBossUpdate().getDeliveryStatus()).isEqualTo(DeliveryStatus.AI_ANALYSIS_FAILED); + } + + @Test + void invalidScoreAndArrayElementTypesAreRejected() { + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.NOT_DELIVERED)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(aiService.sendRequest(any())) + .thenReturn(""" + {"score":101,"decision":"APPLY","summary":"匹配","strengths":[],"risks":[],"greeting":"你好"} + """) + .thenReturn(""" + {"score":80,"decision":"APPLY","summary":"匹配","strengths":[1],"risks":[],"greeting":"你好"} + """); + + JobAiAnalysisService.AnalysisResult invalidScore = service.analyzeJob(bossRequest()); + JobAiAnalysisService.AnalysisResult invalidArray = service.analyzeJob(bossRequest()); + + assertThat(invalidScore.getErrorCode()).isEqualTo("AI_OUTPUT_INVALID_SCORE"); + assertThat(invalidArray.getErrorCode()).isEqualTo("AI_OUTPUT_INVALID_SCHEMA"); + } + + @Test + void invalidJsonAndDecisionAreRejected() { + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.NOT_DELIVERED)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(aiService.sendRequest(any())) + .thenReturn("not-json-at-all") + .thenReturn(""" + {"score":80,"decision":"MAYBE","summary":"匹配","strengths":[],"risks":[],"greeting":"你好"} + """); + + JobAiAnalysisService.AnalysisResult invalidJson = service.analyzeJob(bossRequest()); + JobAiAnalysisService.AnalysisResult invalidDecision = service.analyzeJob(bossRequest()); + + assertThat(invalidJson.getErrorCode()).isEqualTo("AI_OUTPUT_INVALID_JSON"); + assertThat(invalidDecision.getErrorCode()).isEqualTo("AI_OUTPUT_INVALID_DECISION"); + } + + @Test + void rawProviderResponseIsReplacedWithDiagnosticFingerprint() { + String marker = "sensitive-response-marker"; + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.NOT_DELIVERED)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(aiService.sendRequest(any())).thenReturn(""" + {"score":88,"decision":"APPLY","summary":"sensitive-response-marker","strengths":[],"risks":[],"greeting":"你好"} + """); + + JobAiAnalysisService.AnalysisResult result = service.analyzeJob(bossRequest()); + + assertThat(result.isFailure()).isFalse(); + ArgumentCaptor captor = ArgumentCaptor.forClass(JobAiAnalysisEntity.class); + verify(jobAiAnalysisMapper).insert(captor.capture()); + assertThat(captor.getValue().getRawResponse()) + .contains("provider_response_fingerprint", "sha256", "length") + .doesNotContain(marker); + } + + @Test + void persistenceFailureNeverReportsTaskSuccess() { + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.NOT_DELIVERED)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(jobAiAnalysisMapper.insert(any(JobAiAnalysisEntity.class))).thenReturn(0); + when(bossJobDataMapper.update(any(), any(UpdateWrapper.class))).thenReturn(1, 0, 1); + when(aiService.sendRequest(any())).thenReturn(""" + {"score":88,"decision":"APPLY","summary":"匹配","strengths":[],"risks":[],"greeting":"你好"} + """); + + JobAiAnalysisService.AnalysisResult result = service.analyzeJob(bossRequest()); + + assertThat(result.isFailure()).isTrue(); + assertThat(result.getErrorCode()).isEqualTo("AI_PERSISTENCE_FAILED"); + assertThat(result.isProviderOutcomeUnknown()).isTrue(); + assertThat(lastBossUpdate().getDeliveryStatus()).isEqualTo(DeliveryStatus.AI_ANALYSIS_FAILED); + } + + @Test + void platformWriteFailureCanBeConfirmedAndRetriedWithoutGettingStuckAnalyzing() { + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.AI_ANALYZING)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(bossJobDataMapper.update(any(), any(UpdateWrapper.class))).thenReturn(1, 0, 1, 1, 1); + when(aiService.sendRequest(any())).thenReturn(""" + {"score":88,"decision":"APPLY","summary":"匹配","strengths":[],"risks":[],"greeting":"你好"} + """); + + JobAiAnalysisService.AnalysisResult firstResult = service.analyzeJob(bossRequest()); + JobAiAnalysisService.AnalysisResult confirmedRetryResult = service.analyzeJob(bossRequest()); + + assertThat(firstResult.isFailure()).isTrue(); + assertThat(firstResult.getErrorCode()).isEqualTo("AI_PLATFORM_WRITE_FAILED"); + assertThat(firstResult.isProviderOutcomeUnknown()).isTrue(); + assertThat(confirmedRetryResult.isFailure()).isFalse(); + List updates = allBossUpdates(); + assertThat(updates).extracting(BossJobDataEntity::getDeliveryStatus) + .containsSequence( + DeliveryStatus.AI_ANALYZING, + DeliveryStatus.WAITING_CONFIRM, + DeliveryStatus.AI_ANALYSIS_FAILED, + DeliveryStatus.AI_ANALYZING, + DeliveryStatus.WAITING_CONFIRM + ); + verify(aiService, times(2)).sendRequest(any()); + } + + @Test + void providerTimeoutIsPersistedAsUnknownOutcome() { + when(bossJobDataMapper.selectOne(any())).thenReturn(bossJob(DeliveryStatus.NOT_DELIVERED)); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(aiService.sendRequest(any())).thenThrow(new AiProviderException( + AiProviderException.Code.TIMEOUT, + "AI Provider 请求超时(requestId=test-request)", + null, + "test-request", + "", + true, + null + )); + + JobAiAnalysisService.AnalysisResult result = service.analyzeJob(bossRequest()); + + assertThat(result.isFailure()).isTrue(); + assertThat(result.getErrorCode()).isEqualTo("AI_PROVIDER_TIMEOUT"); + assertThat(result.isProviderOutcomeUnknown()).isTrue(); + } + @Test void deliveryFailureStatusKeepsFailureTypeAndReason() { ZhilianService zhilianService = new ZhilianService(null, null, zhilianJobDataMapper, null, profileService); @@ -424,10 +580,14 @@ private PriorityCompanyEntity priorityCompany(String companyName) { } private BossJobDataEntity lastBossUpdate() { + List values = allBossUpdates(); + return values.get(values.size() - 1); + } + + private List allBossUpdates() { ArgumentCaptor captor = ArgumentCaptor.forClass(BossJobDataEntity.class); verify(bossJobDataMapper, atLeastOnce()).update(captor.capture(), any(UpdateWrapper.class)); - List values = captor.getAllValues(); - return values.get(values.size() - 1); + return captor.getAllValues(); } private ZhilianJobDataEntity lastZhilianUpdate() { diff --git a/src/test/java/com/getjobs/application/service/JobAnalysisTaskStoreTest.java b/src/test/java/com/getjobs/application/service/JobAnalysisTaskStoreTest.java index 4e9cfcf..c677932 100644 --- a/src/test/java/com/getjobs/application/service/JobAnalysisTaskStoreTest.java +++ b/src/test/java/com/getjobs/application/service/JobAnalysisTaskStoreTest.java @@ -124,6 +124,27 @@ void leaseHeartbeatStopsAtHardExecutionLimitAndTaskCanExpire() { .contains(taskId); } + @Test + void currentLeaseCanCompleteUnknownAndRequiresConfirmedRetry() { + long taskId = store.submit(request(1L, "boss", "job-provider-timeout", "run-a")).task().id(); + assertThat(store.claim(taskId, "lease-provider-timeout", Duration.ofMinutes(1))).isNotNull(); + + assertThat(store.completeUnknown(taskId, "wrong-lease", "provider result unknown")).isFalse(); + assertThat(store.completeUnknown( + taskId, + "lease-provider-timeout", + "provider result unknown" + )).isTrue(); + + JobAnalysisTaskStore.TaskRecord task = store.findById(taskId); + assertThat(task.status()).isEqualTo("UNKNOWN"); + assertThat(task.completedAt()).isNull(); + assertThat(task.leaseExpiresAt()).isNull(); + assertThat(task.lastError()).contains("provider result unknown"); + assertThat(store.retry(taskId, 1L).accepted()).isFalse(); + assertThat(store.retry(taskId, 1L, true).accepted()).isTrue(); + } + private boolean claimAfterBarrier(long taskId, String lease, CountDownLatch ready, diff --git a/tasks/2026-08-24-p1-2-provider-resilience.md b/tasks/2026-08-24-p1-2-provider-resilience.md new file mode 100644 index 0000000..4ab6340 --- /dev/null +++ b/tasks/2026-08-24-p1-2-provider-resilience.md @@ -0,0 +1,70 @@ +# P1.2 Provider 超时、重试与成本保护 + +## 背景 + +远程 AI API 目前只有连接超时,没有请求总超时;429 不读取 `Retry-After`,错误日志和异常消息会携带完整响应体。岗位分析对空内容、坏 JSON 和缺字段会降级为 `score=0/SKIP`,分析历史或平台缓存写入失败后任务仍可能显示成功。Codex CLI 的槽位等待和进程执行还会各自消耗一整份超时预算。 + +## 目标 + +1. 远程文本和图片请求都具备明确总超时,不再无限占用 AI worker。 +2. 自动重发仅限明确 429,最多一次,遵守且限制 `Retry-After`;500、网络断开和超时不盲目重发。 +3. 单次业务调用最多两次远程请求;reasoning endpoint fallback 与 429 重试共享预算。 +4. 每次调用生成本地 request id;日志和异常只保留安全元数据,不输出完整 Provider body、Prompt 或密钥。 +5. 空内容、不可修复 JSON、缺字段、非法 score/decision 进入明确 AI 失败,不再伪装为业务 `SKIP`。 +6. 分析历史插入失败或平台 CAS 未命中时,持久任务不得标记 SUCCEEDED。 +7. timeout/network/结果写回不确定进入任务 UNKNOWN,显式重试继续沿用 P1.1 的 `confirmUnknown=true` 门禁。 +8. Codex CLI 的槽位等待和进程执行共享一个总 deadline,超时后清理本次进程树。 + +## 允许修改范围 + +- `AiService` 的远程 HTTP 请求、错误分类、请求追踪和日志。 +- `CodexCliService` 的 deadline 与本次子进程清理。 +- `ConfigService`、环境配置页和示例配置中的远程 API 超时设置。 +- `JobAiAnalysisService` 的严格输出校验、脱敏诊断、写入结果判断。 +- `JobAnalysisTaskStore` / `ChromeJobAnalysisQueueService` 的 UNKNOWN 完成语义。 +- AI Controller 的失败响应,以及对应纯本地单元/故障注入测试。 + +## 禁止修改范围 + +- 不调用真实 Provider、Codex CLI 或招聘平台,不启动正式服务,不访问真实数据库。 +- 不修改 Prompt 内容、模型路由启发式、AI 阈值、投递状态机或数据库 Schema。 +- 不加入 Resilience4j、消息队列、多 Provider 智能路由、全局配额/计费系统。 +- 不对 timeout、网络异常、500/502/503/504 自动重试;这些结果可能已产生费用。 +- 不将 request id 描述为 Provider 幂等保证。 + +## 已确定实现要求 + +- 新增 `AI_REQUEST_TIMEOUT_SECONDS`,默认 120 秒,限制 1~1800 秒;同时写入 `HttpRequest.timeout`。 +- 每次远程业务调用使用同一个 UUID request id 和最多 2 次总请求预算。 +- 429 仅在首次请求、可解析且不超过 10 秒的 `Retry-After`(缺失时使用短默认值)下重试一次;不得与 reasoning fallback 叠加超过两次。 +- 错误分类至少包含 rate limit、HTTP 4xx、HTTP 5xx、timeout、network、empty response;异常消息不得拼接原始 body。 +- 错误日志只含 endpoint host、client request id、provider request id、status、body length 与 SHA-256 短摘要。 +- 合法 Markdown JSON 和有限语法修复继续兼容;无法修复或 schema 不完整必须失败。 +- 成功结果不再把完整 `raw_response` 写入数据库,只保留长度、hash 和请求诊断;结构化 score/decision/summary 等字段继续保存。 +- `job_ai_analysis` insert 返回 0/异常,或平台状态更新影响 0 行,必须返回明确错误结果;不得完成为 SUCCEEDED。 +- Codex CLI 总时长以单一 deadline 计算;只终止本次创建的进程及其 descendants。 + +## 验收标准 + +- 本地 HTTP fixture 覆盖:200、429+Retry-After、500、timeout、空 body、reasoning fallback;断言请求次数上限和 request id 一致。 +- 500/timeout 不自动重试;429 最多两次;reasoning fallback 最多两次。 +- 日志/异常/持久化诊断不包含 fixture 中的敏感响应正文。 +- Markdown 合法 JSON 保持成功;空、坏 JSON、缺字段、非法 score/decision 全部进入 AI_ANALYSIS_FAILED。 +- 历史写入失败、平台 update=0 时任务不进入 SUCCEEDED;timeout/network 任务进入 UNKNOWN。 +- Codex 纯单元测试验证 deadline 计算和进程清理 helper,不启动真实 Codex CLI。 +- Java 全量测试、前端 lint/typecheck/build 通过;真实数据库文件指纹和 sidecar 状态不变。 + +## 测试命令 + +```powershell +.\gradlew.bat clean test --console=plain +pnpm --dir front lint +pnpm --dir front exec tsc --noEmit +pnpm --dir front build +``` + +## 返回格式 + +- 超时、重试次数、UNKNOWN/FAILED 判定和成本边界说明。 +- 故障矩阵与测试证据。 +- 真实数据库未修改证据、diff、Commit、Push 与堆叠 PR。