diff --git a/src/main/java/com/getjobs/application/controller/AiConfigController.java b/src/main/java/com/getjobs/application/controller/AiConfigController.java index 42356f2..8c884ef 100644 --- a/src/main/java/com/getjobs/application/controller/AiConfigController.java +++ b/src/main/java/com/getjobs/application/controller/AiConfigController.java @@ -2,6 +2,8 @@ import com.getjobs.application.entity.AiEntity; import com.getjobs.application.service.AiService; +import com.getjobs.application.service.ChromeJobAnalysisQueueService; +import com.getjobs.application.service.JobAnalysisTaskStore; import com.getjobs.application.service.JobAiAnalysisService; import com.getjobs.application.service.ProfileService; import lombok.Data; @@ -34,6 +36,9 @@ public class AiConfigController { @Autowired private ProfileService profileService; + @Autowired + private ChromeJobAnalysisQueueService chromeJobAnalysisQueueService; + /** * 获取AI配置 * @return AI配置信息 @@ -257,6 +262,50 @@ public ResponseEntity> analyzeJob(@RequestBody JobAiAnalysis } } + @GetMapping("/job-analysis/tasks") + public ResponseEntity> listJobAnalysisTasks( + @RequestParam(name = "limit", defaultValue = "50") int limit + ) { + Map response = new HashMap<>(); + try { + long profileId = profileService.getCurrentProfileId(); + response.put("success", true); + response.put("data", chromeJobAnalysisQueueService.listTasks(profileId, limit)); + response.put("queueSize", chromeJobAnalysisQueueService.queueSize(profileId)); + response.put("message", "AI 分析任务读取成功"); + return ResponseEntity.ok(response); + } catch (Exception e) { + log.error("读取 AI 分析任务失败", e); + response.put("success", false); + response.put("message", "读取 AI 分析任务失败: " + e.getMessage()); + return ResponseEntity.internalServerError().body(response); + } + } + + @PostMapping("/job-analysis/tasks/{taskId}/retry") + public ResponseEntity> retryJobAnalysisTask( + @PathVariable long taskId, + @RequestParam(name = "confirmUnknown", defaultValue = "false") boolean confirmUnknown + ) { + Map response = new HashMap<>(); + try { + long profileId = profileService.getCurrentProfileId(); + JobAnalysisTaskStore.RetryResult result = chromeJobAnalysisQueueService.retry( + taskId, profileId, confirmUnknown); + response.put("success", result.accepted()); + response.put("data", result.task() == null ? null : result.task().toView()); + response.put("message", result.message()); + return result.accepted() + ? ResponseEntity.ok(response) + : ResponseEntity.badRequest().body(response); + } catch (Exception e) { + log.error("重试 AI 分析任务失败", e); + response.put("success", false); + response.put("message", "重试 AI 分析任务失败: " + e.getMessage()); + return ResponseEntity.internalServerError().body(response); + } + } + /** * 健康检查接口 * @return 服务状态 diff --git a/src/main/java/com/getjobs/application/controller/BossController.java b/src/main/java/com/getjobs/application/controller/BossController.java index d4ca520..add1be8 100644 --- a/src/main/java/com/getjobs/application/controller/BossController.java +++ b/src/main/java/com/getjobs/application/controller/BossController.java @@ -205,11 +205,11 @@ public ResponseEntity> receiveChromeJobs(@RequestBody Chrome continue; } - saved = bossService.updateDeliveryStatusById(saved.getId(), DeliveryStatus.AI_ANALYZING); JobAiAnalysisService.JobAnalysisRequest analysisRequest = new JobAiAnalysisService.JobAnalysisRequest(); analysisRequest.setProfileId(profileId); analysisRequest.setPlatform("boss"); analysisRequest.setJobKey(saved.getEncryptId()); + analysisRequest.setJobRowId(saved.getId()); analysisRequest.setKeyword(dto.getKeyword() == null ? request.getKeyword() : dto.getKeyword()); analysisRequest.setCompanyName(saved.getCompanyName()); analysisRequest.setJobName(saved.getJobName()); @@ -230,7 +230,6 @@ public ResponseEntity> receiveChromeJobs(@RequestBody Chrome ChromeJobAnalysisQueueService.EnqueueResult enqueueResult = chromeJobAnalysisQueueService.enqueue(job); if (enqueueResult.isRejected()) { - bossService.updateDeliveryStatusById(saved.getId(), firstNonBlank(currentStatus, DeliveryStatus.NOT_DELIVERED)); Map response = decorateListCollectionResponse( bossChromeJobsResponse(false, false, received, insertedOrUpdated, queued, skipped, insufficient, restored, autoDeliver, analyses), listOnlyCollection, listCollected, collectionWarnings diff --git a/src/main/java/com/getjobs/application/controller/ZhilianController.java b/src/main/java/com/getjobs/application/controller/ZhilianController.java index 4f79610..477830e 100644 --- a/src/main/java/com/getjobs/application/controller/ZhilianController.java +++ b/src/main/java/com/getjobs/application/controller/ZhilianController.java @@ -415,15 +415,11 @@ public ResponseEntity> receiveChromeJobs(@RequestBody Chrome sendZhilianProgress(JobProgressMessage.warning("zhilian", message)); continue; } - if (!isFinalZhilianStatus(currentStatus)) { - zhilianService.updateDeliveryStatusById(saved.getId(), DeliveryStatus.AI_ANALYZING); - saved = zhilianService.getZhilianJobById(saved.getId()); - } - JobAiAnalysisService.JobAnalysisRequest analysisRequest = new JobAiAnalysisService.JobAnalysisRequest(); analysisRequest.setProfileId(profileId); analysisRequest.setPlatform("zhilian"); analysisRequest.setJobKey(saved.getJobId()); + analysisRequest.setJobRowId(saved.getId()); analysisRequest.setKeyword(dto.getKeyword() == null ? request.getKeyword() : dto.getKeyword()); analysisRequest.setCompanyName(saved.getCompanyName()); analysisRequest.setJobName(saved.getJobTitle()); @@ -444,7 +440,6 @@ public ResponseEntity> receiveChromeJobs(@RequestBody Chrome ChromeJobAnalysisQueueService.EnqueueResult enqueueResult = chromeJobAnalysisQueueService.enqueue(job); if (enqueueResult.isRejected()) { - zhilianService.updateDeliveryStatusById(saved.getId(), firstNonBlank(currentStatus, DeliveryStatus.NOT_DELIVERED)); Map response = zhilianChromeJobsResponse( false, false, received, savedCount, queued, skipped, insufficient, restored, analyses ); diff --git a/src/main/java/com/getjobs/application/service/BossService.java b/src/main/java/com/getjobs/application/service/BossService.java index e51dc89..492a469 100644 --- a/src/main/java/com/getjobs/application/service/BossService.java +++ b/src/main/java/com/getjobs/application/service/BossService.java @@ -1687,9 +1687,21 @@ public Map clearBossAnalysisData() { conn.setAutoCommit(false); int analysisDeleted; + int tasksDeleted; int jobsDeleted; try (Statement st = conn.createStatement()) { Long profileId = profileService.getCurrentProfileId(); + tasksDeleted = st.executeUpdate("DELETE FROM job_analysis_task WHERE lower(platform)='boss' " + + "AND profile_id=" + profileId + " AND status<>'LEASED'"); + try (java.sql.ResultSet rs = st.executeQuery("SELECT COUNT(*) FROM job_analysis_task " + + "WHERE lower(platform)='boss' AND profile_id=" + profileId + " AND status='LEASED'")) { + if (rs.next() && rs.getLong(1) > 0) { + conn.rollback(); + resp.put("success", false); + resp.put("message", "仍有 Boss AI 分析正在执行,已阻止清空;请等待完成或进入 UNKNOWN 后再试"); + return resp; + } + } analysisDeleted = st.executeUpdate("DELETE FROM job_ai_analysis WHERE lower(platform)='boss' AND profile_id=" + profileId); jobsDeleted = st.executeUpdate("DELETE FROM boss_data WHERE profile_id=" + profileId); } @@ -1699,6 +1711,7 @@ public Map clearBossAnalysisData() { resp.put("message", "Boss投递分析数据已清空"); resp.put("jobsDeleted", jobsDeleted); resp.put("analysisDeleted", analysisDeleted); + resp.put("tasksDeleted", tasksDeleted); resp.put("total", 0); } catch (Exception e) { try { if (conn != null) conn.rollback(); } catch (Exception ignore) {} diff --git a/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java b/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java index 47b3a52..93a4cf2 100644 --- a/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java +++ b/src/main/java/com/getjobs/application/service/ChromeJobAnalysisQueueService.java @@ -1,139 +1,340 @@ package com.getjobs.application.service; import com.getjobs.worker.dto.JobProgressMessage; +import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.Data; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import java.time.Duration; import java.util.Map; import java.util.Objects; +import java.util.UUID; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ThreadFactory; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Consumer; +/** + * 持久 AI 任务的进程内执行器。任务生命周期由 JobAnalysisTaskStore 管理,线程池不是真相源。 + */ @Slf4j @Service public class ChromeJobAnalysisQueueService { private static final int AI_CONCURRENCY = 2; - private static final int QUEUE_CAPACITY = 200; - private static final long COMPLETED_KEY_TTL_MS = TimeUnit.MINUTES.toMillis(30); + private static final int LOCAL_QUEUE_CAPACITY = 200; + private static final Duration LEASE_DURATION = Duration.ofMinutes(5); + private static final long LEASE_HEARTBEAT_SECONDS = 60; private final JobAiAnalysisService jobAiAnalysisService; + private final JobAnalysisTaskStore taskStore; private final ThreadPoolExecutor executor; - private final java.util.Set activeKeys = ConcurrentHashMap.newKeySet(); - private final Map completedKeys = new ConcurrentHashMap<>(); + private final ScheduledExecutorService leaseHeartbeatExecutor; + private final java.util.Set locallyScheduledTaskIds = ConcurrentHashMap.newKeySet(); + private final Map runtimeJobs = new ConcurrentHashMap<>(); + private volatile boolean stopping; - public ChromeJobAnalysisQueueService(JobAiAnalysisService jobAiAnalysisService) { + public ChromeJobAnalysisQueueService(JobAiAnalysisService jobAiAnalysisService, + JobAnalysisTaskStore taskStore) { this.jobAiAnalysisService = jobAiAnalysisService; + this.taskStore = taskStore; this.executor = new ThreadPoolExecutor( AI_CONCURRENCY, AI_CONCURRENCY, 60L, TimeUnit.SECONDS, - new ArrayBlockingQueue<>(QUEUE_CAPACITY), + new ArrayBlockingQueue<>(LOCAL_QUEUE_CAPACITY), new NamedThreadFactory(), new ThreadPoolExecutor.AbortPolicy() ); this.executor.allowCoreThreadTimeOut(false); + this.leaseHeartbeatExecutor = java.util.concurrent.Executors.newSingleThreadScheduledExecutor( + new LeaseHeartbeatThreadFactory()); + } + + @PostConstruct + public void initialize() { + recoverOrphanedAnalyzingTasks(); + reconcileExpiredLeases(); + dispatchPendingTasks(); } public EnqueueResult enqueue(AnalysisJob job) { if (job == null || job.getRequest() == null) { return EnqueueResult.skipped("分析任务为空"); } - purgeCompletedKeys(); - - if (isFinalStatus(job.getCurrentStatus())) { - return EnqueueResult.skipped("岗位已完成AI分析或投递处理"); - } - - String dedupeKey = dedupeKey(job); - if (completedKeys.containsKey(dedupeKey) || !activeKeys.add(dedupeKey)) { - return EnqueueResult.skipped("重复AI分析任务"); + if (DeliveryStatus.isFinalStatus(job.getCurrentStatus())) { + return EnqueueResult.skipped("岗位已完成 AI 分析或投递处理"); } try { - executor.execute(() -> runAnalysis(job, dedupeKey)); + JobAnalysisTaskStore.SubmitResult submitted = taskStore.submit(job.getRequest()); + if (!submitted.accepted() || submitted.task() == null) { + return EnqueueResult.rejected(submitted.message()); + } + JobAnalysisTaskStore.TaskRecord task = submitted.task(); + if (task.statusEnum() == JobAnalysisTaskStore.Status.PENDING) { + if (job.getProgressCallback() != null) { + runtimeJobs.put(task.id(), job); + } + schedule(task.id()); + } + if (!submitted.created()) { + return EnqueueResult.skipped(submitted.message()); + } return EnqueueResult.queued(queueSize()); - } catch (RejectedExecutionException e) { - activeKeys.remove(dedupeKey); - return EnqueueResult.rejected("AI分析队列已满,请稍后再试"); + } catch (IllegalArgumentException e) { + return EnqueueResult.rejected(e.getMessage()); + } catch (Exception e) { + log.error("持久化 Chrome AI 分析任务失败", e); + return EnqueueResult.rejected("AI 分析任务持久化失败,请稍后重试"); } } public int queueSize() { - return executor.getQueue().size() + executor.getActiveCount(); + return taskStore.outstandingCount(); } - private void runAnalysis(AnalysisJob job, String dedupeKey) { - JobAiAnalysisService.JobAnalysisRequest request = job.getRequest(); - Consumer progress = job.getProgressCallback(); - String platform = Objects.toString(request.getPlatform(), ""); - String jobName = Objects.toString(request.getJobName(), ""); + public int queueSize(long profileId) { + return taskStore.outstandingCount(profileId); + } + + public java.util.List listTasks(long profileId, int limit) { + return taskStore.listRecent(profileId, limit); + } + public JobAnalysisTaskStore.RetryResult retry(long taskId, long profileId, boolean confirmUnknown) { + JobAnalysisTaskStore.RetryResult result = taskStore.retry(taskId, profileId, confirmUnknown); + if (result.accepted() && result.task() != null) { + schedule(result.task().id()); + } + return result; + } + + @Scheduled(fixedDelay = 15000) + public void maintainDurableQueue() { + if (stopping) return; + reconcileExpiredLeases(); + dispatchPendingTasks(); + } + + void dispatchPendingTasks() { + if (stopping) return; + for (JobAnalysisTaskStore.TaskRecord task : taskStore.listDuePending(LOCAL_QUEUE_CAPACITY)) { + schedule(task.id()); + } + } + + void reconcileExpiredLeases() { + if (stopping) return; + for (JobAnalysisTaskStore.TaskRecord task : taskStore.listExpiredLeases(LOCAL_QUEUE_CAPACITY)) { + reconcileExpiredLease(task); + } + } + + void recoverOrphanedAnalyzingTasks() { + if (stopping) return; + String reason = "升级或异常退出前的 AI 任务上下文已丢失,已登记为 UNKNOWN;请人工确认后显式重试"; + while (!stopping) { + java.util.List orphaned = + taskStore.listOrphanedAnalyzingRequests(LOCAL_QUEUE_CAPACITY); + if (orphaned.isEmpty()) return; + int created = 0; + for (JobAiAnalysisService.JobAnalysisRequest request : orphaned) { + try { + JobAnalysisTaskStore.SubmitResult recorded = taskStore.recordUnknown(request, reason); + if (recorded.created() && recorded.task() != null) { + created++; + jobAiAnalysisService.markAnalysisInterrupted(request, reason); + } + } catch (Exception e) { + log.warn("登记遗留 AI 分析中岗位失败: platform={}, rowId={}, error={}", + request.getPlatform(), request.getJobRowId(), e.getMessage(), e); + } + } + if (orphaned.size() < LOCAL_QUEUE_CAPACITY || created == 0) return; + } + } + + private void schedule(long taskId) { + if (stopping || !locallyScheduledTaskIds.add(taskId)) return; try { - emit(progress, JobProgressMessage.progress( - platform, - "AI分析中:" + jobName, - job.getCurrent(), - job.getTotal() - )); - JobAiAnalysisService.AnalysisResult result = jobAiAnalysisService.analyzeJob(request); + executor.execute(() -> runPersistedTask(taskId)); + } catch (RejectedExecutionException e) { + locallyScheduledTaskIds.remove(taskId); + log.debug("本地 AI executor 已满,任务 {} 保留为 PENDING", taskId); + } + } + + private void runPersistedTask(long taskId) { + String leaseToken = UUID.randomUUID().toString(); + JobAnalysisTaskStore.TaskRecord claimed = null; + JobAiAnalysisService.JobAnalysisRequest request = null; + ScheduledFuture heartbeat = null; + try { + claimed = taskStore.claim(taskId, leaseToken, LEASE_DURATION); + if (claimed == null) return; + heartbeat = startLeaseHeartbeat(taskId, leaseToken); + + request = taskStore.deserialize(claimed); + AnalysisJob runtimeJob = runtimeJobs.get(taskId); + Consumer progress = runtimeJob == null ? null : runtimeJob.getProgressCallback(); + int current = runtimeJob == null ? 0 : runtimeJob.getCurrent(); + int total = runtimeJob == null ? 0 : runtimeJob.getTotal(); + String platform = Objects.toString(request.getPlatform(), ""); + String jobName = Objects.toString(request.getJobName(), ""); - if (result.isFailure()) { - emit(progress, JobProgressMessage.warning( - platform, - "AI分析失败:" + jobName + "," + Objects.toString(result.getSummary(), "") - )); + emit(progress, JobProgressMessage.progress(platform, "AI分析中:" + jobName, current, total)); + JobAiAnalysisService.AnalysisResult result = jobAiAnalysisService.analyzeJob( + request, + () -> taskStore.isLeaseOwner(taskId, leaseToken), + action -> taskStore.executeWithLease(taskId, leaseToken, action) + ); + if (result.isStaleLease()) { + log.warn("AI 任务 {} 的租约已失效,旧执行结果已丢弃", taskId); + return; + } + boolean failed = result.isFailure(); + String summary = Objects.toString(result.getSummary(), ""); + boolean completed = taskStore.complete(taskId, leaseToken, failed, summary); + if (!completed) { + reconcileLateWorkerResult(claimed, leaseToken, failed, summary); } + if (failed) { + emit(progress, JobProgressMessage.warning(platform, "AI分析失败:" + jobName + "," + summary)); + } emit(progress, JobProgressMessage.progress( platform, - (result.shouldApply() ? DeliveryStatus.WAITING_CONFIRM + ":" : result.isFailure() ? DeliveryStatus.AI_ANALYSIS_FAILED + ":" : "跳过:") + jobName, - job.getCurrent(), - job.getTotal() + (result.shouldApply() + ? DeliveryStatus.WAITING_CONFIRM + ":" + : failed ? DeliveryStatus.AI_ANALYSIS_FAILED + ":" : "跳过:") + jobName, + current, + total )); } catch (Exception e) { - log.warn("{} Chrome后台AI分析任务失败: {}", platform, e.getMessage(), e); - emit(progress, JobProgressMessage.warning( - platform, - "AI分析失败:" + jobName + "," + e.getMessage() - )); + log.warn("Chrome 后台 AI 分析任务 {} 失败: {}", taskId, e.getMessage(), e); + if (claimed != null) { + finishAfterExecutionException(claimed, leaseToken, request, e); + } } finally { - activeKeys.remove(dedupeKey); - completedKeys.put(dedupeKey, System.currentTimeMillis()); + if (heartbeat != null) heartbeat.cancel(false); + runtimeJobs.remove(taskId); + locallyScheduledTaskIds.remove(taskId); + dispatchPendingTasks(); } } - private void emit(Consumer progress, JobProgressMessage message) { - if (progress == null || message == null) return; + private ScheduledFuture startLeaseHeartbeat(long taskId, String leaseToken) { + return leaseHeartbeatExecutor.scheduleAtFixedRate(() -> { + try { + if (!taskStore.renewLease(taskId, leaseToken, LEASE_DURATION)) { + log.warn("AI 任务 {} 的租约续期被拒绝,旧执行结果将被租约校验丢弃", taskId); + } + } catch (Exception e) { + log.warn("AI 任务 {} 的租约续期失败: {}", taskId, e.getMessage()); + } + }, LEASE_HEARTBEAT_SECONDS, LEASE_HEARTBEAT_SECONDS, TimeUnit.SECONDS); + } + + private void finishAfterExecutionException(JobAnalysisTaskStore.TaskRecord task, + String leaseToken, + JobAiAnalysisService.JobAnalysisRequest request, + Exception error) { + String message = firstNonBlank(error.getMessage(), "AI 分析执行异常"); try { - progress.accept(message); - } catch (Exception e) { - log.debug("发送后台AI分析进度失败: {}", e.getMessage()); + JobAiAnalysisService.JobAnalysisRequest recoverableRequest = request == null + ? taskStore.deserialize(task) + : request; + JobAiAnalysisService.PlatformAnalysisState platformState = + jobAiAnalysisService.inspectPlatformAnalysis(recoverableRequest); + if (platformState.completed()) { + boolean failed = platformState.failed(); + if (!taskStore.complete(task.id(), leaseToken, failed, + "执行异常后已从平台结果对账:" + platformState.status())) { + taskStore.reconcileUnknown( + task.id(), leaseToken, + failed ? JobAnalysisTaskStore.Status.FAILED : JobAnalysisTaskStore.Status.SUCCEEDED, + "执行异常后已从平台结果对账:" + platformState.status() + ); + } + return; + } + + boolean platformWriteAllowed = taskStore.executeWithLease( + task.id(), leaseToken, + () -> jobAiAnalysisService.markAnalysisInterrupted(recoverableRequest, message) + ); + if (!platformWriteAllowed) return; + boolean failed = taskStore.complete(task.id(), leaseToken, true, message); + if (!failed) { + taskStore.reconcileUnknown( + task.id(), leaseToken, JobAnalysisTaskStore.Status.FAILED, message); + } + } catch (Exception recoveryError) { + log.error("AI 任务 {} 的异常恢复失败,将保留租约等待过期对账: {}", + task.id(), recoveryError.getMessage(), recoveryError); } } - private boolean isFinalStatus(String status) { - return DeliveryStatus.isFinalStatus(status); + private void reconcileExpiredLease(JobAnalysisTaskStore.TaskRecord task) { + try { + JobAiAnalysisService.JobAnalysisRequest request = taskStore.deserialize(task); + JobAiAnalysisService.PlatformAnalysisState platformState = + jobAiAnalysisService.inspectPlatformAnalysis(request); + if (platformState.completed()) { + JobAnalysisTaskStore.Status target = platformState.failed() + ? JobAnalysisTaskStore.Status.FAILED + : JobAnalysisTaskStore.Status.SUCCEEDED; + taskStore.reconcileExpired( + task.id(), task.leaseOwner(), target, + "启动恢复已从平台结果对账:" + platformState.status() + ); + return; + } + + String reason = "AI 请求执行结果未知,已停止自动重试;请人工确认后显式重试"; + taskStore.reconcileExpired( + task.id(), task.leaseOwner(), JobAnalysisTaskStore.Status.UNKNOWN, reason, + () -> jobAiAnalysisService.markAnalysisInterrupted(request, reason) + ); + } catch (Exception e) { + String reason = "AI 任务恢复失败,已停止自动重试:" + firstNonBlank(e.getMessage(), "未知错误"); + taskStore.reconcileExpired( + task.id(), task.leaseOwner(), JobAnalysisTaskStore.Status.UNKNOWN, reason); + log.warn("恢复过期 AI 任务 {} 失败: {}", task.id(), e.getMessage(), e); + } } - private String dedupeKey(AnalysisJob job) { - JobAiAnalysisService.JobAnalysisRequest request = job.getRequest(); - String runId = Objects.toString(job.getRunId(), "no-run"); - String platform = Objects.toString(request.getPlatform(), ""); - String jobKey = firstNonBlank( - request.getJobKey(), - request.getCompanyName() + "::" + request.getJobName() + private void reconcileLateWorkerResult(JobAnalysisTaskStore.TaskRecord task, + String leaseToken, + boolean failed, + String message) { + JobAnalysisTaskStore.Status target = failed + ? JobAnalysisTaskStore.Status.FAILED + : JobAnalysisTaskStore.Status.SUCCEEDED; + taskStore.reconcileUnknown( + task.id(), leaseToken, target, + "过期执行器结果已完成对账:" + firstNonBlank(message, target.name()) ); - return runId + "::" + platform + "::" + jobKey; + } + + private void emit(Consumer progress, JobProgressMessage message) { + if (progress == null || message == null) return; + try { + progress.accept(message); + } catch (Exception e) { + log.debug("发送后台 AI 分析进度失败: {}", e.getMessage()); + } } private String firstNonBlank(String... values) { @@ -144,13 +345,10 @@ private String firstNonBlank(String... values) { return ""; } - private void purgeCompletedKeys() { - long cutoff = System.currentTimeMillis() - COMPLETED_KEY_TTL_MS; - completedKeys.entrySet().removeIf(entry -> entry.getValue() == null || entry.getValue() < cutoff); - } - @PreDestroy public void shutdown() { + stopping = true; + leaseHeartbeatExecutor.shutdownNow(); executor.shutdown(); try { if (!executor.awaitTermination(10, TimeUnit.SECONDS)) { @@ -203,4 +401,13 @@ public Thread newThread(Runnable runnable) { return thread; } } + + private static class LeaseHeartbeatThreadFactory implements ThreadFactory { + @Override + public Thread newThread(Runnable runnable) { + Thread thread = new Thread(runnable, "Chrome-AI-Lease-Heartbeat"); + thread.setDaemon(true); + return thread; + } + } } diff --git a/src/main/java/com/getjobs/application/service/DatabaseSchemaService.java b/src/main/java/com/getjobs/application/service/DatabaseSchemaService.java index 7a02c4f..d88a517 100644 --- a/src/main/java/com/getjobs/application/service/DatabaseSchemaService.java +++ b/src/main/java/com/getjobs/application/service/DatabaseSchemaService.java @@ -12,6 +12,7 @@ import java.sql.ResultSet; import java.sql.SQLException; import java.sql.Statement; +import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -587,6 +588,14 @@ private static long scalarCount(Statement stmt, String table) throws SQLExceptio } public static void validateSchema(Connection conn) throws SQLException { + validateSchema(conn, true); + } + + public static void validateSchemaBeforeV7(Connection conn) throws SQLException { + validateSchema(conn, false); + } + + private static void validateSchema(Connection conn, boolean requireV7TaskSchema) throws SQLException { List requiredTables = List.of( "profile", "config", "cookie", "ai", "resume_profile", "priority_company", "job_ai_analysis", "job_analysis_task", "boss_config", "boss_data", @@ -608,8 +617,14 @@ public static void validateSchema(Connection conn) throws SQLException { requiredColumns.put("zhilian_data", Set.of("profile_id", "job_id", "delivery_status", "scan_run_id")); requiredColumns.put("liepin_data", Set.of("job_id", "delivered")); requiredColumns.put("job51_data", Set.of("job_id", "delivered")); - requiredColumns.put("job_analysis_task", Set.of("profile_id", "platform", "status", "scan_run_id")); - List requiredIndexes = List.of( + requiredColumns.put("job_analysis_task", requireV7TaskSchema + ? Set.of( + "profile_id", "platform", "status", "scan_run_id", "task_key", "job_key", + "job_row_id", "request_json", "attempt_count", "next_retry_at", "lease_owner", + "lease_expires_at", "last_error", "started_at", "completed_at" + ) + : Set.of("profile_id", "platform", "status", "scan_run_id")); + List requiredIndexes = new ArrayList<>(List.of( "idx_priority_company_profile_name", "idx_boss_blacklist_type_value", "idx_boss_option_type_code", @@ -624,7 +639,15 @@ public static void validateSchema(Connection conn) throws SQLException { "idx_job51_data_company_job", "idx_job_analysis_task_profile_platform_status", "idx_job_analysis_task_scan_run" - ); + )); + if (requireV7TaskSchema) { + requiredIndexes.addAll(List.of( + "idx_job_analysis_task_task_key", + "idx_job_analysis_task_active_job", + "idx_job_analysis_task_dispatch", + "idx_job_analysis_task_lease" + )); + } try (Statement stmt = conn.createStatement()) { for (String table : requiredTables) { diff --git a/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java b/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java index 704db2f..aa619cd 100644 --- a/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java +++ b/src/main/java/com/getjobs/application/service/JobAiAnalysisService.java @@ -36,6 +36,8 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.BooleanSupplier; import java.util.stream.Collectors; @Slf4j @@ -205,19 +207,45 @@ private void evictPriorityCompanyCache(Long profileId) { } public AnalysisResult analyzeJob(JobAnalysisRequest request) { + return analyzeJob(request, () -> true, action -> { + action.run(); + return true; + }); + } + + public AnalysisResult analyzeJob(JobAnalysisRequest request, + BooleanSupplier leaseIsCurrent, + LeaseWriteGuard leaseWriteGuard) { if (request == null) throw new IllegalArgumentException("岗位分析请求不能为空"); + if (!isLeaseCurrent(leaseIsCurrent)) return AnalysisResult.staleLease(); Long profileId = resolveAnalysisProfileId(request); request.setProfileId(profileId); - markPlatformAnalysisStarted(request); + if (!isLeaseCurrent(leaseIsCurrent)) return AnalysisResult.staleLease(); + AtomicBoolean platformReserved = new AtomicBoolean(); + if (!executeLeaseWrite(leaseWriteGuard, + () -> platformReserved.set(markPlatformAnalysisStarted(request)))) { + return AnalysisResult.staleLease(); + } + if (!platformReserved.get()) { + return AnalysisResult.failed( + DeliveryStatus.AI_ANALYSIS_FAILED, + "岗位状态已变化或岗位不存在,未调用 AI Provider" + ); + } boolean priority = isPriorityCompany(request.getCompanyName(), profileId); int threshold = resolveApplyThreshold(profileId, priority); ResumeProfileEntity resume = getResumeProfile(profileId); String resumeText = resume == null ? "" : resume.getResumeText(); if (resumeText == null || resumeText.trim().isEmpty()) { + if (!isLeaseCurrent(leaseIsCurrent)) return AnalysisResult.staleLease(); AnalysisResult result = AnalysisResult.failed(DeliveryStatus.AI_ANALYSIS_FAILED, "请先在AI配置页保存简历内容"); result.setPriorityCompany(priority); - persistAnalysis(request, result, "{\"error\":\"missing resume\"}"); - updatePlatformCache(request, result); + if (!executeLeaseWrite(leaseWriteGuard, () -> { + persistAnalysis(request, result, "{\"error\":\"missing resume\"}"); + updatePlatformCache(request, result); + })) { + return AnalysisResult.staleLease(); + } return result; } @@ -238,20 +266,50 @@ public AnalysisResult analyzeJob(JobAnalysisRequest request) { if ("APPLY".equalsIgnoreCase(result.getDecision()) && result.getScore() < threshold) { result.setDecision("SKIP"); } - persistAnalysis(request, result, raw); - updatePlatformCache(request, result); + if (!isLeaseCurrent(leaseIsCurrent)) return AnalysisResult.staleLease(); + if (!executeLeaseWrite(leaseWriteGuard, () -> { + persistAnalysis(request, result, raw); + updatePlatformCache(request, result); + })) { + return AnalysisResult.staleLease(); + } return result; } 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.setPriorityCompany(priority); result.setThreshold(threshold); - persistAnalysis(request, result, "{\"error\":\"" + escape(e.getMessage()) + "\"}"); - updatePlatformCache(request, result); + if (!executeLeaseWrite(leaseWriteGuard, () -> { + persistAnalysis(request, result, "{\"error\":\"" + escape(e.getMessage()) + "\"}"); + updatePlatformCache(request, result); + })) { + return AnalysisResult.staleLease(); + } return result; } } + private boolean isLeaseCurrent(BooleanSupplier leaseIsCurrent) { + if (leaseIsCurrent == null) return false; + try { + return leaseIsCurrent.getAsBoolean(); + } catch (RuntimeException e) { + log.warn("验证 AI 分析任务租约失败,保守停止结果写入: {}", e.getMessage()); + return false; + } + } + + private boolean executeLeaseWrite(LeaseWriteGuard leaseWriteGuard, Runnable action) { + if (leaseWriteGuard == null || action == null) return false; + try { + return leaseWriteGuard.execute(action); + } catch (RuntimeException e) { + log.warn("AI 分析结果的租约事务失败,保守停止结果写入: {}", e.getMessage()); + return false; + } + } + public List generateBossSearchKeywords(List existingKeywords, int limitCount) { int max = Math.max(1, Math.min(limitCount <= 0 ? 5 : limitCount, 5)); ResumeProfileEntity resume = getResumeProfile(); @@ -478,11 +536,13 @@ public void updatePlatformCache(JobAnalysisRequest request, AnalysisResult resul if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { update.setScanRunId(request.getScanRunId()); } - if (!DeliveryStatus.isDeliveryLocked(existing == null ? null : existing.getDeliveryStatus())) { + if (existing == null || !DeliveryStatus.isFinalStatus(existing.getDeliveryStatus())) { update.setDeliveryStatus(nextStatus); } update.setUpdatedAt(LocalDateTime.now()); - bossJobDataMapper.update(update, bossUpdateWrapper(request)); + 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( @@ -500,30 +560,145 @@ public void updatePlatformCache(JobAnalysisRequest request, AnalysisResult resul if (request.getJobDescription() != null && !request.getJobDescription().isBlank()) { update.setJobDescription(request.getJobDescription()); } - if (!DeliveryStatus.isDeliveryLocked(existing == null ? null : existing.getDeliveryStatus())) { + if (existing == null || !DeliveryStatus.isFinalStatus(existing.getDeliveryStatus())) { update.setDeliveryStatus(nextStatus); } update.setUpdateTime(LocalDateTime.now()); - zhilianJobDataMapper.update(update, zhilianUpdateWrapper(request)); + UpdateWrapper wrapper = zhilianUpdateWrapper(request); + applyExpectedZhilianStatus(wrapper, existing); + zhilianJobDataMapper.update(update, wrapper); + } + } + + /** + * 只读取平台兼容状态,用于进程重启后判断过期租约是否已经完成结果写回。 + */ + public PlatformAnalysisState inspectPlatformAnalysis(JobAnalysisRequest request) { + if (request == null) return PlatformAnalysisState.incomplete("MISSING_REQUEST"); + if ("boss".equalsIgnoreCase(request.getPlatform())) { + BossJobDataEntity existing = findBossJobForAnalysis(request); + if (existing == null) return PlatformAnalysisState.incomplete("MISSING_JOB"); + return platformAnalysisState(existing.getDeliveryStatus()); + } + if ("zhilian".equalsIgnoreCase(request.getPlatform())) { + ZhilianJobDataEntity existing = findZhilianJobForAnalysis(request); + if (existing == null) return PlatformAnalysisState.incomplete("MISSING_JOB"); + return platformAnalysisState(existing.getDeliveryStatus()); + } + return PlatformAnalysisState.incomplete("UNSUPPORTED_PLATFORM"); + } + + /** + * 仅把仍停留在 AI_ANALYZING 的岗位转成明确失败;不会覆盖投递锁或已落库的 AI 结果。 + */ + public boolean markAnalysisInterrupted(JobAnalysisRequest request, String reason) { + if (request == null) return false; + String message = reason == null || reason.isBlank() + ? "AI 分析被中断,结果未知" + : reason.trim(); + if ("boss".equalsIgnoreCase(request.getPlatform())) { + BossJobDataEntity update = new BossJobDataEntity(); + update.setDeliveryStatus(DeliveryStatus.AI_ANALYSIS_FAILED); + update.setAiDecision(DeliveryStatus.AI_ANALYSIS_FAILED); + update.setAiReason(message); + update.setUpdatedAt(LocalDateTime.now()); + UpdateWrapper wrapper = bossUpdateWrapper(request); + wrapper.eq("delivery_status", DeliveryStatus.AI_ANALYZING); + return bossJobDataMapper.update(update, wrapper) == 1; + } + if ("zhilian".equalsIgnoreCase(request.getPlatform())) { + ZhilianJobDataEntity update = new ZhilianJobDataEntity(); + update.setDeliveryStatus(DeliveryStatus.AI_ANALYSIS_FAILED); + update.setAiDecision(DeliveryStatus.AI_ANALYSIS_FAILED); + update.setAiReason(message); + update.setUpdateTime(LocalDateTime.now()); + UpdateWrapper wrapper = zhilianUpdateWrapper(request); + wrapper.eq("delivery_status", DeliveryStatus.AI_ANALYZING); + return zhilianJobDataMapper.update(update, wrapper) == 1; + } + return false; + } + + private PlatformAnalysisState platformAnalysisState(String status) { + String normalizedStatus = status == null ? "" : status.trim(); + if (DeliveryStatus.AI_ANALYZING.equals(normalizedStatus)) { + return PlatformAnalysisState.incomplete(normalizedStatus); } + boolean failed = DeliveryStatus.AI_ANALYSIS_FAILED.equals(normalizedStatus); + boolean completed = failed + || DeliveryStatus.WAITING_CONFIRM.equals(normalizedStatus) + || DeliveryStatus.AI_NOT_MATCH.equals(normalizedStatus); + return new PlatformAnalysisState(completed, failed, + normalizedStatus.isBlank() ? "NO_STATUS" : normalizedStatus); } - private void markPlatformAnalysisStarted(JobAnalysisRequest request) { - if (request == null) return; + private boolean markPlatformAnalysisStarted(JobAnalysisRequest request) { + if (request == null) return false; if ("boss".equalsIgnoreCase(request.getPlatform())) { + if (request.getJobRowId() != null) { + BossJobDataEntity update = new BossJobDataEntity(); + update.setDeliveryStatus(DeliveryStatus.AI_ANALYZING); + update.setUpdatedAt(LocalDateTime.now()); + UpdateWrapper wrapper = bossUpdateWrapper(request); + wrapper.and(w -> w.in("delivery_status", List.of( + DeliveryStatus.NOT_DELIVERED, + DeliveryStatus.LIST_COLLECTED, + DeliveryStatus.AI_ANALYSIS_FAILED + )) + .or() + .isNull("delivery_status")); + return bossJobDataMapper.update(update, wrapper) == 1; + } BossJobDataEntity existing = findBossJobForAnalysis(request); - if (DeliveryStatus.isDeliveryLocked(existing == null ? null : existing.getDeliveryStatus())) return; + if (DeliveryStatus.isDeliveryLocked(existing == null ? null : existing.getDeliveryStatus())) return false; BossJobDataEntity update = new BossJobDataEntity(); update.setDeliveryStatus(DeliveryStatus.AI_ANALYZING); update.setUpdatedAt(LocalDateTime.now()); bossJobDataMapper.update(update, bossUpdateWrapper(request)); + return true; } else if ("zhilian".equalsIgnoreCase(request.getPlatform())) { + if (request.getJobRowId() != null) { + ZhilianJobDataEntity update = new ZhilianJobDataEntity(); + update.setDeliveryStatus(DeliveryStatus.AI_ANALYZING); + update.setUpdateTime(LocalDateTime.now()); + UpdateWrapper wrapper = zhilianUpdateWrapper(request); + wrapper.and(w -> w.in("delivery_status", List.of( + DeliveryStatus.NOT_DELIVERED, + DeliveryStatus.LIST_COLLECTED, + DeliveryStatus.AI_ANALYSIS_FAILED + )) + .or() + .isNull("delivery_status")); + return zhilianJobDataMapper.update(update, wrapper) == 1; + } ZhilianJobDataEntity existing = findZhilianJobForAnalysis(request); - if (DeliveryStatus.isDeliveryLocked(existing == null ? null : existing.getDeliveryStatus())) return; + if (DeliveryStatus.isDeliveryLocked(existing == null ? null : existing.getDeliveryStatus())) return false; ZhilianJobDataEntity update = new ZhilianJobDataEntity(); update.setDeliveryStatus(DeliveryStatus.AI_ANALYZING); update.setUpdateTime(LocalDateTime.now()); zhilianJobDataMapper.update(update, zhilianUpdateWrapper(request)); + return true; + } + return request.getJobRowId() == null; + } + + private void applyExpectedBossStatus(UpdateWrapper wrapper, + BossJobDataEntity existing) { + if (existing == null) return; + if (existing.getDeliveryStatus() == null) { + wrapper.isNull("delivery_status"); + } else { + wrapper.eq("delivery_status", existing.getDeliveryStatus()); + } + } + + private void applyExpectedZhilianStatus(UpdateWrapper wrapper, + ZhilianJobDataEntity existing) { + if (existing == null) return; + if (existing.getDeliveryStatus() == null) { + wrapper.isNull("delivery_status"); + } else { + wrapper.eq("delivery_status", existing.getDeliveryStatus()); } } @@ -532,12 +707,14 @@ private UpdateWrapper bossUpdateWrapper(JobAnalysisRequest re if (request.getProfileId() != null) { uw.eq("profile_id", request.getProfileId()); } - if (request.getJobKey() != null && !request.getJobKey().isBlank()) { + if (request.getJobRowId() != null) { + uw.eq("id", request.getJobRowId()); + } else if (request.getJobKey() != null && !request.getJobKey().isBlank()) { uw.eq("encrypt_id", request.getJobKey()); } else { uw.eq("company_name", request.getCompanyName()).eq("job_name", request.getJobName()); } - if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { + if (request.getJobRowId() == null && request.getScanRunId() != null && !request.getScanRunId().isBlank()) { uw.eq("scan_run_id", request.getScanRunId()); } return uw; @@ -548,12 +725,14 @@ private UpdateWrapper zhilianUpdateWrapper(JobAnalysisRequ if (request.getProfileId() != null) { uw.eq("profile_id", request.getProfileId()); } - if (request.getJobKey() != null && !request.getJobKey().isBlank()) { + if (request.getJobRowId() != null) { + uw.eq("id", request.getJobRowId()); + } else if (request.getJobKey() != null && !request.getJobKey().isBlank()) { uw.eq("job_id", request.getJobKey()); } else { uw.eq("company_name", request.getCompanyName()).eq("job_title", request.getJobName()); } - if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { + if (request.getJobRowId() == null && request.getScanRunId() != null && !request.getScanRunId().isBlank()) { uw.eq("scan_run_id", request.getScanRunId()); } return uw; @@ -564,12 +743,14 @@ private BossJobDataEntity findBossJobForAnalysis(JobAnalysisRequest request) { if (request.getProfileId() != null) { wrapper.eq("profile_id", request.getProfileId()); } - if (request.getJobKey() != null && !request.getJobKey().isBlank()) { + if (request.getJobRowId() != null) { + wrapper.eq("id", request.getJobRowId()); + } else if (request.getJobKey() != null && !request.getJobKey().isBlank()) { wrapper.eq("encrypt_id", request.getJobKey()); } else { wrapper.eq("company_name", request.getCompanyName()).eq("job_name", request.getJobName()); } - if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { + if (request.getJobRowId() == null && request.getScanRunId() != null && !request.getScanRunId().isBlank()) { wrapper.eq("scan_run_id", request.getScanRunId()); } wrapper.last("LIMIT 1"); @@ -581,12 +762,14 @@ private ZhilianJobDataEntity findZhilianJobForAnalysis(JobAnalysisRequest reques if (request.getProfileId() != null) { wrapper.eq("profile_id", request.getProfileId()); } - if (request.getJobKey() != null && !request.getJobKey().isBlank()) { + if (request.getJobRowId() != null) { + wrapper.eq("id", request.getJobRowId()); + } else if (request.getJobKey() != null && !request.getJobKey().isBlank()) { wrapper.eq("job_id", request.getJobKey()); } else { wrapper.eq("company_name", request.getCompanyName()).eq("job_title", request.getJobName()); } - if (request.getScanRunId() != null && !request.getScanRunId().isBlank()) { + if (request.getJobRowId() == null && request.getScanRunId() != null && !request.getScanRunId().isBlank()) { wrapper.eq("scan_run_id", request.getScanRunId()); } wrapper.last("LIMIT 1"); @@ -631,6 +814,7 @@ public static class JobAnalysisRequest { private Long profileId; private String platform; private String jobKey; + private Long jobRowId; private String keyword; private String companyName; private String jobName; @@ -643,6 +827,17 @@ public static class JobAnalysisRequest { private String scanRunId; } + public record PlatformAnalysisState(boolean completed, boolean failed, String status) { + public static PlatformAnalysisState incomplete(String status) { + return new PlatformAnalysisState(false, false, status); + } + } + + @FunctionalInterface + public interface LeaseWriteGuard { + boolean execute(Runnable action); + } + @Data public static class AnalysisResult { private Integer score; @@ -653,6 +848,7 @@ public static class AnalysisResult { private String greeting; private Boolean priorityCompany; private Integer threshold; + private boolean staleLease; public boolean shouldApply() { return "APPLY".equalsIgnoreCase(decision); @@ -679,5 +875,11 @@ public static AnalysisResult failed(String decision, String message) { result.setGreeting(""); return result; } + + public static AnalysisResult staleLease() { + AnalysisResult result = failed(DeliveryStatus.AI_ANALYSIS_FAILED, "AI 分析任务租约已失效,已丢弃旧执行结果"); + result.setStaleLease(true); + return result; + } } } diff --git a/src/main/java/com/getjobs/application/service/JobAnalysisTaskStore.java b/src/main/java/com/getjobs/application/service/JobAnalysisTaskStore.java new file mode 100644 index 0000000..eecf1b4 --- /dev/null +++ b/src/main/java/com/getjobs/application/service/JobAnalysisTaskStore.java @@ -0,0 +1,714 @@ +package com.getjobs.application.service; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import jakarta.annotation.PostConstruct; +import lombok.RequiredArgsConstructor; +import org.springframework.context.annotation.DependsOn; +import org.springframework.dao.DataAccessException; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.core.RowMapper; +import org.springframework.stereotype.Service; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.support.TransactionTemplate; + +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.time.Duration; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.time.format.DateTimeParseException; +import java.util.HexFormat; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Set; + +/** + * Chrome 岗位 AI 分析任务的持久事实源。线程池只是消费者,任务是否存在以本表为准。 + */ +@Service +@RequiredArgsConstructor +@DependsOn("databaseSchemaService") +public class JobAnalysisTaskStore { + public static final int MAX_OUTSTANDING_TASKS = 200; + public static final Duration MAX_EXECUTION_DURATION = Duration.ofMinutes(65); + private static final Set SUPPORTED_PLATFORMS = Set.of("boss", "zhilian"); + private static final DateTimeFormatter DB_TIME = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS"); + + private final JdbcTemplate jdbcTemplate; + private final PlatformTransactionManager transactionManager; + private final ObjectMapper objectMapper; + + private static final RowMapper TASK_MAPPER = (rs, rowNum) -> new TaskRecord( + rs.getLong("id"), + nullableLong(rs.getObject("profile_id")), + rs.getString("platform"), + rs.getString("job_key"), + nullableLong(rs.getObject("job_row_id")), + rs.getString("scan_run_id"), + rs.getString("status"), + rs.getInt("attempt_count"), + rs.getString("request_json"), + rs.getString("lease_owner"), + rs.getString("lease_expires_at"), + rs.getString("last_error"), + rs.getString("created_at"), + rs.getString("updated_at"), + rs.getString("started_at"), + rs.getString("completed_at") + ); + + @PostConstruct + public void validateSchema() { + List requiredColumns = List.of( + "task_key", "job_key", "job_row_id", "request_json", "attempt_count", + "next_retry_at", "lease_owner", "lease_expires_at", "last_error", + "started_at", "completed_at" + ); + for (String column : requiredColumns) { + Integer count = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM pragma_table_info('job_analysis_task') WHERE name=?", + Integer.class, + column + ); + if (count == null || count != 1) { + throw new IllegalStateException("AI 分析任务 schema 不完整,缺少字段: " + column); + } + } + for (String index : List.of( + "idx_job_analysis_task_task_key", + "idx_job_analysis_task_active_job", + "idx_job_analysis_task_dispatch", + "idx_job_analysis_task_lease" + )) { + Integer count = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM sqlite_master WHERE type='index' AND name=?", + Integer.class, + index + ); + if (count == null || count != 1) { + throw new IllegalStateException("AI 分析任务 schema 不完整,缺少索引: " + index); + } + } + } + + public SubmitResult submit(JobAiAnalysisService.JobAnalysisRequest request) { + validateRequest(request); + String platform = normalizePlatform(request.getPlatform()); + String jobKey = stableJobKey(request); + String taskKey = taskKey(request, platform, jobKey); + String requestJson = serialize(request); + TransactionTemplate transaction = new TransactionTemplate(transactionManager); + return transaction.execute(status -> { + TaskRecord exact = findByTaskKey(taskKey); + if (exact != null) { + return SubmitResult.existing(exact, duplicateMessage(exact)); + } + TaskRecord active = findActive(request.getProfileId(), platform, jobKey); + if (active != null) { + return SubmitResult.existing(active, "该岗位已有待执行或执行中的 AI 任务"); + } + if (outstandingCount() >= MAX_OUTSTANDING_TASKS) { + return SubmitResult.rejected("持久 AI 任务队列已满,请等待现有任务完成"); + } + + String now = dbTime(LocalDateTime.now()); + int inserted; + try { + inserted = jdbcTemplate.update("INSERT OR IGNORE INTO job_analysis_task (" + + "profile_id, platform, scan_run_id, status, total_count, processed_count, " + + "success_count, failed_count, message, created_at, updated_at, task_key, job_key, " + + "job_row_id, request_json, attempt_count) " + + "VALUES (?, ?, ?, 'PENDING', 1, 0, 0, 0, ?, ?, ?, ?, ?, ?, ?, 0)", + request.getProfileId(), platform, blankToNull(request.getScanRunId()), + "已持久化,等待 AI 分析", now, now, taskKey, jobKey, + request.getJobRowId(), requestJson); + } catch (DataAccessException e) { + inserted = 0; + } + TaskRecord stored = findByTaskKey(taskKey); + if (stored == null) { + stored = findActive(request.getProfileId(), platform, jobKey); + } + if (inserted == 1 && stored != null) { + return SubmitResult.created(stored); + } + if (stored != null) { + return SubmitResult.existing(stored, "该岗位任务已被并发创建"); + } + status.setRollbackOnly(); + return SubmitResult.rejected("AI 分析任务持久化失败"); + }); + } + + public SubmitResult recordUnknown(JobAiAnalysisService.JobAnalysisRequest request, String message) { + validateRequest(request); + String platform = normalizePlatform(request.getPlatform()); + String jobKey = stableJobKey(request); + String taskKey = taskKey(request, platform, jobKey); + String requestJson = serialize(request); + TransactionTemplate transaction = new TransactionTemplate(transactionManager); + return transaction.execute(status -> { + TaskRecord exact = findByTaskKey(taskKey); + if (exact != null) return SubmitResult.existing(exact, duplicateMessage(exact)); + TaskRecord active = findActive(request.getProfileId(), platform, jobKey); + if (active != null) return SubmitResult.existing(active, "该岗位已有持久 AI 任务"); + + String now = dbTime(LocalDateTime.now()); + int inserted = jdbcTemplate.update("INSERT OR IGNORE INTO job_analysis_task (" + + "profile_id, platform, scan_run_id, status, total_count, processed_count, " + + "success_count, failed_count, message, created_at, updated_at, task_key, job_key, " + + "job_row_id, request_json, attempt_count, last_error) " + + "VALUES (?, ?, ?, 'UNKNOWN', 1, 0, 0, 0, ?, ?, ?, ?, ?, ?, ?, 0, ?)", + request.getProfileId(), platform, blankToNull(request.getScanRunId()), + message, now, now, taskKey, jobKey, request.getJobRowId(), requestJson, message); + TaskRecord stored = findByTaskKey(taskKey); + if (inserted == 1 && stored != null) return SubmitResult.created(stored); + if (stored != null) return SubmitResult.existing(stored, "该岗位恢复任务已存在"); + status.setRollbackOnly(); + return SubmitResult.rejected("无法登记遗留 AI 分析中任务"); + }); + } + + public List listOrphanedAnalyzingRequests(int limit) { + int safeLimit = Math.max(1, Math.min(limit, MAX_OUTSTANDING_TASKS)); + List requests = new java.util.ArrayList<>(); + requests.addAll(jdbcTemplate.query("SELECT b.id, b.profile_id, b.encrypt_id AS job_key, " + + "b.source_keyword AS keyword, b.company_name, b.job_name, b.salary, b.location, " + + "b.experience, b.degree, b.introduce AS company_info, b.job_description, b.scan_run_id " + + "FROM boss_data b WHERE b.delivery_status=? AND b.profile_id IS NOT NULL " + + "AND NOT EXISTS (SELECT 1 FROM job_analysis_task t WHERE t.task_key IS NOT NULL " + + "AND lower(t.platform)='boss' AND t.profile_id=b.profile_id AND t.job_row_id=b.id " + + ") " + + "ORDER BY b.id LIMIT ?", + (rs, rowNum) -> analysisRequest( + "boss", rs.getLong("id"), rs.getLong("profile_id"), rs.getString("job_key"), + rs.getString("keyword"), rs.getString("company_name"), rs.getString("job_name"), + rs.getString("salary"), rs.getString("location"), rs.getString("experience"), + rs.getString("degree"), rs.getString("company_info"), rs.getString("job_description"), + rs.getString("scan_run_id") + ), DeliveryStatus.AI_ANALYZING, safeLimit)); + int remaining = safeLimit - requests.size(); + if (remaining > 0) { + requests.addAll(jdbcTemplate.query("SELECT z.id, z.profile_id, z.job_id AS job_key, " + + "z.company_name, z.job_title AS job_name, z.salary, z.location, z.experience, " + + "z.degree, z.job_description, z.scan_run_id FROM zhilian_data z " + + "WHERE z.delivery_status=? AND z.profile_id IS NOT NULL " + + "AND NOT EXISTS (SELECT 1 FROM job_analysis_task t WHERE t.task_key IS NOT NULL " + + "AND lower(t.platform)='zhilian' AND t.profile_id=z.profile_id AND t.job_row_id=z.id " + + ") " + + "ORDER BY z.id LIMIT ?", + (rs, rowNum) -> analysisRequest( + "zhilian", rs.getLong("id"), rs.getLong("profile_id"), rs.getString("job_key"), + "", rs.getString("company_name"), rs.getString("job_name"), + rs.getString("salary"), rs.getString("location"), rs.getString("experience"), + rs.getString("degree"), "", rs.getString("job_description"), + rs.getString("scan_run_id") + ), DeliveryStatus.AI_ANALYZING, remaining)); + } + return List.copyOf(requests); + } + + public List listDuePending(int limit) { + int safeLimit = Math.max(1, Math.min(limit, MAX_OUTSTANDING_TASKS)); + return jdbcTemplate.query("SELECT " + selectColumns() + " FROM job_analysis_task " + + "WHERE task_key IS NOT NULL AND request_json IS NOT NULL AND status='PENDING' " + + "AND (next_retry_at IS NULL OR next_retry_at<=?) ORDER BY id LIMIT ?", + TASK_MAPPER, dbTime(LocalDateTime.now()), safeLimit); + } + + public List listExpiredLeases(int limit) { + int safeLimit = Math.max(1, Math.min(limit, MAX_OUTSTANDING_TASKS)); + return jdbcTemplate.query("SELECT " + selectColumns() + " FROM job_analysis_task " + + "WHERE task_key IS NOT NULL AND request_json IS NOT NULL AND status='LEASED' " + + "AND lease_expires_at IS NOT NULL AND lease_expires_at<=? ORDER BY lease_expires_at LIMIT ?", + TASK_MAPPER, dbTime(LocalDateTime.now()), safeLimit); + } + + public TaskRecord claim(long taskId, String leaseToken, Duration leaseDuration) { + if (leaseToken == null || leaseToken.isBlank()) { + throw new IllegalArgumentException("leaseToken 不能为空"); + } + Duration safeDuration = leaseDuration == null || leaseDuration.isNegative() || leaseDuration.isZero() + ? Duration.ofMinutes(15) + : leaseDuration; + LocalDateTime now = LocalDateTime.now(); + String nowText = dbTime(now); + String leaseExpiry = dbTime(min(now.plus(safeDuration), now.plus(MAX_EXECUTION_DURATION))); + int changed = jdbcTemplate.update("UPDATE job_analysis_task SET status='LEASED', lease_owner=?, " + + "lease_expires_at=?, attempt_count=attempt_count+1, started_at=?, updated_at=?, " + + "message='AI 分析执行中', last_error=NULL WHERE id=? AND status='PENDING' " + + "AND task_key IS NOT NULL AND request_json IS NOT NULL " + + "AND (next_retry_at IS NULL OR next_retry_at<=?)", + leaseToken, leaseExpiry, nowText, nowText, taskId, nowText); + return changed == 1 ? findById(taskId) : null; + } + + public boolean renewLease(long taskId, String leaseToken, Duration leaseDuration) { + if (leaseToken == null || leaseToken.isBlank()) return false; + Duration safeDuration = leaseDuration == null || leaseDuration.isNegative() || leaseDuration.isZero() + ? Duration.ofMinutes(5) + : leaseDuration; + TaskRecord current = findById(taskId); + if (current == null || current.statusEnum() != Status.LEASED + || !leaseToken.equals(current.leaseOwner())) { + return false; + } + LocalDateTime startedAt = parseDbTime(current.startedAt()); + if (startedAt == null) return false; + LocalDateTime now = LocalDateTime.now(); + LocalDateTime hardDeadline = startedAt.plus(MAX_EXECUTION_DURATION); + if (!now.isBefore(hardDeadline)) return false; + String nowText = dbTime(now); + String leaseExpiry = dbTime(min(now.plus(safeDuration), hardDeadline)); + return jdbcTemplate.update("UPDATE job_analysis_task SET lease_expires_at=?, updated_at=? " + + "WHERE id=? AND status='LEASED' AND lease_owner=? AND started_at=?", + leaseExpiry, nowText, taskId, leaseToken, current.startedAt()) == 1; + } + + public boolean isLeaseOwner(long taskId, String leaseToken) { + if (leaseToken == null || leaseToken.isBlank()) return false; + Integer count = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM job_analysis_task WHERE id=? AND status='LEASED' AND lease_owner=? " + + "AND lease_expires_at IS NOT NULL AND lease_expires_at>?", + Integer.class, + taskId, + leaseToken, + dbTime(LocalDateTime.now()) + ); + return count != null && count == 1; + } + + /** + * 在持有 SQLite 写事务期间验证租约并执行结果写入,避免“校验通过后租约立刻失效”的竞态。 + */ + public boolean executeWithLease(long taskId, String leaseToken, Runnable action) { + if (leaseToken == null || leaseToken.isBlank() || action == null) return false; + TransactionTemplate transaction = new TransactionTemplate(transactionManager); + Boolean executed = transaction.execute(status -> { + String now = dbTime(LocalDateTime.now()); + int locked = jdbcTemplate.update("UPDATE job_analysis_task SET updated_at=updated_at " + + "WHERE id=? AND status='LEASED' AND lease_owner=? " + + "AND lease_expires_at IS NOT NULL AND lease_expires_at>?", + taskId, leaseToken, now); + if (locked != 1) return false; + action.run(); + return true; + }); + return Boolean.TRUE.equals(executed); + } + + public boolean complete(long taskId, String leaseToken, boolean failed, String message) { + String target = failed ? Status.FAILED.name() : Status.SUCCEEDED.name(); + String now = dbTime(LocalDateTime.now()); + int changed = jdbcTemplate.update("UPDATE job_analysis_task SET status=?, processed_count=1, " + + "success_count=?, failed_count=?, message=?, last_error=?, completed_at=?, 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>?", + target, failed ? 0 : 1, failed ? 1 : 0, + firstNonBlank(message, failed ? "AI 分析失败" : "AI 分析完成"), + failed ? firstNonBlank(message, "AI 分析失败") : null, + now, now, taskId, leaseToken, now); + return changed == 1; + } + + public boolean reconcileExpired(long taskId, + String leaseToken, + Status target, + String message) { + return reconcileExpired(taskId, leaseToken, target, message, null); + } + + public boolean reconcileExpired(long taskId, + String leaseToken, + Status target, + String message, + Runnable afterReconcile) { + if (target != Status.SUCCEEDED && target != Status.FAILED && target != Status.UNKNOWN) { + throw new IllegalArgumentException("过期租约只能对账为 SUCCEEDED、FAILED 或 UNKNOWN"); + } + TransactionTemplate transaction = new TransactionTemplate(transactionManager); + Boolean reconciled = transaction.execute(status -> { + String now = dbTime(LocalDateTime.now()); + boolean terminal = target == Status.SUCCEEDED || target == Status.FAILED; + int changed = jdbcTemplate.update("UPDATE job_analysis_task SET status=?, processed_count=?, " + + "success_count=?, failed_count=?, message=?, last_error=?, completed_at=?, 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<=?", + target.name(), terminal ? 1 : 0, + target == Status.SUCCEEDED ? 1 : 0, + target == Status.FAILED ? 1 : 0, + message, + target == Status.SUCCEEDED ? null : message, + terminal ? now : null, + now, taskId, leaseToken, now); + if (changed != 1) return false; + if (afterReconcile != null) afterReconcile.run(); + return true; + }); + return Boolean.TRUE.equals(reconciled); + } + + public boolean reconcileUnknown(long taskId, + String leaseToken, + Status target, + String message) { + if (target != Status.SUCCEEDED && target != Status.FAILED) { + throw new IllegalArgumentException("UNKNOWN 只能对账为 SUCCEEDED 或 FAILED"); + } + String now = dbTime(LocalDateTime.now()); + int changed = jdbcTemplate.update("UPDATE job_analysis_task SET status=?, processed_count=1, " + + "success_count=?, failed_count=?, message=?, last_error=?, completed_at=?, updated_at=? " + + "WHERE id=? AND status='UNKNOWN' AND lease_owner=?", + target.name(), target == Status.SUCCEEDED ? 1 : 0, target == Status.FAILED ? 1 : 0, + message, target == Status.FAILED ? message : null, now, now, taskId, leaseToken); + return changed == 1; + } + + public RetryResult retry(long taskId, long profileId) { + return retry(taskId, profileId, false); + } + + public RetryResult retry(long taskId, long profileId, boolean confirmUnknown) { + TransactionTemplate transaction = new TransactionTemplate(transactionManager); + return transaction.execute(status -> { + TaskRecord current = findByIdAndProfile(taskId, profileId); + if (current == null) { + return RetryResult.rejected("AI 分析任务不存在或不属于当前档案"); + } + Status state = current.statusEnum(); + if (state != Status.FAILED && state != Status.UNKNOWN) { + return RetryResult.rejected("仅 FAILED 或 UNKNOWN 任务允许显式重试"); + } + if (state == Status.UNKNOWN && !confirmUnknown) { + return RetryResult.rejected("UNKNOWN 任务可能已经产生外部调用,请确认平台结果并使用 confirmUnknown=true 后再重试"); + } + TaskRecord active = findActive(current.profileId(), current.platform(), current.jobKey()); + if (active != null && active.id() != current.id()) { + return RetryResult.rejected("该岗位已有待执行或执行中的 AI 任务"); + } + String now = dbTime(LocalDateTime.now()); + int changed = jdbcTemplate.update("UPDATE job_analysis_task SET status='PENDING', " + + "processed_count=0, success_count=0, failed_count=0, message='用户已显式重试,等待 AI 分析', " + + "next_retry_at=NULL, lease_owner=NULL, lease_expires_at=NULL, last_error=NULL, " + + "started_at=NULL, completed_at=NULL, updated_at=? WHERE id=? AND profile_id=? " + + "AND status IN ('FAILED', 'UNKNOWN')", + now, taskId, profileId); + if (changed != 1) { + status.setRollbackOnly(); + return RetryResult.rejected("任务状态已被并发更新,请刷新后重试"); + } + return RetryResult.accepted(findById(taskId)); + }); + } + + public List listRecent(long profileId, int limit) { + int safeLimit = Math.max(1, Math.min(limit, 200)); + return jdbcTemplate.query("SELECT " + selectColumns() + " FROM job_analysis_task " + + "WHERE profile_id=? AND task_key IS NOT NULL ORDER BY id DESC LIMIT ?", + TASK_MAPPER, profileId, safeLimit).stream().map(TaskRecord::toView).toList(); + } + + public int outstandingCount() { + Integer count = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM job_analysis_task WHERE task_key IS NOT NULL AND status IN ('PENDING','LEASED')", + Integer.class + ); + return count == null ? 0 : count; + } + + public int outstandingCount(long profileId) { + Integer count = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM job_analysis_task WHERE profile_id=? AND task_key IS NOT NULL " + + "AND status IN ('PENDING','LEASED')", + Integer.class, + profileId + ); + return count == null ? 0 : count; + } + + public JobAiAnalysisService.JobAnalysisRequest deserialize(TaskRecord task) { + if (task == null || task.requestJson() == null || task.requestJson().isBlank()) { + throw new IllegalArgumentException("任务缺少可恢复的请求快照"); + } + try { + return objectMapper.readValue(task.requestJson(), JobAiAnalysisService.JobAnalysisRequest.class); + } catch (JsonProcessingException e) { + throw new IllegalStateException("AI 分析任务请求快照损坏", e); + } + } + + TaskRecord findById(long id) { + List rows = jdbcTemplate.query( + "SELECT " + selectColumns() + " FROM job_analysis_task WHERE id=?", + TASK_MAPPER, + id + ); + return rows.isEmpty() ? null : rows.get(0); + } + + private TaskRecord findByIdAndProfile(long id, long profileId) { + List rows = jdbcTemplate.query( + "SELECT " + selectColumns() + " FROM job_analysis_task WHERE id=? AND profile_id=?", + TASK_MAPPER, + id, + profileId + ); + return rows.isEmpty() ? null : rows.get(0); + } + + private TaskRecord findByTaskKey(String taskKey) { + List rows = jdbcTemplate.query( + "SELECT " + selectColumns() + " FROM job_analysis_task WHERE task_key=?", + TASK_MAPPER, + taskKey + ); + return rows.isEmpty() ? null : rows.get(0); + } + + private TaskRecord findActive(Long profileId, String platform, String jobKey) { + List rows = jdbcTemplate.query( + "SELECT " + selectColumns() + " FROM job_analysis_task WHERE profile_id=? AND platform=? " + + "AND job_key=? AND task_key IS NOT NULL AND status IN ('PENDING','LEASED') ORDER BY id DESC LIMIT 1", + TASK_MAPPER, + profileId, + platform, + jobKey + ); + return rows.isEmpty() ? null : rows.get(0); + } + + private String taskKey(JobAiAnalysisService.JobAnalysisRequest request, String platform, String jobKey) { + Map inputs = new LinkedHashMap<>(); + inputs.put("profileId", request.getProfileId()); + inputs.put("platform", platform); + inputs.put("jobKey", jobKey); + inputs.put("keyword", canonical(request.getKeyword())); + inputs.put("companyName", canonical(request.getCompanyName())); + inputs.put("jobName", canonical(request.getJobName())); + inputs.put("salary", canonical(request.getSalary())); + inputs.put("location", canonical(request.getLocation())); + inputs.put("experience", canonical(request.getExperience())); + inputs.put("degree", canonical(request.getDegree())); + inputs.put("companyInfo", canonical(request.getCompanyInfo())); + inputs.put("jobDescription", canonical(request.getJobDescription())); + try { + String digest = sha256(objectMapper.writeValueAsString(inputs)); + return "ai:v1:" + request.getProfileId() + ":" + platform + ":" + digest; + } catch (JsonProcessingException e) { + throw new IllegalStateException("无法生成 AI 分析任务摘要", e); + } + } + + private JobAiAnalysisService.JobAnalysisRequest analysisRequest(String platform, + long rowId, + long profileId, + String jobKey, + String keyword, + String companyName, + String jobName, + String salary, + String location, + String experience, + String degree, + String companyInfo, + String jobDescription, + String scanRunId) { + JobAiAnalysisService.JobAnalysisRequest request = new JobAiAnalysisService.JobAnalysisRequest(); + request.setProfileId(profileId); + request.setPlatform(platform); + request.setJobKey(jobKey); + request.setJobRowId(rowId); + request.setKeyword(keyword); + request.setCompanyName(companyName); + request.setJobName(jobName); + request.setSalary(salary); + request.setLocation(location); + request.setExperience(experience); + request.setDegree(degree); + request.setCompanyInfo(companyInfo); + request.setJobDescription(jobDescription); + request.setScanRunId(scanRunId); + return request; + } + + private String serialize(JobAiAnalysisService.JobAnalysisRequest request) { + try { + return objectMapper.writeValueAsString(request); + } catch (JsonProcessingException e) { + throw new IllegalArgumentException("AI 分析请求无法持久化", e); + } + } + + private void validateRequest(JobAiAnalysisService.JobAnalysisRequest request) { + if (request == null) throw new IllegalArgumentException("AI 分析任务不能为空"); + if (request.getProfileId() == null) throw new IllegalArgumentException("AI 分析任务缺少档案 ID"); + String platform = normalizePlatform(request.getPlatform()); + if (!SUPPORTED_PLATFORMS.contains(platform)) { + throw new IllegalArgumentException("持久 AI 队列只支持 boss/zhilian"); + } + if (stableJobKey(request).isBlank()) { + throw new IllegalArgumentException("AI 分析任务缺少稳定岗位标识"); + } + } + + private String normalizePlatform(String platform) { + return platform == null ? "" : platform.trim().toLowerCase(Locale.ROOT); + } + + private String stableJobKey(JobAiAnalysisService.JobAnalysisRequest request) { + String direct = firstNonBlank(request.getJobKey()); + if (!direct.isBlank()) return direct; + String company = canonical(request.getCompanyName()); + String jobName = canonical(request.getJobName()); + return company.isBlank() && jobName.isBlank() ? "" : company + "::" + jobName; + } + + private String canonical(String value) { + return value == null ? "" : value.trim().replace("\r\n", "\n").replace('\r', '\n'); + } + + private static String dbTime(LocalDateTime value) { + return value == null ? null : DB_TIME.format(value); + } + + private static LocalDateTime parseDbTime(String value) { + if (value == null || value.isBlank()) return null; + for (DateTimeFormatter formatter : List.of(DB_TIME, DateTimeFormatter.ISO_LOCAL_DATE_TIME)) { + try { + return LocalDateTime.parse(value.trim(), formatter); + } catch (DateTimeParseException ignored) { + // 兼容升级前 SQLite CURRENT_TIMESTAMP 的秒级格式。 + } + } + try { + return LocalDateTime.parse(value.trim(), DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")); + } catch (DateTimeParseException ignored) { + return null; + } + } + + private static LocalDateTime min(LocalDateTime first, LocalDateTime second) { + return first.isBefore(second) ? first : second; + } + + private String sha256(String value) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + return HexFormat.of().formatHex(digest.digest(value.getBytes(StandardCharsets.UTF_8))); + } catch (Exception e) { + throw new IllegalStateException("当前运行时不支持 SHA-256", e); + } + } + + private String duplicateMessage(TaskRecord task) { + return switch (task.statusEnum()) { + case PENDING, LEASED -> "重复 AI 分析任务,已复用现有任务"; + case SUCCEEDED -> "相同岗位输入已经完成 AI 分析"; + case FAILED -> "相同岗位输入上次分析失败,需要显式重试"; + case UNKNOWN -> "相同岗位输入结果未知,需要先人工确认再显式重试"; + }; + } + + private String selectColumns() { + return "id, profile_id, platform, job_key, job_row_id, scan_run_id, status, attempt_count, " + + "request_json, lease_owner, lease_expires_at, last_error, created_at, updated_at, started_at, completed_at"; + } + + private static Long nullableLong(Object value) { + return value == null ? null : ((Number) value).longValue(); + } + + private static String blankToNull(String value) { + return value == null || value.isBlank() ? null : value.trim(); + } + + private static String firstNonBlank(String... values) { + if (values == null) return ""; + for (String value : values) { + if (value != null && !value.isBlank()) return value.trim(); + } + return ""; + } + + public enum Status { + PENDING, + LEASED, + SUCCEEDED, + FAILED, + UNKNOWN + } + + public record TaskRecord( + long id, + Long profileId, + String platform, + String jobKey, + Long jobRowId, + String scanRunId, + String status, + int attemptCount, + @JsonIgnore String requestJson, + @JsonIgnore String leaseOwner, + String leaseExpiresAt, + String lastError, + String createdAt, + String updatedAt, + String startedAt, + String completedAt + ) { + public Status statusEnum() { + return Status.valueOf(status); + } + + public TaskView toView() { + return new TaskView(id, profileId, platform, jobKey, jobRowId, scanRunId, status, + attemptCount, leaseExpiresAt, lastError, createdAt, updatedAt, startedAt, completedAt); + } + } + + public record TaskView( + long id, + Long profileId, + String platform, + String jobKey, + Long jobRowId, + String scanRunId, + String status, + int attemptCount, + String leaseExpiresAt, + String lastError, + String createdAt, + String updatedAt, + String startedAt, + String completedAt + ) { + } + + public record SubmitResult(boolean accepted, boolean created, TaskRecord task, String message) { + static SubmitResult created(TaskRecord task) { + return new SubmitResult(true, true, task, "AI 分析任务已持久化"); + } + + static SubmitResult existing(TaskRecord task, String message) { + return new SubmitResult(true, false, task, message); + } + + static SubmitResult rejected(String message) { + return new SubmitResult(false, false, null, message); + } + } + + public record RetryResult(boolean accepted, TaskRecord task, String message) { + static RetryResult accepted(TaskRecord task) { + return new RetryResult(true, task, "AI 分析任务已重新进入等待队列"); + } + + static RetryResult rejected(String message) { + return new RetryResult(false, null, message); + } + } +} diff --git a/src/main/java/com/getjobs/application/service/ProfileService.java b/src/main/java/com/getjobs/application/service/ProfileService.java index 616796a..01491e1 100644 --- a/src/main/java/com/getjobs/application/service/ProfileService.java +++ b/src/main/java/com/getjobs/application/service/ProfileService.java @@ -8,7 +8,9 @@ import org.springframework.context.annotation.DependsOn; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.stereotype.Service; +import org.springframework.transaction.PlatformTransactionManager; import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionTemplate; import java.time.LocalDateTime; import java.util.LinkedHashMap; @@ -27,11 +29,13 @@ public class ProfileService { "boss_data", "zhilian_data", "job_ai_analysis", + "job_analysis_task", "priority_company" ); private final ProfileMapper profileMapper; private final JdbcTemplate jdbcTemplate; + private final PlatformTransactionManager transactionManager; @Transactional(readOnly = true) public List listProfiles() { @@ -103,8 +107,14 @@ public ProfileEntity activateProfile(Long id) { return entity; } - @Transactional public DeleteProfileResult deleteProfile(Long id, boolean force) { + TransactionTemplate transaction = new TransactionTemplate(transactionManager); + return transaction.execute(status -> deleteProfileInTransaction(id, force, status)); + } + + private DeleteProfileResult deleteProfileInTransaction(Long id, + boolean force, + org.springframework.transaction.TransactionStatus status) { ProfileEntity entity = requireProfile(id); Long count = profileMapper.selectCount(null); if (count == null || count <= 1) { @@ -131,6 +141,22 @@ public DeleteProfileResult deleteProfile(Long id, boolean force) { boolean wasActive = entity.getIsActive() != null && entity.getIsActive() == 1; if (force) { + jdbcTemplate.update("DELETE FROM job_analysis_task WHERE profile_id=? AND status<>'LEASED'", id); + Long leasedTasks = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM job_analysis_task WHERE profile_id=? AND status='LEASED'", + Long.class, + id + ); + if (leasedTasks != null && leasedTasks > 0) { + status.setRollbackOnly(); + return new DeleteProfileResult( + false, + "该档案仍有 AI 分析正在执行,已阻止删除;请等待完成或进入 UNKNOWN 后再试。", + impactCounts, + getCurrentProfile(), + true + ); + } deleteProfileRelatedData(id); } profileMapper.deleteById(id); diff --git a/src/main/java/com/getjobs/application/service/ZhilianService.java b/src/main/java/com/getjobs/application/service/ZhilianService.java index 0bacf2e..eb7c367 100644 --- a/src/main/java/com/getjobs/application/service/ZhilianService.java +++ b/src/main/java/com/getjobs/application/service/ZhilianService.java @@ -821,9 +821,21 @@ public Map clearZhilianAnalysisData() { conn.setAutoCommit(false); int analysisDeleted; + int tasksDeleted; int jobsDeleted; try (Statement st = conn.createStatement()) { Long profileId = profileService.getCurrentProfileId(); + tasksDeleted = st.executeUpdate("DELETE FROM job_analysis_task WHERE lower(platform)='zhilian' " + + "AND profile_id=" + profileId + " AND status<>'LEASED'"); + try (java.sql.ResultSet rs = st.executeQuery("SELECT COUNT(*) FROM job_analysis_task " + + "WHERE lower(platform)='zhilian' AND profile_id=" + profileId + " AND status='LEASED'")) { + if (rs.next() && rs.getLong(1) > 0) { + conn.rollback(); + resp.put("success", false); + resp.put("message", "仍有智联 AI 分析正在执行,已阻止清空;请等待完成或进入 UNKNOWN 后再试"); + return resp; + } + } analysisDeleted = st.executeUpdate("DELETE FROM job_ai_analysis WHERE lower(platform)='zhilian' AND profile_id=" + profileId); jobsDeleted = st.executeUpdate("DELETE FROM zhilian_data WHERE profile_id=" + profileId); } @@ -833,6 +845,7 @@ public Map clearZhilianAnalysisData() { resp.put("message", "智联投递分析数据已清空"); resp.put("jobsDeleted", jobsDeleted); resp.put("analysisDeleted", analysisDeleted); + resp.put("tasksDeleted", tasksDeleted); resp.put("total", 0); } catch (Exception e) { try { if (conn != null) conn.rollback(); } catch (Exception ignore) {} diff --git a/src/main/java/db/migration/V5__consolidate_legacy_schema.java b/src/main/java/db/migration/V5__consolidate_legacy_schema.java index 14419aa..7f07006 100644 --- a/src/main/java/db/migration/V5__consolidate_legacy_schema.java +++ b/src/main/java/db/migration/V5__consolidate_legacy_schema.java @@ -22,7 +22,7 @@ public void migrate(Context context) throws Exception { executeResource(context, "db/migration/V4__add_job_analysis_task.sql"); DatabaseSchemaService.migrateLegacySchema(context.getConnection()); executeResource(context, "db/migration/V2__add_indexes.sql"); - DatabaseSchemaService.validateSchema(context.getConnection()); + DatabaseSchemaService.validateSchemaBeforeV7(context.getConnection()); } private void executeResource(Context context, String resourcePath) throws Exception { diff --git a/src/main/java/db/migration/V7__persist_job_analysis_tasks.java b/src/main/java/db/migration/V7__persist_job_analysis_tasks.java new file mode 100644 index 0000000..7e7d440 --- /dev/null +++ b/src/main/java/db/migration/V7__persist_job_analysis_tasks.java @@ -0,0 +1,56 @@ +package db.migration; + +import org.flywaydb.core.api.migration.BaseJavaMigration; +import org.flywaydb.core.api.migration.Context; + +import java.sql.ResultSet; +import java.sql.Statement; + +/** + * 将 V4 的批次统计表演进为逐岗位的持久 AI 任务表;旧聚合行原样保留且不会被调度。 + */ +public class V7__persist_job_analysis_tasks extends BaseJavaMigration { + @Override + public void migrate(Context context) throws Exception { + try (Statement statement = context.getConnection().createStatement()) { + addColumn(statement, "task_key", "TEXT"); + addColumn(statement, "job_key", "TEXT"); + addColumn(statement, "job_row_id", "INTEGER"); + addColumn(statement, "request_json", "TEXT"); + addColumn(statement, "attempt_count", "INTEGER NOT NULL DEFAULT 0"); + addColumn(statement, "next_retry_at", "DATETIME"); + addColumn(statement, "lease_owner", "TEXT"); + addColumn(statement, "lease_expires_at", "DATETIME"); + addColumn(statement, "last_error", "TEXT"); + addColumn(statement, "started_at", "DATETIME"); + addColumn(statement, "completed_at", "DATETIME"); + + statement.execute("CREATE UNIQUE INDEX IF NOT EXISTS idx_job_analysis_task_task_key " + + "ON job_analysis_task(task_key) WHERE task_key IS NOT NULL"); + statement.execute("CREATE UNIQUE INDEX IF NOT EXISTS idx_job_analysis_task_active_job " + + "ON job_analysis_task(profile_id, platform, job_key) " + + "WHERE task_key IS NOT NULL AND status IN ('PENDING', 'LEASED')"); + statement.execute("CREATE INDEX IF NOT EXISTS idx_job_analysis_task_dispatch " + + "ON job_analysis_task(status, next_retry_at, id)"); + statement.execute("CREATE INDEX IF NOT EXISTS idx_job_analysis_task_lease " + + "ON job_analysis_task(status, lease_expires_at)"); + } + } + + private void addColumn(Statement statement, String column, String definition) throws Exception { + if (!columnExists(statement, column)) { + statement.execute("ALTER TABLE job_analysis_task ADD COLUMN " + column + " " + definition); + } + } + + private boolean columnExists(Statement statement, String column) throws Exception { + try (ResultSet resultSet = statement.executeQuery("PRAGMA table_info('job_analysis_task')")) { + while (resultSet.next()) { + if (column.equalsIgnoreCase(resultSet.getString("name"))) { + return true; + } + } + } + return false; + } +} diff --git a/src/test/java/com/getjobs/application/controller/AiConfigControllerJobTaskTest.java b/src/test/java/com/getjobs/application/controller/AiConfigControllerJobTaskTest.java new file mode 100644 index 0000000..fb3a958 --- /dev/null +++ b/src/test/java/com/getjobs/application/controller/AiConfigControllerJobTaskTest.java @@ -0,0 +1,61 @@ +package com.getjobs.application.controller; + +import com.getjobs.application.service.ChromeJobAnalysisQueueService; +import com.getjobs.application.service.JobAnalysisTaskStore; +import com.getjobs.application.service.ProfileService; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.http.ResponseEntity; +import org.springframework.test.util.ReflectionTestUtils; + +import java.util.List; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class AiConfigControllerJobTaskTest { + private AiConfigController controller; + private ProfileService profileService; + private ChromeJobAnalysisQueueService queueService; + + @BeforeEach + void setUp() { + controller = new AiConfigController(); + profileService = mock(ProfileService.class); + queueService = mock(ChromeJobAnalysisQueueService.class); + ReflectionTestUtils.setField(controller, "profileService", profileService); + ReflectionTestUtils.setField(controller, "chromeJobAnalysisQueueService", queueService); + when(profileService.getCurrentProfileId()).thenReturn(7L); + } + + @Test + void taskListIsScopedToCurrentProfile() { + when(queueService.listTasks(7L, 20)).thenReturn(List.of()); + when(queueService.queueSize(7L)).thenReturn(3); + + ResponseEntity> response = controller.listJobAnalysisTasks(20); + + assertThat(response.getStatusCode().is2xxSuccessful()).isTrue(); + assertThat(response.getBody()) + .containsEntry("success", true) + .containsEntry("queueSize", 3); + verify(queueService).listTasks(7L, 20); + } + + @Test + void retryUsesCurrentProfileAndReturnsRejectedTransition() { + when(queueService.retry(12L, 7L, false)).thenReturn( + new JobAnalysisTaskStore.RetryResult(false, null, "仅 FAILED 或 UNKNOWN 任务允许显式重试")); + + ResponseEntity> response = controller.retryJobAnalysisTask(12L, false); + + assertThat(response.getStatusCode().is4xxClientError()).isTrue(); + assertThat(response.getBody()) + .containsEntry("success", false) + .containsEntry("message", "仅 FAILED 或 UNKNOWN 任务允许显式重试"); + verify(queueService).retry(12L, 7L, false); + } +} diff --git a/src/test/java/com/getjobs/application/service/AnalysisClearSequenceSafetyTest.java b/src/test/java/com/getjobs/application/service/AnalysisClearSequenceSafetyTest.java index 8f9319e..5da373a 100644 --- a/src/test/java/com/getjobs/application/service/AnalysisClearSequenceSafetyTest.java +++ b/src/test/java/com/getjobs/application/service/AnalysisClearSequenceSafetyTest.java @@ -33,6 +33,10 @@ void clearingAnalysisPreservesAttemptHistoryAndDoesNotReuseBossOrZhilianRowIds() "(request_key, platform, profile_id, job_key, job_row_id, state, requested_at, updated_at) " + "VALUES ('boss-old-attempt', 'boss', 1, 'boss-old', 10, 'UNKNOWN', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), " + "('zhilian-old-attempt', 'zhilian', 1, 'zhilian-old', 20, 'UNKNOWN', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)"); + jdbcTemplate.update("INSERT INTO job_analysis_task " + + "(profile_id, platform, status, task_key, job_key, job_row_id, request_json, created_at, updated_at) " + + "VALUES (1, 'boss', 'PENDING', 'boss-ai-old', 'boss-old', 10, '{}', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), " + + "(1, 'zhilian', 'FAILED', 'zhilian-ai-old', 'zhilian-old', 20, '{}', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)"); ProfileService profileService = mock(ProfileService.class); when(profileService.getCurrentProfileId()).thenReturn(1L); @@ -50,5 +54,37 @@ void clearingAnalysisPreservesAttemptHistoryAndDoesNotReuseBossOrZhilianRowIds() .isGreaterThan(20L); assertThat(jdbcTemplate.queryForObject("SELECT COUNT(*) FROM delivery_attempt", Integer.class)) .isEqualTo(2); + assertThat(jdbcTemplate.queryForObject("SELECT COUNT(*) FROM job_analysis_task WHERE task_key IS NOT NULL", Integer.class)) + .isZero(); + } + + @Test + void clearingAnalysisIsBlockedWhileAiTaskLeaseIsActive() { + DriverManagerDataSource dataSource = new DriverManagerDataSource( + "jdbc:sqlite:" + tempDir.resolve("clear-active-lease.db").toAbsolutePath()); + Flyway.configure() + .dataSource(dataSource) + .locations("classpath:db/migration") + .load() + .migrate(); + JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource); + jdbcTemplate.update("INSERT INTO profile(id, name, is_active) VALUES (1, 'profile', 1)"); + jdbcTemplate.update("INSERT INTO boss_data(id, profile_id, encrypt_id, delivery_status) " + + "VALUES (10, 1, 'boss-active', 'AI分析中')"); + jdbcTemplate.update("INSERT INTO job_analysis_task " + + "(profile_id, platform, status, task_key, job_key, job_row_id, request_json, " + + "lease_owner, lease_expires_at, created_at, updated_at) " + + "VALUES (1, 'boss', 'LEASED', 'boss-ai-active', 'boss-active', 10, '{}', " + + "'lease', '2099-01-01 00:00:00', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)"); + + ProfileService profileService = mock(ProfileService.class); + when(profileService.getCurrentProfileId()).thenReturn(1L); + BossService bossService = new BossService(null, null, null, null, null, dataSource, profileService); + + assertThat(bossService.clearBossAnalysisData()) + .containsEntry("success", false); + assertThat(jdbcTemplate.queryForObject("SELECT COUNT(*) FROM boss_data", Integer.class)).isEqualTo(1); + assertThat(jdbcTemplate.queryForObject("SELECT COUNT(*) FROM job_analysis_task WHERE status='LEASED'", Integer.class)) + .isEqualTo(1); } } diff --git a/src/test/java/com/getjobs/application/service/ChromeJobAnalysisQueueServiceTest.java b/src/test/java/com/getjobs/application/service/ChromeJobAnalysisQueueServiceTest.java new file mode 100644 index 0000000..69ef022 --- /dev/null +++ b/src/test/java/com/getjobs/application/service/ChromeJobAnalysisQueueServiceTest.java @@ -0,0 +1,251 @@ +package com.getjobs.application.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import org.flywaydb.core.Flyway; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.datasource.DataSourceTransactionManager; +import org.springframework.jdbc.datasource.DriverManagerDataSource; + +import java.nio.file.Path; +import java.time.Duration; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.doCallRealMethod; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.timeout; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.mockito.Mockito.spy; + +class ChromeJobAnalysisQueueServiceTest { + @TempDir + Path tempDir; + + private JdbcTemplate jdbcTemplate; + private JobAnalysisTaskStore store; + private JobAiAnalysisService analysisService; + private ChromeJobAnalysisQueueService queue; + + @BeforeEach + void setUp() { + DriverManagerDataSource dataSource = new DriverManagerDataSource( + "jdbc:sqlite:" + tempDir.resolve("queue.db").toAbsolutePath()); + Flyway.configure() + .dataSource(dataSource) + .locations("classpath:db/migration") + .load() + .migrate(); + jdbcTemplate = new JdbcTemplate(dataSource); + store = new JobAnalysisTaskStore( + jdbcTemplate, + new DataSourceTransactionManager(dataSource), + new ObjectMapper() + ); + store.validateSchema(); + analysisService = mock(JobAiAnalysisService.class); + } + + @AfterEach + void tearDown() { + if (queue != null) queue.shutdown(); + } + + @Test + void duplicateEnqueueInvokesProviderOnlyOnce() { + when(analysisService.analyzeJob(any(), any(), any())).thenReturn(successResult()); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + ChromeJobAnalysisQueueService.AnalysisJob job = job(request("boss", "job-duplicate", "run-a")); + + ChromeJobAnalysisQueueService.EnqueueResult first = queue.enqueue(job); + ChromeJobAnalysisQueueService.EnqueueResult duplicate = queue.enqueue( + job(request("boss", "job-duplicate", "run-b"))); + + assertThat(first.isQueued()).isTrue(); + assertThat(duplicate.isQueued()).isFalse(); + verify(analysisService, timeout(3000).times(1)).analyzeJob(any(), any(), any()); + awaitStatus(firstTaskId(), "SUCCEEDED"); + } + + @Test + void startupDispatchesPersistedPendingTask() { + long taskId = store.submit(request("boss", "job-restart", "run-before-restart")).task().id(); + when(analysisService.analyzeJob(any(), any(), any())).thenReturn(successResult()); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + queue.initialize(); + + verify(analysisService, timeout(3000).times(1)).analyzeJob(any(), any(), any()); + awaitStatus(taskId, "SUCCEEDED"); + assertThat(store.findById(taskId).attemptCount()).isEqualTo(1); + } + + @Test + void expiredLeaseWithPersistedPlatformResultIsReconciledWithoutProviderCall() { + long taskId = leaseExpiredTask("boss", "job-reconciled"); + when(analysisService.inspectPlatformAnalysis(any())) + .thenReturn(new JobAiAnalysisService.PlatformAnalysisState(true, false, DeliveryStatus.WAITING_CONFIRM)); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + queue.reconcileExpiredLeases(); + + assertThat(store.findById(taskId).status()).isEqualTo("SUCCEEDED"); + verify(analysisService, never()).analyzeJob(any(), any(), any()); + verify(analysisService, never()).markAnalysisInterrupted(any(), any()); + } + + @Test + void expiredUnresolvedLeaseBecomesUnknownAndDoesNotRetryProvider() { + long taskId = leaseExpiredTask("zhilian", "job-unknown"); + when(analysisService.inspectPlatformAnalysis(any())) + .thenReturn(JobAiAnalysisService.PlatformAnalysisState.incomplete(DeliveryStatus.AI_ANALYZING)); + when(analysisService.markAnalysisInterrupted(any(), any())).thenReturn(true); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + queue.reconcileExpiredLeases(); + + assertThat(store.findById(taskId).status()).isEqualTo("UNKNOWN"); + verify(analysisService, never()).analyzeJob(any(), any(), any()); + verify(analysisService).markAnalysisInterrupted(any(), any()); + } + + @Test + void startupRegistersLegacyAnalyzingRowAsUnknownWithoutProviderCall() { + jdbcTemplate.update("INSERT INTO profile(id, name, is_active) VALUES (1, 'profile', 1)"); + jdbcTemplate.update("INSERT INTO boss_data(id, profile_id, encrypt_id, company_name, job_name, " + + "delivery_status, job_description, scan_run_id) VALUES " + + "(30, 1, 'legacy-analyzing', '测试公司', 'Java 工程师', ?, '岗位描述', 'legacy-run')", + DeliveryStatus.AI_ANALYZING); + when(analysisService.markAnalysisInterrupted(any(), any())).thenReturn(true); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + queue.initialize(); + queue.recoverOrphanedAnalyzingTasks(); + + assertThat(jdbcTemplate.queryForObject( + "SELECT status FROM job_analysis_task WHERE platform='boss' AND job_row_id=30", + String.class)).isEqualTo("UNKNOWN"); + verify(analysisService, never()).analyzeJob(any(), any(), any()); + verify(analysisService, times(1)).markAnalysisInterrupted(any(), any()); + } + + @Test + void unexpectedExecutionFailureWritesExplicitTaskAndPlatformFailure() { + when(analysisService.analyzeJob(any(), any(), any())).thenThrow(new IllegalStateException("executor failed")); + when(analysisService.inspectPlatformAnalysis(any())) + .thenReturn(JobAiAnalysisService.PlatformAnalysisState.incomplete(DeliveryStatus.AI_ANALYZING)); + when(analysisService.markAnalysisInterrupted(any(), any())).thenReturn(true); + queue = new ChromeJobAnalysisQueueService(analysisService, store); + + ChromeJobAnalysisQueueService.EnqueueResult submitted = queue.enqueue( + job(request("boss", "job-exception", "run-a"))); + + awaitStatus(submittedTaskId("job-exception"), "FAILED"); + verify(analysisService).markAnalysisInterrupted(any(), any()); + } + + @Test + void completionWriteExceptionReconcilesPersistedPlatformSuccess() { + JobAnalysisTaskStore flakyStore = spy(store); + doThrow(new IllegalStateException("first completion write failed")) + .doCallRealMethod() + .when(flakyStore) + .complete(anyLong(), anyString(), anyBoolean(), anyString()); + when(analysisService.analyzeJob(any(), any(), any())).thenReturn(successResult()); + when(analysisService.inspectPlatformAnalysis(any())) + .thenReturn(new JobAiAnalysisService.PlatformAnalysisState( + true, false, DeliveryStatus.WAITING_CONFIRM)); + queue = new ChromeJobAnalysisQueueService(analysisService, flakyStore); + + queue.enqueue(job(request("boss", "job-completion-recovery", "run-a"))); + + awaitStatus(submittedTaskId("job-completion-recovery"), "SUCCEEDED"); + verify(analysisService, timeout(3000).times(1)).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(); + jdbcTemplate.update( + "UPDATE job_analysis_task SET lease_expires_at='2000-01-01 00:00:00' WHERE id=?", + taskId + ); + return taskId; + } + + private long firstTaskId() { + return jdbcTemplate.queryForObject( + "SELECT id FROM job_analysis_task WHERE task_key IS NOT NULL ORDER BY id LIMIT 1", + Long.class + ); + } + + private long submittedTaskId(String jobKey) { + return jdbcTemplate.queryForObject( + "SELECT id FROM job_analysis_task WHERE job_key=? ORDER BY id DESC LIMIT 1", + Long.class, + jobKey + ); + } + + private void awaitStatus(long taskId, String expected) { + long deadline = System.currentTimeMillis() + 3000; + while (System.currentTimeMillis() < deadline) { + if (expected.equals(store.findById(taskId).status())) return; + try { + Thread.sleep(20); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError("等待任务状态时被中断", e); + } + } + assertThat(store.findById(taskId).status()).isEqualTo(expected); + } + + private ChromeJobAnalysisQueueService.AnalysisJob job(JobAiAnalysisService.JobAnalysisRequest request) { + ChromeJobAnalysisQueueService.AnalysisJob job = new ChromeJobAnalysisQueueService.AnalysisJob(); + job.setRunId(request.getScanRunId()); + job.setCurrentStatus(DeliveryStatus.NOT_DELIVERED); + job.setCurrent(1); + job.setTotal(1); + job.setRequest(request); + return job; + } + + private JobAiAnalysisService.JobAnalysisRequest request(String platform, String jobKey, String runId) { + JobAiAnalysisService.JobAnalysisRequest request = new JobAiAnalysisService.JobAnalysisRequest(); + request.setProfileId(1L); + request.setPlatform(platform); + request.setJobKey(jobKey); + request.setJobRowId("boss".equals(platform) ? 10L : 20L); + request.setKeyword("Java"); + request.setCompanyName("测试公司"); + request.setJobName("Java 工程师"); + request.setSalary("20-30K"); + request.setLocation("深圳"); + request.setExperience("3-5年"); + request.setDegree("本科"); + request.setCompanyInfo("互联网"); + request.setJobDescription("负责 Spring Boot 服务开发"); + request.setScanRunId(runId); + return request; + } + + private JobAiAnalysisService.AnalysisResult successResult() { + JobAiAnalysisService.AnalysisResult result = new JobAiAnalysisService.AnalysisResult(); + result.setScore(90); + result.setDecision("APPLY"); + result.setSummary("匹配"); + return result; + } +} diff --git a/src/test/java/com/getjobs/application/service/DatabaseMigrationRehearsalTest.java b/src/test/java/com/getjobs/application/service/DatabaseMigrationRehearsalTest.java index 2aa66c3..f1112a0 100644 --- a/src/test/java/com/getjobs/application/service/DatabaseMigrationRehearsalTest.java +++ b/src/test/java/com/getjobs/application/service/DatabaseMigrationRehearsalTest.java @@ -58,7 +58,7 @@ void migratesIsolatedCopyWithoutChangingSourceDatabase() throws Exception { DatabaseSchemaService.validateSchema(connection); assertThat(scalarText(connection, "PRAGMA integrity_check")).isEqualTo("ok"); assertThat(scalarLong(connection, - "SELECT COUNT(*) FROM flyway_schema_history WHERE success=1 AND version='6'")) + "SELECT COUNT(*) FROM flyway_schema_history WHERE success=1 AND version='7'")) .isEqualTo(1L); assertThat(scalarLong(connection, "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='delivery_attempt'")) diff --git a/src/test/java/com/getjobs/application/service/DatabaseMigrationTest.java b/src/test/java/com/getjobs/application/service/DatabaseMigrationTest.java index d399716..cecb0ab 100644 --- a/src/test/java/com/getjobs/application/service/DatabaseMigrationTest.java +++ b/src/test/java/com/getjobs/application/service/DatabaseMigrationTest.java @@ -18,7 +18,7 @@ class DatabaseMigrationTest { Path tempDir; @Test - void freshDatabaseMigratesThroughV6AndMatchesSchemaContract() throws Exception { + void freshDatabaseMigratesThroughV7AndMatchesSchemaContract() throws Exception { String url = sqliteUrl(tempDir.resolve("fresh.db")); Flyway flyway = flyway(url); @@ -27,7 +27,7 @@ void freshDatabaseMigratesThroughV6AndMatchesSchemaContract() throws Exception { try (Connection connection = DriverManager.getConnection(url)) { DatabaseSchemaService.validateSchema(connection); assertThat(scalar(connection, - "SELECT COUNT(*) FROM flyway_schema_history WHERE success=1 AND version='6'")) + "SELECT COUNT(*) FROM flyway_schema_history WHERE success=1 AND version='7'")) .isEqualTo(1L); assertThat(columns(connection, "ai")).contains("apply_threshold", "priority_apply_threshold"); assertThat(columns(connection, "boss_data")) @@ -35,6 +35,35 @@ void freshDatabaseMigratesThroughV6AndMatchesSchemaContract() throws Exception { assertThat(columns(connection, "liepin_data")).contains("delivery_status"); assertThat(columns(connection, "job51_data")).contains("delivery_status"); assertThat(tableExists(connection, "delivery_attempt")).isTrue(); + assertThat(columns(connection, "job_analysis_task")) + .contains("task_key", "job_key", "job_row_id", "request_json", "attempt_count", + "lease_owner", "lease_expires_at", "last_error", "started_at", "completed_at"); + } + } + + @Test + void v7PreservesLegacyAggregateRowsAndLeavesThemUndispatchable() throws Exception { + String url = sqliteUrl(tempDir.resolve("legacy-ai-task.db")); + Flyway beforeV7 = Flyway.configure() + .dataSource(url, null, null) + .locations("classpath:db/migration") + .target("6") + .load(); + beforeV7.migrate(); + try (Connection connection = DriverManager.getConnection(url); Statement statement = connection.createStatement()) { + statement.execute("INSERT INTO job_analysis_task(platform, scan_run_id, status, total_count, created_at) " + + "VALUES ('boss', 'legacy-run', 'SUCCEEDED', 12, CURRENT_TIMESTAMP)"); + } + + flyway(url).migrate(); + + try (Connection connection = DriverManager.getConnection(url)) { + assertThat(scalar(connection, "SELECT COUNT(*) FROM job_analysis_task WHERE scan_run_id='legacy-run'")) + .isEqualTo(1L); + assertThat(scalar(connection, "SELECT COUNT(*) FROM job_analysis_task WHERE task_key IS NOT NULL")) + .isZero(); + assertThat(columns(connection, "job_analysis_task")) + .contains("task_key", "request_json", "lease_expires_at"); } } diff --git a/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java b/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java index 406940f..f527a68 100644 --- a/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java +++ b/src/test/java/com/getjobs/application/service/JobAiAnalysisServiceStatusTest.java @@ -21,11 +21,14 @@ import java.nio.charset.StandardCharsets; import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -109,6 +112,98 @@ void deliveredJobKeepsStatusWhenAnalyzedAgain() { assertThat(lastBossUpdate().getDeliveryStatus()).isNull(); } + @Test + void restartInspectionRecognizesPersistedPlatformResult() { + BossJobDataEntity completed = bossJob(DeliveryStatus.WAITING_CONFIRM); + completed.setAiDecision("APPLY"); + when(bossJobDataMapper.selectOne(any())).thenReturn(completed); + + JobAiAnalysisService.PlatformAnalysisState state = service.inspectPlatformAnalysis(bossRequest()); + + assertThat(state.completed()).isTrue(); + assertThat(state.failed()).isFalse(); + assertThat(state.status()).isEqualTo(DeliveryStatus.WAITING_CONFIRM); + } + + @Test + void restartInspectionIgnoresStaleDecisionWhileTaskIsStillAnalyzing() { + BossJobDataEntity analyzing = bossJob(DeliveryStatus.AI_ANALYZING); + analyzing.setAiDecision(DeliveryStatus.AI_ANALYSIS_FAILED); + when(bossJobDataMapper.selectOne(any())).thenReturn(analyzing); + + JobAiAnalysisService.PlatformAnalysisState state = service.inspectPlatformAnalysis(bossRequest()); + + assertThat(state.completed()).isFalse(); + assertThat(state.failed()).isFalse(); + } + + @Test + void restartInspectionDoesNotTreatOldDecisionAsCompletedWithoutResultStatus() { + BossJobDataEntity pending = bossJob(DeliveryStatus.NOT_DELIVERED); + pending.setAiDecision(DeliveryStatus.AI_ANALYSIS_FAILED); + when(bossJobDataMapper.selectOne(any())).thenReturn(pending); + + JobAiAnalysisService.PlatformAnalysisState state = service.inspectPlatformAnalysis(bossRequest()); + + assertThat(state.completed()).isFalse(); + assertThat(state.failed()).isFalse(); + } + + @Test + void interruptedRecoveryOnlyWritesExplicitAiFailure() { + when(zhilianJobDataMapper.update(any(), any(UpdateWrapper.class))).thenReturn(1); + + boolean changed = service.markAnalysisInterrupted(zhilianRequest(), "租约过期,结果未知"); + + assertThat(changed).isTrue(); + ZhilianJobDataEntity update = lastZhilianUpdate(); + assertThat(update.getDeliveryStatus()).isEqualTo(DeliveryStatus.AI_ANALYSIS_FAILED); + assertThat(update.getAiDecision()).isEqualTo(DeliveryStatus.AI_ANALYSIS_FAILED); + assertThat(update.getAiReason()).contains("结果未知"); + } + + @Test + void durableTaskDoesNotCallProviderWhenExactJobCannotBeReserved() { + JobAiAnalysisService.JobAnalysisRequest request = bossRequest(); + request.setJobRowId(99L); + when(bossJobDataMapper.update(any(), any(UpdateWrapper.class))).thenReturn(0); + + JobAiAnalysisService.AnalysisResult result = service.analyzeJob(request); + + assertThat(result.isFailure()).isTrue(); + assertThat(result.getSummary()).contains("未调用 AI Provider"); + verify(aiService, never()).sendRequest(any()); + } + + @Test + void leaseTransactionRejectsLateProviderResultBeforeAnyResultWrite() { + JobAiAnalysisService.JobAnalysisRequest request = bossRequest(); + request.setJobRowId(99L); + when(bossJobDataMapper.update(any(), any(UpdateWrapper.class))).thenReturn(1); + when(resumeProfileMapper.selectOne(any())).thenReturn(resume()); + when(aiService.sendRequest(any())).thenReturn(""" + {"score":90,"decision":"APPLY","summary":"旧租约结果","strengths":[],"risks":[],"greeting":"你好"} + """); + AtomicInteger guardedWrites = new AtomicInteger(); + + JobAiAnalysisService.AnalysisResult result = service.analyzeJob( + request, + () -> true, + action -> { + if (guardedWrites.incrementAndGet() == 1) { + action.run(); + return true; + } + return false; + } + ); + + assertThat(result.isStaleLease()).isTrue(); + verify(aiService).sendRequest(any()); + verify(jobAiAnalysisMapper, never()).insert(any(com.getjobs.application.entity.JobAiAnalysisEntity.class)); + verify(bossJobDataMapper, times(1)).update(any(), any(UpdateWrapper.class)); + } + @Test void manualZhilianAnalyzeApplyEndsWaitingConfirm() { when(zhilianJobDataMapper.selectOne(any())).thenReturn(zhilianJob(DeliveryStatus.NOT_DELIVERED)); diff --git a/src/test/java/com/getjobs/application/service/JobAnalysisTaskStoreTest.java b/src/test/java/com/getjobs/application/service/JobAnalysisTaskStoreTest.java new file mode 100644 index 0000000..4e9cfcf --- /dev/null +++ b/src/test/java/com/getjobs/application/service/JobAnalysisTaskStoreTest.java @@ -0,0 +1,157 @@ +package com.getjobs.application.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import org.flywaydb.core.Flyway; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.datasource.DataSourceTransactionManager; +import org.springframework.jdbc.datasource.DriverManagerDataSource; + +import java.nio.file.Path; +import java.time.Duration; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.assertj.core.api.Assertions.assertThat; + +class JobAnalysisTaskStoreTest { + @TempDir + Path tempDir; + + private JdbcTemplate jdbcTemplate; + private JobAnalysisTaskStore store; + + @BeforeEach + void setUp() { + DriverManagerDataSource dataSource = new DriverManagerDataSource( + "jdbc:sqlite:" + tempDir.resolve("analysis-task.db").toAbsolutePath()); + Flyway.configure() + .dataSource(dataSource) + .locations("classpath:db/migration") + .load() + .migrate(); + jdbcTemplate = new JdbcTemplate(dataSource); + store = new JobAnalysisTaskStore( + jdbcTemplate, + new DataSourceTransactionManager(dataSource), + new ObjectMapper() + ); + store.validateSchema(); + } + + @Test + void stableTaskKeyDeduplicatesAcrossRunIdsButSeparatesProfiles() { + JobAnalysisTaskStore.SubmitResult first = store.submit(request(1L, "boss", "job-1", "run-a")); + JobAnalysisTaskStore.SubmitResult duplicate = store.submit(request(1L, "boss", "job-1", "run-b")); + JobAnalysisTaskStore.SubmitResult otherProfile = store.submit(request(2L, "boss", "job-1", "run-b")); + + assertThat(first.created()).isTrue(); + assertThat(duplicate.created()).isFalse(); + assertThat(duplicate.task().id()).isEqualTo(first.task().id()); + assertThat(otherProfile.created()).isTrue(); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM job_analysis_task WHERE task_key IS NOT NULL", Integer.class)).isEqualTo(2); + } + + @Test + void concurrentConsumersCanOnlyClaimOnce() throws Exception { + long taskId = store.submit(request(1L, "boss", "job-claim", "run-a")).task().id(); + ExecutorService executor = Executors.newFixedThreadPool(2); + CountDownLatch ready = new CountDownLatch(2); + CountDownLatch start = new CountDownLatch(1); + try { + Future first = executor.submit(() -> claimAfterBarrier(taskId, "lease-a", ready, start)); + Future second = executor.submit(() -> claimAfterBarrier(taskId, "lease-b", ready, start)); + ready.await(); + start.countDown(); + + assertThat(List.of(first.get(), second.get())).containsExactlyInAnyOrder(true, false); + assertThat(store.findById(taskId).attemptCount()).isEqualTo(1); + } finally { + executor.shutdownNow(); + } + } + + @Test + void failedAndUnknownTasksRequireExplicitRetryAndIncrementAttemptsOnClaim() { + long failedId = store.submit(request(1L, "boss", "job-failed", "run-a")).task().id(); + assertThat(store.claim(failedId, "lease-failed", Duration.ofMinutes(1))).isNotNull(); + assertThat(store.renewLease(failedId, "wrong-lease", Duration.ofMinutes(1))).isFalse(); + assertThat(store.renewLease(failedId, "lease-failed", Duration.ofMinutes(1))).isTrue(); + assertThat(store.complete(failedId, "lease-failed", true, "provider failed")).isTrue(); + assertThat(store.retry(failedId, 1L).accepted()).isTrue(); + assertThat(store.claim(failedId, "lease-retry", Duration.ofMinutes(1))).isNotNull(); + assertThat(store.findById(failedId).attemptCount()).isEqualTo(2); + assertThat(store.complete(failedId, "lease-retry", false, "ok")).isTrue(); + assertThat(store.retry(failedId, 1L).accepted()).isFalse(); + + long unknownId = store.submit(request(1L, "zhilian", "job-unknown", "run-a")).task().id(); + assertThat(store.claim(unknownId, "lease-unknown", Duration.ofMinutes(1))).isNotNull(); + jdbcTemplate.update("UPDATE job_analysis_task SET lease_expires_at='2000-01-01 00:00:00' WHERE id=?", unknownId); + assertThat(store.reconcileExpired( + unknownId, + "lease-unknown", + JobAnalysisTaskStore.Status.UNKNOWN, + "result unknown" + )).isTrue(); + assertThat(store.retry(unknownId, 2L, true).accepted()).isFalse(); + assertThat(store.retry(unknownId, 1L).accepted()).isFalse(); + assertThat(store.retry(unknownId, 1L, true).accepted()).isTrue(); + assertThat(store.findById(unknownId).status()).isEqualTo("PENDING"); + } + + @Test + void leaseHeartbeatStopsAtHardExecutionLimitAndTaskCanExpire() { + long taskId = store.submit(request(1L, "boss", "job-hard-timeout", "run-a")).task().id(); + assertThat(store.claim(taskId, "lease-hard-timeout", Duration.ofMinutes(5))).isNotNull(); + jdbcTemplate.update("UPDATE job_analysis_task SET started_at='2000-01-01 00:00:00.000', " + + "lease_expires_at='2000-01-01 00:00:01.000' WHERE id=?", + taskId); + AtomicBoolean writeExecuted = new AtomicBoolean(); + + assertThat(store.renewLease(taskId, "lease-hard-timeout", Duration.ofMinutes(5))).isFalse(); + assertThat(store.isLeaseOwner(taskId, "lease-hard-timeout")).isFalse(); + assertThat(store.executeWithLease(taskId, "lease-hard-timeout", () -> writeExecuted.set(true))).isFalse(); + assertThat(writeExecuted).isFalse(); + assertThat(store.complete(taskId, "lease-hard-timeout", false, "late result")).isFalse(); + assertThat(store.listExpiredLeases(10)).extracting(JobAnalysisTaskStore.TaskRecord::id) + .contains(taskId); + } + + private boolean claimAfterBarrier(long taskId, + String lease, + CountDownLatch ready, + CountDownLatch start) throws Exception { + ready.countDown(); + start.await(); + return store.claim(taskId, lease, Duration.ofMinutes(1)) != null; + } + + private JobAiAnalysisService.JobAnalysisRequest request(long profileId, + String platform, + String jobKey, + String runId) { + JobAiAnalysisService.JobAnalysisRequest request = new JobAiAnalysisService.JobAnalysisRequest(); + request.setProfileId(profileId); + request.setPlatform(platform); + request.setJobKey(jobKey); + request.setJobRowId("boss".equals(platform) ? 10L : 20L); + request.setKeyword("Java"); + request.setCompanyName("测试公司"); + request.setJobName("Java 工程师"); + request.setSalary("20-30K"); + request.setLocation("深圳"); + request.setExperience("3-5年"); + request.setDegree("本科"); + request.setCompanyInfo("互联网"); + request.setJobDescription("负责 Spring Boot 服务开发"); + request.setScanRunId(runId); + return request; + } +} diff --git a/src/test/java/com/getjobs/application/service/ProfileServiceAiTaskSafetyTest.java b/src/test/java/com/getjobs/application/service/ProfileServiceAiTaskSafetyTest.java new file mode 100644 index 0000000..28a960c --- /dev/null +++ b/src/test/java/com/getjobs/application/service/ProfileServiceAiTaskSafetyTest.java @@ -0,0 +1,69 @@ +package com.getjobs.application.service; + +import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; +import com.getjobs.application.entity.ProfileEntity; +import com.getjobs.application.mapper.ProfileMapper; +import org.flywaydb.core.Flyway; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.datasource.DataSourceTransactionManager; +import org.springframework.jdbc.datasource.DriverManagerDataSource; + +import java.nio.file.Path; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class ProfileServiceAiTaskSafetyTest { + @TempDir + Path tempDir; + + @Test + void forceDeleteRollsBackWhenAiTaskLeaseIsActive() { + DriverManagerDataSource dataSource = new DriverManagerDataSource( + "jdbc:sqlite:" + tempDir.resolve("profile-delete.db").toAbsolutePath()); + Flyway.configure() + .dataSource(dataSource) + .locations("classpath:db/migration") + .load() + .migrate(); + JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource); + jdbcTemplate.update("INSERT INTO profile(id, name, is_active) VALUES (1, 'active', 1), (2, 'other', 0)"); + jdbcTemplate.update("INSERT INTO boss_data(id, profile_id, encrypt_id, delivery_status) " + + "VALUES (10, 1, 'boss-active', 'AI分析中')"); + jdbcTemplate.update("INSERT INTO job_analysis_task " + + "(profile_id, platform, status, task_key, job_key, job_row_id, request_json, " + + "lease_owner, lease_expires_at, created_at, updated_at) " + + "VALUES (1, 'boss', 'LEASED', 'profile-delete-active', 'boss-active', 10, '{}', " + + "'lease', '2099-01-01 00:00:00', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)"); + + ProfileEntity active = new ProfileEntity(); + active.setId(1L); + active.setName("active"); + active.setIsActive(1); + ProfileMapper profileMapper = mock(ProfileMapper.class); + when(profileMapper.selectById(1L)).thenReturn(active); + when(profileMapper.selectCount(isNull())).thenReturn(2L); + when(profileMapper.selectOne(any(QueryWrapper.class))).thenReturn(active); + ProfileService service = new ProfileService( + profileMapper, + jdbcTemplate, + new DataSourceTransactionManager(dataSource) + ); + + ProfileService.DeleteProfileResult result = service.deleteProfile(1L, true); + + assertThat(result).isNotNull(); + assertThat(result.success()).isFalse(); + assertThat(result.message()).contains("AI 分析正在执行"); + assertThat(jdbcTemplate.queryForObject("SELECT COUNT(*) FROM boss_data WHERE profile_id=1", Integer.class)) + .isEqualTo(1); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM job_analysis_task WHERE profile_id=1 AND status='LEASED'", Integer.class)) + .isEqualTo(1); + } +} diff --git a/tasks/2026-08-24-p1-1-ai-task-recovery.md b/tasks/2026-08-24-p1-1-ai-task-recovery.md new file mode 100644 index 0000000..0f5354c --- /dev/null +++ b/tasks/2026-08-24-p1-1-ai-task-recovery.md @@ -0,0 +1,68 @@ +# P1.1 AI 任务持久化与重启恢复 + +## 背景 + +当前 `ChromeJobAnalysisQueueService` 只使用进程内线程池、`activeKeys` 和 30 分钟 `completedKeys`。Controller 会先把岗位写成 `AI分析中`,再入内存队列;进程在排队、Provider 请求或结果写回之间退出时,任务与去重信息丢失,岗位可能永久残留在 `AI分析中`。V4 已创建 `job_analysis_task`,但主链尚未使用。 + +## 目标 + +1. 将 Boss/智联 Chrome AI 队列写入 SQLite,线程池只作为持久任务的执行器。 +2. 使用稳定 `task_key` 跨 runId、跨重启去重;重复入队不重复调用 Provider。 +3. 使用 `PENDING / LEASED / SUCCEEDED / FAILED / UNKNOWN` 和租约 CAS。 +4. 重启后自动恢复从未开始的 `PENDING`;过期 `LEASED` 只进入 `UNKNOWN`,不得自动重发计费请求。 +5. 若 Provider/结果已完成但任务终态尚未写入,依据平台读模型完成对账,不覆盖已有结果。 +6. 提供当前 Profile 下的只读任务查询和 `UNKNOWN/FAILED` 显式重试;重试必须由用户主动调用。 + +## 允许修改范围 + +- Flyway V7:演进现有 `job_analysis_task`,增加逐岗位任务、租约、重试和错误字段/索引。 +- 新增任务存储服务;改造 `ChromeJobAnalysisQueueService` 使用持久任务。 +- `JobAiAnalysisService` 增加只读完成判断和保守的中断状态写回。 +- AI Controller 增加当前 Profile 的任务查询/显式重试 API。 +- 对应 migration、存储、队列、重启恢复和 Controller 隔离测试。 + +## 禁止修改范围 + +- 不调用真实 AI Provider,不启动服务,不访问或迁移 `db/getjobs.db`。 +- 不修改 Provider 协议、Prompt、总超时、429/500 退避或计费策略;这些属于 P1.2。 +- 不持久化 SSE callback,不整体重写 Controller/Worker,不处理旧 Playwright 同步 AI 路径。 +- 不自动重试 `LEASED/UNKNOWN`,不声称外部 Provider 可实现 exactly-once。 + +## 已确定实现要求 + +- `task_key` 不包含 runId;至少绑定 profile、platform、jobKey 与岗位分析输入摘要。 +- enqueue 的“持久化成功”和 executor 的“已开始执行”分开;executor 满时任务保留 `PENDING`。 +- claim、lease、完成、失败、UNKNOWN 均使用条件 UPDATE;旧 lease 或旧 worker 不能覆盖新 attempt。 +- 执行期间用独立心跳续租,避免合法长耗时 Provider 被误判过期;单次执行最长续租 65 分钟,防止无请求超时的 Provider 永久占用任务租约。 +- Provider 返回后、写入分析结果前必须再次校验当前 lease token;旧 worker 丢失租约后只丢弃结果,不能覆盖显式重试后的新 attempt。 +- `PENDING` 可安全自动恢复;`LEASED` 过期后先查看平台状态:结果已经落库则完成任务,否则在同一事务中标记 `UNKNOWN` 并把残留 `AI分析中` 变为明确失败。 +- 升级前遗留且没有持久任务的 `AI分析中` 岗位仅在服务启动、尚未接收同步请求时分批登记为 `UNKNOWN`,不得直接补发 Provider 请求。 +- `FAILED/UNKNOWN` 显式重试复用同一业务 task,但 attempt 递增、清除旧 lease,并重新进入 `PENDING`。 +- `UNKNOWN` 重试必须额外传入 `confirmUnknown=true`,明确承担可能重复计费的风险。 +- 旧 V4 聚合行没有 `task_key/request_json`,保留但不调度、不删除。 + +## 验收标准 + +- fresh/legacy SQLite 均能迁移到 V7;V4 旧行保留。 +- 重复 enqueue 只生成一条可执行任务;重启恢复 PENDING 不重复建任务。 +- 两个消费者并发 claim 时只有一个成功。 +- 过期 LEASED 在平台已有最终结果时对账完成;仍为 `AI分析中` 时进入 UNKNOWN,Provider 调用次数为 0。 +- 遗留无任务的 `AI分析中` 行会在启动阶段分批进入 UNKNOWN,Provider 调用次数为 0;定时维护不得误伤同步 `/analyze-job`;合法长耗时任务通过心跳保持租约,但不能超过硬上限。 +- 旧租约 Provider 延迟返回时不会写入分析表或覆盖平台状态;超过硬上限的任务可自然过期并进入 UNKNOWN 对账。 +- FAILED/UNKNOWN 只有显式 retry 才重新进入 PENDING;SUCCEEDED/LEASED 不允许重试。 +- 完整后端测试与隔离迁移演练通过;真实数据库 hash/mtime/sidecar 不变。 + +## 测试命令 + +```powershell +.\gradlew.bat clean test --console=plain +$env:P0_REHEARSAL_DB = (Resolve-Path db/getjobs.db).Path +.\gradlew.bat test --tests com.getjobs.application.service.DatabaseMigrationRehearsalTest +Remove-Item Env:P0_REHEARSAL_DB +``` + +## 返回格式 + +- 状态机、租约、重启恢复和重复计费边界说明。 +- migration 与并发/崩溃点测试证据。 +- 原数据库未修改证据、diff、Commit、Push 与堆叠 PR。