From 020bc46046daef523d610b1c2a4cb60c33284bb8 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:15:17 -0800 Subject: [PATCH 01/19] fix(ingestion): reindex chunks on metadata drift Chunk deduplication only checked hash markers derived from URL, chunk index, and chunk text. When title or package metadata changed without content changes, ingestion skipped those chunks and left stale metadata in the vector layer. This change persists title/package metadata in hash marker payloads and compares stored marker metadata during dedup checks. Existing hashes are now reprocessed when metadata changes, while unchanged hashes still skip as before. - Add metadata-aware hash marker read/write logic in LocalStoreService - Add hasHashMetadataChanged(...) and metadata parsing for existing markers - Update ChunkProcessingService to reprocess hash hits on metadata drift - Pass title/package metadata when marking ingested hashes in ingestion flows - Add regression tests for both metadata-drift reingest and unchanged skip --- .../service/ChunkProcessingService.java | 24 +++- .../service/DocsIngestionService.java | 12 +- .../javachat/service/LocalStoreService.java | 125 +++++++++++++++++- .../LocalDocsFileIngestionProcessor.java | 12 +- .../service/ChunkProcessingServiceTest.java | 75 +++++++++++ 5 files changed, 237 insertions(+), 11 deletions(-) create mode 100644 src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java diff --git a/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java b/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java index c94a6d9c..c3634593 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java +++ b/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java @@ -38,7 +38,17 @@ public ChunkProcessingService( this.chunker = Objects.requireNonNull(chunker, "chunker"); this.hasher = Objects.requireNonNull(hasher, "hasher"); this.documentFactory = Objects.requireNonNull(documentFactory, "documentFactory"); - this.hashIngestionLookup = requiredLocalStore::isHashIngested; + this.hashIngestionLookup = new HashIngestionLookup() { + @Override + public boolean isHashIngested(String hash) { + return requiredLocalStore.isHashIngested(hash); + } + + @Override + public boolean hasMetadataChanged(String hash, String title, String packageName) { + return requiredLocalStore.hasHashMetadataChanged(hash, title, packageName); + } + }; this.chunkTextStore = requiredLocalStore::saveChunkText; this.pdfExtractor = Objects.requireNonNull(pdfExtractor, "pdfExtractor"); } @@ -71,7 +81,8 @@ public ChunkProcessingOutcome processAndStoreChunks(String text, String url, Str allChunkHashes.add(hash); // Skip if already processed (deduplication) - if (hashIngestionLookup.isHashIngested(hash)) { + boolean hashAlreadyIngested = hashIngestionLookup.isHashIngested(hash); + if (hashAlreadyIngested && !hashIngestionLookup.hasMetadataChanged(hash, title, packageName)) { skipped++; continue; } @@ -172,7 +183,8 @@ public ChunkProcessingOutcome processPdfAndStoreWithPages( totalChunks++; String hash = hasher.generateChunkHash(url, globalIndex, chunkText); allChunkHashes.add(hash); - if (!hashIngestionLookup.isHashIngested(hash)) { + boolean hashAlreadyIngested = hashIngestionLookup.isHashIngested(hash); + if (!hashAlreadyIngested || hashIngestionLookup.hasMetadataChanged(hash, title, packageName)) { Document doc = documentFactory.createDocumentWithPages( chunkText, url, title, globalIndex, packageName, hash, pageIndex + 1, pageIndex + 1); pageDocuments.add(doc); @@ -266,12 +278,16 @@ public boolean skippedAllChunks() { /** * Reads whether a chunk hash has already been indexed. */ - @FunctionalInterface private interface HashIngestionLookup { /** * Returns true when the chunk hash already has an ingest marker. */ boolean isHashIngested(String hash); + + /** + * Returns true when stored metadata for an ingested hash differs from current metadata. + */ + boolean hasMetadataChanged(String hash, String title, String packageName); } /** diff --git a/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java b/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java index d63a31f4..d9adeed6 100644 --- a/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java +++ b/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java @@ -157,14 +157,24 @@ private void markDocumentsIngested(List discoveredLinks = new java.util.ArrayList<>(); diff --git a/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java b/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java index 16f0276f..c7844557 100644 --- a/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java +++ b/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java @@ -10,6 +10,7 @@ import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; import java.util.ArrayList; +import java.util.Base64; import java.util.HexFormat; import java.util.List; import java.util.Optional; @@ -34,6 +35,9 @@ public class LocalStoreService { private static final String FILE_MARKER_EXTENSION = ".marker"; private static final String FILE_MARKER_HASH_PREFIX = "hash="; private static final String FILE_MARKER_FINGERPRINT_PREFIX = "fingerprint="; + private static final String HASH_MARKER_INGESTED_FLAG = "1"; + private static final String HASH_MARKER_TITLE_PREFIX = "titleB64="; + private static final String HASH_MARKER_PACKAGE_PREFIX = "packageB64="; private final String snapshotDirConfig; private final String parsedDirConfig; @@ -102,21 +106,73 @@ public void saveChunkText(String url, int index, String text, String hash) throw * Returns true when an ingest marker exists for the given chunk hash. */ public boolean isHashIngested(String hash) { - Path marker = indexDir.resolve(hash); - return Files.exists(marker); + return Files.exists(hashMarkerPath(hash)); + } + + /** + * Returns true when the stored metadata for an ingested chunk hash differs from the current metadata. + * + *

Older markers may not contain metadata fields. In that case this method returns true so callers + * can perform a one-time metadata refresh and backfill marker metadata.

+ * + * @param hash chunk hash marker key + * @param title current document title + * @param packageName current package name + * @return true when metadata has changed since the chunk was marked ingested + */ + public boolean hasHashMetadataChanged(String hash, String title, String packageName) { + Path markerPath = hashMarkerPath(hash); + if (!Files.exists(markerPath)) { + return false; + } + try { + HashMarkerMetadata storedMetadata = readHashMarkerMetadata(markerPath); + String normalizedTitle = normalizeHashMetadataText(title); + String normalizedPackageName = normalizeHashMetadataText(packageName); + return !normalizedTitle.equals(storedMetadata.title()) + || !normalizedPackageName.equals(storedMetadata.packageName()); + } catch (IOException markerReadFailure) { + throw new IllegalStateException("Failed to read hash ingestion marker for hash: " + hash, markerReadFailure); + } } /** * Writes an ingest marker for the given chunk hash when not already present. */ public void markHashIngested(String hash) throws IOException { - Path marker = indexDir.resolve(hash); - if (!Files.exists(marker)) { - Files.writeString(marker, "1", StandardCharsets.UTF_8); + markHashIngested(hash, "", ""); + } + + /** + * Writes or updates an ingest marker for the given chunk hash and associated metadata. + * + *

When a marker already exists and metadata changed, this method updates the marker payload so + * future dedup checks can detect metadata drift accurately.

+ * + * @param hash chunk hash marker key + * @param title document title associated with the ingested chunk + * @param packageName package name associated with the ingested chunk + * @throws IOException if marker write fails + */ + public void markHashIngested(String hash, String title, String packageName) throws IOException { + Path markerPath = hashMarkerPath(hash); + String normalizedTitle = normalizeHashMetadataText(title); + String normalizedPackageName = normalizeHashMetadataText(packageName); + String markerPayload = buildHashMarkerPayload(normalizedTitle, normalizedPackageName); + + if (!Files.exists(markerPath)) { + Files.writeString(markerPath, markerPayload, StandardCharsets.UTF_8); // Update progress after successful ingest if (progressTracker != null) { progressTracker.markChunkIndexed(); } + return; + } + + HashMarkerMetadata storedMetadata = readHashMarkerMetadata(markerPath); + if (!normalizedTitle.equals(storedMetadata.title()) + || !normalizedPackageName.equals(storedMetadata.packageName())) { + Files.writeString(markerPath, markerPayload, StandardCharsets.UTF_8); } } @@ -160,6 +216,10 @@ private Path fileMarkerPath(String url) { return indexDir.resolve(FILE_MARKER_PREFIX + safeName(url) + FILE_MARKER_EXTENSION); } + private Path hashMarkerPath(String hash) { + return indexDir.resolve(hash); + } + /** * Records a file-level ingestion marker keyed by URL, including the chunk hashes created for the file. * @@ -363,6 +423,61 @@ private String buildFileMarkerPayload( return payload.toString(); } + private String buildHashMarkerPayload(String title, String packageName) { + String encodedTitle = Base64.getEncoder().encodeToString(title.getBytes(StandardCharsets.UTF_8)); + String encodedPackageName = Base64.getEncoder().encodeToString(packageName.getBytes(StandardCharsets.UTF_8)); + StringBuilder markerPayload = new StringBuilder(); + markerPayload.append(HASH_MARKER_INGESTED_FLAG).append('\n'); + markerPayload.append(HASH_MARKER_TITLE_PREFIX).append(encodedTitle).append('\n'); + markerPayload.append(HASH_MARKER_PACKAGE_PREFIX) + .append(encodedPackageName) + .append('\n'); + return markerPayload.toString(); + } + + private HashMarkerMetadata readHashMarkerMetadata(Path markerPath) throws IOException { + String markerPayload = Files.readString(markerPath, StandardCharsets.UTF_8); + String resolvedTitle = ""; + String resolvedPackageName = ""; + for (String markerLine : markerPayload.split("\n")) { + String trimmedMarkerLine = markerLine == null ? "" : markerLine.trim(); + if (trimmedMarkerLine.startsWith(HASH_MARKER_TITLE_PREFIX)) { + String encodedTitle = trimmedMarkerLine.substring(HASH_MARKER_TITLE_PREFIX.length()).trim(); + resolvedTitle = decodeHashMetadataField(encodedTitle); + } else if (trimmedMarkerLine.startsWith(HASH_MARKER_PACKAGE_PREFIX)) { + String encodedPackageName = trimmedMarkerLine.substring(HASH_MARKER_PACKAGE_PREFIX.length()).trim(); + resolvedPackageName = decodeHashMetadataField(encodedPackageName); + } + } + return new HashMarkerMetadata(resolvedTitle, resolvedPackageName); + } + + private String decodeHashMetadataField(String encodedMetadataField) throws IOException { + if (encodedMetadataField == null || encodedMetadataField.isBlank()) { + return ""; + } + try { + byte[] decodedBytes = Base64.getDecoder().decode(encodedMetadataField); + return new String(decodedBytes, StandardCharsets.UTF_8); + } catch (IllegalArgumentException invalidBase64Exception) { + throw new IOException("Invalid hash marker metadata encoding", invalidBase64Exception); + } + } + + private String normalizeHashMetadataText(String metadataText) { + if (metadataText == null) { + return ""; + } + return metadataText.trim(); + } + + private record HashMarkerMetadata(String title, String packageName) { + private HashMarkerMetadata { + title = title == null ? "" : title; + packageName = packageName == null ? "" : packageName; + } + } + /** * File-level ingestion marker contents. * diff --git a/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java b/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java index ab03a4bb..51ed5cce 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java +++ b/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java @@ -370,14 +370,24 @@ private void markDocumentsIngested(List documents) { if (hashMetadata == null) { continue; } + String title = metadataText(doc, "title"); + String packageName = metadataText(doc, "package"); try { - localStore.markHashIngested(hashMetadata.toString()); + localStore.markHashIngested(hashMetadata.toString(), title, packageName); } catch (IOException markHashException) { throw new IllegalStateException("Failed to mark hash as ingested: " + hashMetadata, markHashException); } } } + private String metadataText(Document document, String metadataKey) { + Object metadataRaw = document.getMetadata().get(metadataKey); + if (metadataRaw == null) { + return ""; + } + return metadataRaw.toString(); + } + private void markFileIngested( String url, long fileSizeBytes, diff --git a/src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java b/src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java new file mode 100644 index 00000000..8d92766c --- /dev/null +++ b/src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java @@ -0,0 +1,75 @@ +package com.williamcallahan.javachat.service; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.anyInt; +import static org.mockito.Mockito.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.io.IOException; +import java.util.List; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** + * Verifies chunk processing deduplication behavior for hash and metadata updates. + */ +class ChunkProcessingServiceTest { + + private static final String SOURCE_URL = "https://example.test/reference"; + private static final String UPDATED_TITLE = "Updated Reference Title"; + private static final String PACKAGE_NAME = "com.example.reference"; + private static final String CHUNK_TEXT = "Chunk text used for deterministic hash checks."; + + private Chunker chunker; + private ContentHasher contentHasher; + private LocalStoreService localStoreService; + private ChunkProcessingService chunkProcessingService; + + @BeforeEach + void setUp() { + chunker = mock(Chunker.class); + contentHasher = new ContentHasher(); + DocumentFactory documentFactory = new DocumentFactory(contentHasher); + localStoreService = mock(LocalStoreService.class); + PdfContentExtractor pdfContentExtractor = mock(PdfContentExtractor.class); + chunkProcessingService = new ChunkProcessingService( + chunker, contentHasher, documentFactory, localStoreService, pdfContentExtractor); + } + + @Test + void processAndStoreChunks_reingestsWhenMetadataChangedForExistingHash() throws IOException { + when(chunker.chunkByTokens(anyString(), anyInt(), anyInt())).thenReturn(List.of(CHUNK_TEXT)); + + String chunkHash = contentHasher.generateChunkHash(SOURCE_URL, 0, CHUNK_TEXT); + when(localStoreService.isHashIngested(chunkHash)).thenReturn(true); + when(localStoreService.hasHashMetadataChanged(chunkHash, UPDATED_TITLE, PACKAGE_NAME)).thenReturn(true); + + ChunkProcessingService.ChunkProcessingOutcome processingOutcome = + chunkProcessingService.processAndStoreChunks("ignored", SOURCE_URL, UPDATED_TITLE, PACKAGE_NAME); + + assertEquals(1, processingOutcome.documents().size()); + assertEquals(0, processingOutcome.skippedChunks()); + verify(localStoreService).saveChunkText(SOURCE_URL, 0, CHUNK_TEXT, chunkHash); + } + + @Test + void processAndStoreChunks_skipsWhenHashIngestedAndMetadataUnchanged() throws IOException { + when(chunker.chunkByTokens(anyString(), anyInt(), anyInt())).thenReturn(List.of(CHUNK_TEXT)); + + String chunkHash = contentHasher.generateChunkHash(SOURCE_URL, 0, CHUNK_TEXT); + when(localStoreService.isHashIngested(chunkHash)).thenReturn(true); + when(localStoreService.hasHashMetadataChanged(chunkHash, UPDATED_TITLE, PACKAGE_NAME)).thenReturn(false); + + ChunkProcessingService.ChunkProcessingOutcome processingOutcome = + chunkProcessingService.processAndStoreChunks("ignored", SOURCE_URL, UPDATED_TITLE, PACKAGE_NAME); + + assertTrue(processingOutcome.documents().isEmpty()); + assertEquals(1, processingOutcome.skippedChunks()); + verify(localStoreService, never()).saveChunkText(SOURCE_URL, 0, CHUNK_TEXT, chunkHash); + } +} + From b8be7967c1c0c58ec807fe7dbda40a9c769366ed Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:15:50 -0800 Subject: [PATCH 02/19] fix(health): prevent stuck Qdrant DOWN state Health recovery logic could stall after transient outages: health snapshots did not trigger retries, unhealthy scheduled checks were effectively bypassed, and parallel retry paths could inflate failure state while risking backoff overflow. This change ensures retry evaluation happens when health is polled, gates health checks to one in-flight check per service status, and uses overflow-safe capped backoff progression. - Trigger retry evaluation from QdrantHealthIndicator.health() - Run unhealthy retry evaluation in scheduled health checks - Add in-progress check gating in ExternalServiceHealth.ServiceStatus - Replace Math.pow duration growth with overflow-safe capped doubling - Add tests for retry trigger behavior and concurrency/overflow handling --- .../config/QdrantHealthIndicator.java | 3 + .../service/ExternalServiceHealth.java | 136 +++++++++++++----- .../config/QdrantHealthIndicatorTest.java | 35 +++++ .../service/ExternalServiceHealthTest.java | 71 +++++++++ 4 files changed, 209 insertions(+), 36 deletions(-) create mode 100644 src/test/java/com/williamcallahan/javachat/config/QdrantHealthIndicatorTest.java create mode 100644 src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java diff --git a/src/main/java/com/williamcallahan/javachat/config/QdrantHealthIndicator.java b/src/main/java/com/williamcallahan/javachat/config/QdrantHealthIndicator.java index 06c18b9e..7188f661 100644 --- a/src/main/java/com/williamcallahan/javachat/config/QdrantHealthIndicator.java +++ b/src/main/java/com/williamcallahan/javachat/config/QdrantHealthIndicator.java @@ -40,6 +40,9 @@ public QdrantHealthIndicator(ExternalServiceHealth externalServiceHealth) { */ @Override public Health health() { + // Trigger retry checks when unhealthy and backoff has elapsed. + externalServiceHealth.isHealthy(ExternalServiceHealth.SERVICE_QDRANT); + ExternalServiceHealth.HealthSnapshot healthSnapshot = externalServiceHealth.getHealthSnapshot(ExternalServiceHealth.SERVICE_QDRANT); diff --git a/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java b/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java index 4bcbba5b..63bfbb51 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java +++ b/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java @@ -36,6 +36,7 @@ public class ExternalServiceHealth { // Health snapshot message templates private static final String HEALTHY_MSG_TEMPLATE = "Healthy (checked %s ago)"; private static final String UNHEALTHY_CHECKING_MSG = "Unhealthy (checking now...)"; + private static final String UNHEALTHY_RETRY_DUE_MSG = "Unhealthy (retry due)"; private static final String UNHEALTHY_NEXT_CHECK_TEMPLATE = "Unhealthy (failed %d times, next check in %s)"; private static final String UNKNOWN_SERVICE_MSG = "Unknown service"; @@ -131,8 +132,7 @@ public boolean isHealthy(String serviceName) { } // If unhealthy, check if we should retry based on backoff - Instant nextCheck = status.lastCheck.plus(status.currentBackoff); - if (Instant.now().isAfter(nextCheck)) { + if (status.shouldRetryNow(Instant.now())) { // Time to retry - trigger async health check if (SERVICE_QDRANT.equals(serviceName)) { checkQdrantHealthAsync(); @@ -164,7 +164,7 @@ public HealthSnapshot getHealthSnapshot(String serviceName) { timeUntilNextCheck = Duration.between(Instant.now(), status.lastCheck.plus(status.currentBackoff)); if (timeUntilNextCheck.isNegative()) { - message = UNHEALTHY_CHECKING_MSG; + message = status.isCheckInProgress() ? UNHEALTHY_CHECKING_MSG : UNHEALTHY_RETRY_DUE_MSG; timeUntilNextCheck = Duration.ZERO; } else { message = String.format( @@ -183,9 +183,14 @@ public HealthSnapshot getHealthSnapshot(String serviceName) { @Scheduled(fixedDelay = 3600000) // 1 hour public void scheduledQdrantHealthCheck() { ServiceStatus status = serviceStatuses.get(SERVICE_QDRANT); - if (status != null && status.isHealthy.get()) { + if (status == null) { + return; + } + if (status.isHealthy.get()) { checkQdrantHealth(); + return; } + isHealthy(SERVICE_QDRANT); } private void checkQdrantHealthAsync() { @@ -195,6 +200,9 @@ private void checkQdrantHealthAsync() { private void checkQdrantConnectivity() { ServiceStatus status = serviceStatuses.get(SERVICE_QDRANT); if (status == null) return; + if (!status.tryStartCheck(Instant.now())) { + return; + } String base = buildQdrantRestBaseUrl(); String connectivityPath = qdrantSsl ? "/collections" : "/health"; @@ -205,22 +213,34 @@ private void checkQdrantConnectivity() { requestSpec = requestSpec.header("api-key", qdrantApiKey); } - requestSpec - .retrieve() - .toBodilessEntity() - .timeout(HEALTH_CHECK_TIMEOUT) - .subscribe(response -> status.markHealthy(), error -> { - status.markUnhealthy(); - log.warn( - "Qdrant connectivity check failed (exception type: {}) - Will retry in {}", - error.getClass().getSimpleName(), - formatDuration(status.currentBackoff)); - }); + try { + requestSpec + .retrieve() + .toBodilessEntity() + .timeout(HEALTH_CHECK_TIMEOUT) + .subscribe(response -> status.markHealthy(), error -> { + status.markUnhealthy(); + log.warn( + "Qdrant connectivity check failed (exception type: {}) - Will retry in {}", + error.getClass().getSimpleName(), + formatDuration(status.currentBackoff)); + }); + } catch (RuntimeException connectivityException) { + status.markUnhealthy(); + log.warn( + "Qdrant connectivity check failed before subscription (exception type: {}) - Will retry in {}", + connectivityException.getClass().getSimpleName(), + formatDuration(status.currentBackoff), + connectivityException); + } } private void checkQdrantHealth() { ServiceStatus status = serviceStatuses.get(SERVICE_QDRANT); if (status == null) return; + if (!status.tryStartCheck(Instant.now())) { + return; + } if (qdrantCollections.isEmpty()) { status.markUnhealthy(); @@ -246,20 +266,29 @@ private void checkQdrantHealth() { .then()); } - Mono.whenDelayError(checks) - .timeout(HEALTH_CHECK_TIMEOUT) - .subscribe( - ignored -> { - status.markHealthy(); - log.debug("Qdrant health check succeeded (all collections present)"); - }, - error -> { - status.markUnhealthy(); - log.warn( - "Qdrant health check failed (exception type: {}) - Will retry in {}", - error.getClass().getSimpleName(), - formatDuration(status.currentBackoff)); - }); + try { + Mono.whenDelayError(checks) + .timeout(HEALTH_CHECK_TIMEOUT) + .subscribe( + ignored -> { + status.markHealthy(); + log.debug("Qdrant health check succeeded (all collections present)"); + }, + error -> { + status.markUnhealthy(); + log.warn( + "Qdrant health check failed (exception type: {}) - Will retry in {}", + error.getClass().getSimpleName(), + formatDuration(status.currentBackoff)); + }); + } catch (RuntimeException healthCheckException) { + status.markUnhealthy(); + log.warn( + "Qdrant health check failed before subscription (exception type: {}) - Will retry in {}", + healthCheckException.getClass().getSimpleName(), + formatDuration(status.currentBackoff), + healthCheckException); + } } private String buildQdrantRestBaseUrl() { @@ -343,6 +372,7 @@ private int mapGrpcPortToRestPort(int grpcPort) { */ private static class ServiceStatus { final AtomicBoolean isHealthy = new AtomicBoolean(false); + final AtomicBoolean checkInProgress = new AtomicBoolean(false); final AtomicInteger consecutiveFailures = new AtomicInteger(0); volatile Instant lastCheck = Instant.now(); volatile Duration currentBackoff = INITIAL_CHECK_INTERVAL; @@ -356,20 +386,16 @@ void markHealthy() { consecutiveFailures.set(0); currentBackoff = HEALTHY_CHECK_INTERVAL; lastCheck = Instant.now(); + checkInProgress.set(false); } void markUnhealthy() { isHealthy.set(false); int failures = consecutiveFailures.incrementAndGet(); - // Exponential backoff: 1min, 2min, 4min, 8min, ..., max 1 day - Duration newBackoff = INITIAL_CHECK_INTERVAL.multipliedBy((long) Math.pow(2, failures - 1)); - if (newBackoff.compareTo(MAX_CHECK_INTERVAL) > 0) { - newBackoff = MAX_CHECK_INTERVAL; - } - - currentBackoff = newBackoff; + currentBackoff = computeBackoffDuration(failures); lastCheck = Instant.now(); + checkInProgress.set(false); } void reset() { @@ -377,6 +403,44 @@ void reset() { consecutiveFailures.set(0); currentBackoff = INITIAL_CHECK_INTERVAL; lastCheck = Instant.EPOCH; // Force immediate check + checkInProgress.set(false); + } + + boolean tryStartCheck(Instant checkStartTime) { + if (!checkInProgress.compareAndSet(false, true)) { + return false; + } + lastCheck = checkStartTime; + return true; + } + + boolean shouldRetryNow(Instant evaluationTime) { + if (checkInProgress.get()) { + return false; + } + Instant nextCheck = lastCheck.plus(currentBackoff); + return !evaluationTime.isBefore(nextCheck); + } + + boolean isCheckInProgress() { + return checkInProgress.get(); + } + + private Duration computeBackoffDuration(int failureCount) { + Duration resolvedBackoff = INITIAL_CHECK_INTERVAL; + for (int failureIndex = 1; failureIndex < failureCount; failureIndex++) { + Duration doubledBackoff; + try { + doubledBackoff = resolvedBackoff.multipliedBy(2); + } catch (ArithmeticException overflowFailure) { + return MAX_CHECK_INTERVAL; + } + if (doubledBackoff.compareTo(MAX_CHECK_INTERVAL) >= 0) { + return MAX_CHECK_INTERVAL; + } + resolvedBackoff = doubledBackoff; + } + return resolvedBackoff; } } diff --git a/src/test/java/com/williamcallahan/javachat/config/QdrantHealthIndicatorTest.java b/src/test/java/com/williamcallahan/javachat/config/QdrantHealthIndicatorTest.java new file mode 100644 index 00000000..6dac9683 --- /dev/null +++ b/src/test/java/com/williamcallahan/javachat/config/QdrantHealthIndicatorTest.java @@ -0,0 +1,35 @@ +package com.williamcallahan.javachat.config; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.williamcallahan.javachat.service.ExternalServiceHealth; +import java.time.Duration; +import org.junit.jupiter.api.Test; +import org.springframework.boot.actuate.health.Health; +import org.springframework.boot.actuate.health.Status; + +/** + * Verifies that actuator health checks trigger retry evaluation in ExternalServiceHealth. + */ +class QdrantHealthIndicatorTest { + + @Test + void health_triggersRetryEvaluationBeforeSnapshotRendering() { + ExternalServiceHealth externalServiceHealth = mock(ExternalServiceHealth.class); + QdrantHealthIndicator qdrantHealthIndicator = new QdrantHealthIndicator(externalServiceHealth); + + ExternalServiceHealth.HealthSnapshot unhealthySnapshot = new ExternalServiceHealth.HealthSnapshot( + ExternalServiceHealth.SERVICE_QDRANT, false, "Unhealthy (checking now...)", Duration.ZERO); + when(externalServiceHealth.getHealthSnapshot(ExternalServiceHealth.SERVICE_QDRANT)) + .thenReturn(unhealthySnapshot); + + Health health = qdrantHealthIndicator.health(); + + verify(externalServiceHealth, times(1)).isHealthy(ExternalServiceHealth.SERVICE_QDRANT); + assertEquals(Status.DOWN, health.getStatus()); + } +} diff --git a/src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java b/src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java new file mode 100644 index 00000000..f1bef082 --- /dev/null +++ b/src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java @@ -0,0 +1,71 @@ +package com.williamcallahan.javachat.service; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Constructor; +import java.lang.reflect.Field; +import java.lang.reflect.Method; +import java.time.Duration; +import java.time.Instant; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.Test; + +/** + * Verifies retry-gate and backoff behavior in external service health tracking. + */ +class ExternalServiceHealthTest { + + @Test + void serviceStatus_marksUnhealthyWithoutDurationOverflow() throws Exception { + Object serviceStatus = newServiceStatus(); + AtomicInteger consecutiveFailures = readConsecutiveFailures(serviceStatus); + consecutiveFailures.set(62); + + Method markUnhealthyMethod = serviceStatus.getClass().getDeclaredMethod("markUnhealthy"); + markUnhealthyMethod.setAccessible(true); + markUnhealthyMethod.invoke(serviceStatus); + + Field currentBackoffField = serviceStatus.getClass().getDeclaredField("currentBackoff"); + currentBackoffField.setAccessible(true); + Duration currentBackoff = (Duration) currentBackoffField.get(serviceStatus); + assertEquals(Duration.ofDays(1), currentBackoff); + } + + @Test + void serviceStatus_allowsOnlyOneConcurrentCheckUntilCompletion() throws Exception { + Object serviceStatus = newServiceStatus(); + Method tryStartCheckMethod = serviceStatus.getClass().getDeclaredMethod("tryStartCheck", Instant.class); + tryStartCheckMethod.setAccessible(true); + + Instant startInstant = Instant.now(); + boolean firstCheckStarted = (boolean) tryStartCheckMethod.invoke(serviceStatus, startInstant); + boolean secondCheckStarted = (boolean) tryStartCheckMethod.invoke(serviceStatus, startInstant.plusSeconds(1)); + + assertTrue(firstCheckStarted); + assertFalse(secondCheckStarted); + + Method markHealthyMethod = serviceStatus.getClass().getDeclaredMethod("markHealthy"); + markHealthyMethod.setAccessible(true); + markHealthyMethod.invoke(serviceStatus); + + boolean checkStartedAfterCompletion = + (boolean) tryStartCheckMethod.invoke(serviceStatus, startInstant.plusSeconds(2)); + assertTrue(checkStartedAfterCompletion); + } + + private Object newServiceStatus() throws Exception { + Class serviceStatusClass = + Class.forName("com.williamcallahan.javachat.service.ExternalServiceHealth$ServiceStatus"); + Constructor constructor = serviceStatusClass.getDeclaredConstructor(String.class); + constructor.setAccessible(true); + return constructor.newInstance(ExternalServiceHealth.SERVICE_QDRANT); + } + + private AtomicInteger readConsecutiveFailures(Object serviceStatus) throws Exception { + Field consecutiveFailuresField = serviceStatus.getClass().getDeclaredField("consecutiveFailures"); + consecutiveFailuresField.setAccessible(true); + return (AtomicInteger) consecutiveFailuresField.get(serviceStatus); + } +} From 70a8b1c2f326f06107cc9dfe138e17891a64d1d8 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:16:07 -0800 Subject: [PATCH 03/19] refactor(chat): unify session state updates atomically Chat memory maintained two separate per-session structures (messages and turns) with non-atomic updates. Under concurrency, ordering could diverge, and session validation semantics relied on side-effect-prone access patterns. This refactor introduces a single per-session conversation object with synchronized atomic updates, plus explicit session recognition checks for validation endpoints. - Replace dual-map storage with SessionConversation aggregate per session - Update user/assistant writes to atomically update history and turns together - Add hasSession(String) for side-effect-free recognition checks - Update /api/chat/session/validate to avoid unknown-session creation side effects - Clarify session response docs and add concurrency/session validation tests --- .../javachat/service/ChatMemoryService.java | 80 ++++++++++++------- .../javachat/web/ChatController.java | 8 +- .../web/SessionValidationResponse.java | 2 +- .../service/ChatMemoryServiceTest.java | 45 +++++++++++ .../ChatControllerSessionValidationTest.java | 59 ++++++++++++++ 5 files changed, 160 insertions(+), 34 deletions(-) create mode 100644 src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java diff --git a/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java b/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java index 02d49241..ae4f0c21 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java +++ b/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java @@ -2,7 +2,6 @@ import com.williamcallahan.javachat.model.ChatTurn; import java.util.ArrayList; -import java.util.Collections; import java.util.List; import java.util.Objects; import java.util.concurrent.ConcurrentHashMap; @@ -23,10 +22,7 @@ public class ChatMemoryService { private static final String REQUIRE_SESSION_ID = "sessionId"; - // Use synchronizedList wrapper to ensure thread-safe list operations. - // ConcurrentHashMap only protects map operations, not the contained lists. - private final ConcurrentMap> sessionToMessages = new ConcurrentHashMap<>(); - private final ConcurrentMap> sessionToTurns = new ConcurrentHashMap<>(); + private final ConcurrentMap sessionConversations = new ConcurrentHashMap<>(); /** * Returns a thread-safe snapshot of the history for the given session. @@ -34,20 +30,11 @@ public class ChatMemoryService { */ public List getHistory(String sessionId) { Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); - List history = - sessionToMessages.computeIfAbsent(sessionId, _ -> Collections.synchronizedList(new ArrayList<>())); - // Return a snapshot to avoid ConcurrentModificationException during iteration - synchronized (history) { - return new ArrayList<>(history); + SessionConversation sessionConversation = sessionConversations.get(sessionId); + if (sessionConversation == null) { + return List.of(); } - } - - /** - * Returns the internal synchronized list for direct modification. - * Use with care - prefer addUser/addAssistant for adding messages. - */ - List getHistoryInternal(String sessionId) { - return sessionToMessages.computeIfAbsent(sessionId, _ -> Collections.synchronizedList(new ArrayList<>())); + return sessionConversation.historySnapshot(); } /** @@ -58,8 +45,9 @@ List getHistoryInternal(String sessionId) { */ public void addUser(String sessionId, String text) { Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); - getHistoryInternal(sessionId).add(new UserMessage(text)); - getTurnsInternal(sessionId).add(new ChatTurn("user", text)); + SessionConversation sessionConversation = + sessionConversations.computeIfAbsent(sessionId, missingSessionId -> new SessionConversation()); + sessionConversation.addUserMessage(text); } /** @@ -70,8 +58,9 @@ public void addUser(String sessionId, String text) { */ public void addAssistant(String sessionId, String text) { Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); - getHistoryInternal(sessionId).add(new AssistantMessage(text)); - getTurnsInternal(sessionId).add(new ChatTurn("assistant", text)); + SessionConversation sessionConversation = + sessionConversations.computeIfAbsent(sessionId, missingSessionId -> new SessionConversation()); + sessionConversation.addAssistantMessage(text); } /** @@ -81,8 +70,7 @@ public void addAssistant(String sessionId, String text) { */ public void clear(String sessionId) { Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); - sessionToMessages.remove(sessionId); - sessionToTurns.remove(sessionId); + sessionConversations.remove(sessionId); } /** @@ -90,18 +78,22 @@ public void clear(String sessionId) { */ public List getTurns(String sessionId) { Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); - List turns = - sessionToTurns.computeIfAbsent(sessionId, _ -> Collections.synchronizedList(new ArrayList<>())); - synchronized (turns) { - return new ArrayList<>(turns); + SessionConversation sessionConversation = sessionConversations.get(sessionId); + if (sessionConversation == null) { + return List.of(); } + return sessionConversation.turnSnapshot(); } /** - * Returns the internal synchronized list for direct modification. + * Returns true when the server currently recognizes the given session identifier. + * + * @param sessionId session identifier + * @return true when the session has been created in memory */ - List getTurnsInternal(String sessionId) { - return sessionToTurns.computeIfAbsent(sessionId, _ -> Collections.synchronizedList(new ArrayList<>())); + public boolean hasSession(String sessionId) { + Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); + return sessionConversations.containsKey(sessionId); } // TODO: Persist chat history embeddings to Qdrant for long-term memory (future feature) @@ -114,4 +106,30 @@ List getTurnsInternal(String sessionId) { // - Embedding strategy for chat turns (user messages + AI responses) // - Semantic similarity search for relevant historical context // - Privacy and data retention compliance + + /** + * Holds per-session conversation state and synchronizes updates across message and turn views. + */ + private static final class SessionConversation { + private final List historyMessages = new ArrayList<>(); + private final List turnHistory = new ArrayList<>(); + + synchronized void addUserMessage(String text) { + historyMessages.add(new UserMessage(text)); + turnHistory.add(new ChatTurn("user", text)); + } + + synchronized void addAssistantMessage(String text) { + historyMessages.add(new AssistantMessage(text)); + turnHistory.add(new ChatTurn("assistant", text)); + } + + synchronized List historySnapshot() { + return List.copyOf(historyMessages); + } + + synchronized List turnSnapshot() { + return List.copyOf(turnHistory); + } + } } diff --git a/src/main/java/com/williamcallahan/javachat/web/ChatController.java b/src/main/java/com/williamcallahan/javachat/web/ChatController.java index 357ed30d..1d55fb28 100644 --- a/src/main/java/com/williamcallahan/javachat/web/ChatController.java +++ b/src/main/java/com/williamcallahan/javachat/web/ChatController.java @@ -310,11 +310,15 @@ public ResponseEntity validateSession( if (sessionId == null || sessionId.isEmpty()) { return ResponseEntity.badRequest().body(new SessionValidationResponse("", 0, false, SESSION_ID_REQUIRED)); } + boolean sessionRecognized = chatMemory.hasSession(sessionId); + if (!sessionRecognized) { + return ResponseEntity.ok(new SessionValidationResponse(sessionId, 0, false, "Session not found on server")); + } var turns = chatMemory.getTurns(sessionId); int turnCount = turns.size(); boolean exists = turnCount > 0; - return ResponseEntity.ok(new SessionValidationResponse( - sessionId, turnCount, exists, exists ? "Session found" : "Session not found on server")); + String validationMessage = exists ? "Session found" : "Session found but empty"; + return ResponseEntity.ok(new SessionValidationResponse(sessionId, turnCount, exists, validationMessage)); } /** diff --git a/src/main/java/com/williamcallahan/javachat/web/SessionValidationResponse.java b/src/main/java/com/williamcallahan/javachat/web/SessionValidationResponse.java index b774fe80..0fb091fb 100644 --- a/src/main/java/com/williamcallahan/javachat/web/SessionValidationResponse.java +++ b/src/main/java/com/williamcallahan/javachat/web/SessionValidationResponse.java @@ -6,7 +6,7 @@ * * @param sessionId the session identifier that was validated * @param turnCount number of conversation turns (user + assistant messages) on server - * @param exists true if the session has any history on the server + * @param exists true if the session has any persisted conversation turns on the server * @param message human-readable status message */ public record SessionValidationResponse(String sessionId, int turnCount, boolean exists, String message) {} diff --git a/src/test/java/com/williamcallahan/javachat/service/ChatMemoryServiceTest.java b/src/test/java/com/williamcallahan/javachat/service/ChatMemoryServiceTest.java index 9711bbde..912ce2a2 100644 --- a/src/test/java/com/williamcallahan/javachat/service/ChatMemoryServiceTest.java +++ b/src/test/java/com/williamcallahan/javachat/service/ChatMemoryServiceTest.java @@ -3,6 +3,10 @@ import static org.junit.jupiter.api.Assertions.*; import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; @@ -75,6 +79,7 @@ void shouldReturnEmptyHistoryForUnknownSession() { assertNotNull(history); assertTrue(history.isEmpty()); + assertFalse(chatMemoryService.hasSession("nonexistent-session")); } @Test @@ -88,6 +93,7 @@ void shouldClearSessionHistory() { List history = chatMemoryService.getHistory(sessionId); assertTrue(history.isEmpty()); + assertFalse(chatMemoryService.hasSession(sessionId)); } @Test @@ -135,4 +141,43 @@ void shouldTrackTurnsAlongsideMessages() { assertEquals("assistant", turns.get(1).getRole()); assertEquals("Assistant answer", turns.get(1).getText()); } + + @Test + @DisplayName("Should keep message history and turn history aligned under concurrent writes") + void shouldKeepHistoryAndTurnsAlignedUnderConcurrency() throws InterruptedException { + String sessionId = "concurrent-session"; + int workerCount = 40; + ExecutorService workerPool = Executors.newFixedThreadPool(workerCount); + CountDownLatch startSignal = new CountDownLatch(1); + CountDownLatch completionSignal = new CountDownLatch(workerCount); + + for (int workerIndex = 0; workerIndex < workerCount; workerIndex++) { + final int messageIndex = workerIndex; + workerPool.submit(() -> { + try { + startSignal.await(); + chatMemoryService.addUser(sessionId, "message-" + messageIndex); + } catch (InterruptedException interruption) { + Thread.currentThread().interrupt(); + } finally { + completionSignal.countDown(); + } + }); + } + + startSignal.countDown(); + boolean completed = completionSignal.await(10, TimeUnit.SECONDS); + workerPool.shutdownNow(); + assertTrue(completed, "Concurrent writes did not complete within timeout"); + + List history = chatMemoryService.getHistory(sessionId); + var turns = chatMemoryService.getTurns(sessionId); + assertEquals(workerCount, history.size()); + assertEquals(workerCount, turns.size()); + for (int messageIndex = 0; messageIndex < workerCount; messageIndex++) { + String historyText = ((UserMessage) history.get(messageIndex)).getText(); + String turnText = turns.get(messageIndex).getText(); + assertEquals(historyText, turnText); + } + } } diff --git a/src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java b/src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java new file mode 100644 index 00000000..acd5e243 --- /dev/null +++ b/src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java @@ -0,0 +1,59 @@ +package com.williamcallahan.javachat.web; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.williamcallahan.javachat.config.AppProperties; +import com.williamcallahan.javachat.service.ChatMemoryService; +import org.junit.jupiter.api.Test; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; + +/** + * Verifies chat session validation semantics for unknown and recognized sessions. + */ +class ChatControllerSessionValidationTest { + + @Test + void validateSession_doesNotCreateUnknownSessionAndReportsRecognizedSessionHistory() { + ChatMemoryService chatMemoryService = new ChatMemoryService(); + ChatController chatController = new ChatController( + null, chatMemoryService, null, null, null, new ExceptionResponseBuilder(), new AppProperties()); + String unknownSessionId = "unknown-session-id"; + + assertFalse(chatMemoryService.hasSession(unknownSessionId)); + + ResponseEntity unknownSessionResponse = + chatController.validateSession(unknownSessionId); + assertEquals(HttpStatus.OK, unknownSessionResponse.getStatusCode()); + SessionValidationResponse unknownSessionBody = unknownSessionResponse.getBody(); + assertNotNull(unknownSessionBody); + assertFalse(unknownSessionBody.exists()); + assertEquals("Session not found on server", unknownSessionBody.message()); + assertFalse(chatMemoryService.hasSession(unknownSessionId)); + + chatMemoryService.addUser("recognized-session-id", "Stored message"); + ResponseEntity recognizedSessionResponse = + chatController.validateSession("recognized-session-id"); + SessionValidationResponse recognizedSessionBody = recognizedSessionResponse.getBody(); + assertNotNull(recognizedSessionBody); + assertTrue(recognizedSessionBody.exists()); + assertEquals("Session found", recognizedSessionBody.message()); + } + + @Test + void validateSession_returnsBadRequestWhenSessionIdIsBlank() { + ChatMemoryService chatMemoryService = new ChatMemoryService(); + ChatController chatController = new ChatController( + null, chatMemoryService, null, null, null, new ExceptionResponseBuilder(), new AppProperties()); + + ResponseEntity response = chatController.validateSession(""); + SessionValidationResponse body = response.getBody(); + assertNotNull(body); + assertEquals(HttpStatus.BAD_REQUEST, response.getStatusCode()); + assertFalse(body.exists()); + assertEquals("Session ID is required", body.message()); + } +} From de8de1d556d8fa20d8af9ff793cdb1a64159385c Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:17:27 -0800 Subject: [PATCH 04/19] fix(web): harden exception detail building for nulls Some HTTP and provider exception types can expose null status text, headers, or response payload values. The previous formatter called methods on those values without guarding for null, which could throw while we were already handling an upstream failure. This change hardens exception detail assembly so error reporting stays reliable under malformed or partial exception metadata. - Guard RestClient status text and response body before blank checks - Guard WebClient status text, response body, and headers before use - Guard OpenAI exception headers/body values before appending - Add regression test covering null status text handling --- .../web/ExceptionResponseBuilder.java | 20 ++++++++++--------- .../web/ExceptionResponseBuilderTest.java | 15 ++++++++++++++ 2 files changed, 26 insertions(+), 9 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/web/ExceptionResponseBuilder.java b/src/main/java/com/williamcallahan/javachat/web/ExceptionResponseBuilder.java index 97d858a3..25ea2b97 100644 --- a/src/main/java/com/williamcallahan/javachat/web/ExceptionResponseBuilder.java +++ b/src/main/java/com/williamcallahan/javachat/web/ExceptionResponseBuilder.java @@ -88,11 +88,11 @@ public String describeException(Exception exception) { private void appendRestClientDetails(StringBuilder details, RestClientResponseException exception) { details.append(" [httpStatus=").append(exception.getStatusCode().value()); String statusText = exception.getStatusText(); - if (!statusText.isBlank()) { + if (statusText != null && !statusText.isBlank()) { details.append(" ").append(statusText); } String responseBody = exception.getResponseBodyAsString(); - if (!responseBody.isBlank()) { + if (responseBody != null && !responseBody.isBlank()) { details.append(", body=").append(responseBody); } HttpHeaders headers = @@ -106,15 +106,15 @@ private void appendRestClientDetails(StringBuilder details, RestClientResponseEx private void appendWebClientDetails(StringBuilder details, WebClientResponseException exception) { details.append(" [httpStatus=").append(exception.getStatusCode().value()); String statusText = exception.getStatusText(); - if (!statusText.isBlank()) { + if (statusText != null && !statusText.isBlank()) { details.append(" ").append(statusText); } String responseBody = exception.getResponseBodyAsString(); - if (!responseBody.isBlank()) { + if (responseBody != null && !responseBody.isBlank()) { details.append(", body=").append(responseBody); } HttpHeaders headers = exception.getHeaders(); - if (!headers.isEmpty()) { + if (headers != null && !headers.isEmpty()) { details.append(", headers=").append(headers); } details.append("]"); @@ -123,13 +123,15 @@ private void appendWebClientDetails(StringBuilder details, WebClientResponseExce private void appendOpenAiDetails(StringBuilder details, OpenAIServiceException exception) { details.append(" [httpStatus=").append(exception.statusCode()); var headers = exception.headers(); - if (!headers.isEmpty()) { + if (headers != null && !headers.isEmpty()) { details.append(", headers=").append(headers); } var bodyJson = exception.body(); - String body = bodyJson.toString(); - if (!body.isBlank()) { - details.append(", body=").append(body); + if (bodyJson != null) { + String body = bodyJson.toString(); + if (!body.isBlank()) { + details.append(", body=").append(body); + } } exception.code().ifPresent(code -> details.append(", code=").append(code)); exception.param().ifPresent(param -> details.append(", param=").append(param)); diff --git a/src/test/java/com/williamcallahan/javachat/web/ExceptionResponseBuilderTest.java b/src/test/java/com/williamcallahan/javachat/web/ExceptionResponseBuilderTest.java index eaade794..49fa763a 100644 --- a/src/test/java/com/williamcallahan/javachat/web/ExceptionResponseBuilderTest.java +++ b/src/test/java/com/williamcallahan/javachat/web/ExceptionResponseBuilderTest.java @@ -1,5 +1,6 @@ package com.williamcallahan.javachat.web; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertTrue; import java.nio.charset.StandardCharsets; @@ -31,4 +32,18 @@ void describeException_includesHttpStatusAndBody() { assertTrue(details.contains("body=problem"), details); assertTrue(details.contains("headers="), details); } + + @Test + void describeException_handlesNullStatusTextWithoutThrowing() { + HttpHeaders headers = new HttpHeaders(); + HttpClientErrorException exception = HttpClientErrorException.create( + HttpStatus.BAD_REQUEST, + null, + headers, + "problem".getBytes(StandardCharsets.UTF_8), + StandardCharsets.UTF_8); + + ExceptionResponseBuilder builder = new ExceptionResponseBuilder(); + assertDoesNotThrow(() -> builder.describeException(exception)); + } } From 7851c46fbe74d5e2af6d279af837d6e2009e995a Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:17:42 -0800 Subject: [PATCH 05/19] fix(rate-limit): preserve failure streak counters Rate-limit window expiry was clearing consecutive failure counters during availability checks, and rate-limit events were not always incrementing total failure metrics. That made resilience telemetry and backoff behavior less trustworthy after throttling events. This change keeps failure history intact across expiry boundaries and records every rate-limit event consistently. - Increment total failure count when recording rate-limit events - Stop resetting consecutive failures when availability window expires - Add focused tests for expiry behavior and total failure increments --- .../javachat/service/RateLimitState.java | 2 +- .../javachat/service/RateLimitStateTest.java | 64 +++++++++++++++++++ 2 files changed, 65 insertions(+), 1 deletion(-) create mode 100644 src/test/java/com/williamcallahan/javachat/service/RateLimitStateTest.java diff --git a/src/main/java/com/williamcallahan/javachat/service/RateLimitState.java b/src/main/java/com/williamcallahan/javachat/service/RateLimitState.java index 680a9fd0..e3025f99 100644 --- a/src/main/java/com/williamcallahan/javachat/service/RateLimitState.java +++ b/src/main/java/com/williamcallahan/javachat/service/RateLimitState.java @@ -91,6 +91,7 @@ public void recordRateLimit(String provider, Instant resetTime, String rateLimit state.rateLimitedUntil = resetTime; int failures = state.consecutiveFailures.incrementAndGet(); + state.totalFailures.incrementAndGet(); state.lastFailure = Instant.now(); // Implement exponential backoff for repeated failures @@ -140,7 +141,6 @@ public boolean isAvailable(String provider) { // Clear rate limit if it has expired if (state.rateLimitedUntil != null && Instant.now().isAfter(state.rateLimitedUntil)) { state.rateLimitedUntil = null; - state.consecutiveFailures.set(0); safeSaveState(); } diff --git a/src/test/java/com/williamcallahan/javachat/service/RateLimitStateTest.java b/src/test/java/com/williamcallahan/javachat/service/RateLimitStateTest.java new file mode 100644 index 00000000..d16b15dc --- /dev/null +++ b/src/test/java/com/williamcallahan/javachat/service/RateLimitStateTest.java @@ -0,0 +1,64 @@ +package com.williamcallahan.javachat.service; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; +import java.lang.reflect.Field; +import java.time.Duration; +import java.time.Instant; +import java.util.Map; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** + * Verifies rate-limit counters and backoff state transitions. + */ +class RateLimitStateTest { + + private static final String PROVIDER_NAME = "provider-under-test"; + + private RateLimitState rateLimitState; + + @BeforeEach + void setUp() { + ObjectMapper objectMapper = new ObjectMapper(); + objectMapper.registerModule(new JavaTimeModule()); + rateLimitState = new RateLimitState(objectMapper); + } + + @Test + void isAvailable_doesNotResetConsecutiveFailuresWhenWindowExpires() throws Exception { + Instant expiredResetTime = Instant.now().minus(Duration.ofSeconds(5)); + rateLimitState.recordRateLimit(PROVIDER_NAME, expiredResetTime, "1m"); + + RateLimitState.ProviderState providerState = providerState(PROVIDER_NAME); + assertEquals(1, providerState.getConsecutiveFailures()); + + assertTrue(rateLimitState.isAvailable(PROVIDER_NAME)); + assertEquals(1, providerState.getConsecutiveFailures()); + } + + @Test + void recordRateLimit_incrementsTotalFailuresCounter() throws Exception { + rateLimitState.recordRateLimit(PROVIDER_NAME, null, "1m"); + + RateLimitState.ProviderState providerState = providerState(PROVIDER_NAME); + assertEquals(1, providerState.getTotalFailures()); + } + + private RateLimitState.ProviderState providerState(String providerName) throws Exception { + Field providerStatesField = RateLimitState.class.getDeclaredField("providerStates"); + providerStatesField.setAccessible(true); + Object providerStatesRaw = providerStatesField.get(rateLimitState); + if (!(providerStatesRaw instanceof Map providerStateMap)) { + throw new IllegalStateException("providerStates field does not contain a map"); + } + Object providerStateRaw = providerStateMap.get(providerName); + if (!(providerStateRaw instanceof RateLimitState.ProviderState providerState)) { + throw new IllegalStateException("No provider state found for provider: " + providerName); + } + return providerState; + } +} From 4f3966ccd4df3b8e1779fd27a6cd086cd2d6aa7b Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:17:50 -0800 Subject: [PATCH 06/19] fix(embedding): align dimension defaults and hints Embedding configuration default dimensions and request hint behavior were out of sync for mixed provider/model usage. The client also sent dimensions indiscriminately, which can break providers that do not accept that field for their embedding models. This change aligns defaults with configured expectations and applies dimension hints only for models that support override semantics. - Set embedding default dimensions to 4096 in AppProperties - Add model-aware dimension override gating for text-embedding-3 models - Keep request construction typed while conditionally applying dimensions - Add regression tests for include/omit dimension request behavior --- .../javachat/config/AppProperties.java | 2 +- .../OpenAiCompatibleEmbeddingClient.java | 18 ++++-- .../OpenAiCompatibleEmbeddingClientTest.java | 63 +++++++++++++++++++ 3 files changed, 78 insertions(+), 5 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/config/AppProperties.java b/src/main/java/com/williamcallahan/javachat/config/AppProperties.java index e86e5af4..e0cdcf75 100644 --- a/src/main/java/com/williamcallahan/javachat/config/AppProperties.java +++ b/src/main/java/com/williamcallahan/javachat/config/AppProperties.java @@ -476,7 +476,7 @@ QdrantCollections validateConfiguration() { /** Embedding vector configuration. */ public static class Embeddings { - private int dimensions = 1536; + private int dimensions = 4096; public int getDimensions() { return dimensions; diff --git a/src/main/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClient.java b/src/main/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClient.java index 9ac9a2aa..880fb2b9 100644 --- a/src/main/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClient.java +++ b/src/main/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClient.java @@ -80,10 +80,12 @@ public List embed(List texts) { if (texts == null || texts.isEmpty()) { return List.of(); } - EmbeddingCreateParams params = EmbeddingCreateParams.builder() - .model(modelName) - .inputOfArrayOfStrings(texts) - .build(); + EmbeddingCreateParams.Builder embeddingRequestBuilder = + EmbeddingCreateParams.builder().model(modelName).inputOfArrayOfStrings(texts); + if (supportsDimensionOverride(modelName)) { + embeddingRequestBuilder.dimensions((long) dimensionsHint); + } + EmbeddingCreateParams params = embeddingRequestBuilder.build(); RequestOptions requestOptions = RequestOptions.builder().timeout(embeddingTimeout()).build(); return executeWithRetry(params, requestOptions, texts.size()); @@ -291,6 +293,14 @@ public int dimensions() { return dimensionsHint; } + private static boolean supportsDimensionOverride(String embeddingModelName) { + if (embeddingModelName == null || embeddingModelName.isBlank()) { + return false; + } + String normalizedModelName = embeddingModelName.trim().toLowerCase(Locale.ROOT); + return normalizedModelName.startsWith("text-embedding-3"); + } + private float[] toFloatVector(List embeddingEntries) { if (embeddingEntries == null || embeddingEntries.isEmpty()) { throw new EmbeddingServiceUnavailableException("Remote embedding response missing embedding values"); diff --git a/src/test/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClientTest.java b/src/test/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClientTest.java index 8e0b89cf..c6123f8b 100644 --- a/src/test/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClientTest.java +++ b/src/test/java/com/williamcallahan/javachat/service/OpenAiCompatibleEmbeddingClientTest.java @@ -12,9 +12,11 @@ import com.openai.client.OpenAIClient; import com.openai.core.RequestOptions; import com.openai.models.embeddings.CreateEmbeddingResponse; +import com.openai.models.embeddings.EmbeddingCreateParams; import com.openai.services.blocking.EmbeddingService; import java.util.List; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; /** * Verifies OpenAI embedding responses preserve request ordering. @@ -137,4 +139,65 @@ void retriesTransientResponseValidationFailuresAndRecovers() { verify(embeddingService, times(2)).create(any(), any(RequestOptions.class)); } } + + @Test + void embed_includesDimensionsForTextEmbedding3Models() { + OpenAIClient client = mock(OpenAIClient.class); + EmbeddingService embeddingService = mock(EmbeddingService.class); + when(client.embeddings()).thenReturn(embeddingService); + + CreateEmbeddingResponse response = CreateEmbeddingResponse.builder() + .model("text-embedding-3-small") + .usage(CreateEmbeddingResponse.Usage.builder() + .promptTokens(1L) + .totalTokens(1L) + .build()) + .data(List.of(com.openai.models.embeddings.Embedding.builder() + .index(0L) + .embedding(List.of(0.4f, 0.6f)) + .build())) + .build(); + when(embeddingService.create(any(), any(RequestOptions.class))).thenReturn(response); + + try (OpenAiCompatibleEmbeddingClient clientAdapter = OpenAiCompatibleEmbeddingClient.create( + client, "text-embedding-3-small", EXPECTED_EMBEDDING_DIMENSION)) { + clientAdapter.embed(List.of("dimension check")); + + ArgumentCaptor requestCaptor = ArgumentCaptor.forClass(EmbeddingCreateParams.class); + verify(embeddingService).create(requestCaptor.capture(), any(RequestOptions.class)); + assertTrue(requestCaptor.getValue().dimensions().isPresent()); + assertEquals( + EXPECTED_EMBEDDING_DIMENSION, + requestCaptor.getValue().dimensions().get().intValue()); + } + } + + @Test + void embed_omitsDimensionsForNonTextEmbedding3Models() { + OpenAIClient client = mock(OpenAIClient.class); + EmbeddingService embeddingService = mock(EmbeddingService.class); + when(client.embeddings()).thenReturn(embeddingService); + + CreateEmbeddingResponse response = CreateEmbeddingResponse.builder() + .model("text-embedding-qwen3-embedding-8b") + .usage(CreateEmbeddingResponse.Usage.builder() + .promptTokens(1L) + .totalTokens(1L) + .build()) + .data(List.of(com.openai.models.embeddings.Embedding.builder() + .index(0L) + .embedding(List.of(0.4f, 0.6f)) + .build())) + .build(); + when(embeddingService.create(any(), any(RequestOptions.class))).thenReturn(response); + + try (OpenAiCompatibleEmbeddingClient clientAdapter = OpenAiCompatibleEmbeddingClient.create( + client, "text-embedding-qwen3-embedding-8b", EXPECTED_EMBEDDING_DIMENSION)) { + clientAdapter.embed(List.of("dimension check")); + + ArgumentCaptor requestCaptor = ArgumentCaptor.forClass(EmbeddingCreateParams.class); + verify(embeddingService).create(requestCaptor.capture(), any(RequestOptions.class)); + assertTrue(requestCaptor.getValue().dimensions().isEmpty()); + } + } } From 2fa2b02446d3fd8b4475c748e68d74a52a9d72b1 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:18:00 -0800 Subject: [PATCH 07/19] fix(frontend): strengthen session ID generation Session identifiers were built from timestamp plus Math.random-only suffixes, which are weaker and less stable across runtimes with varying crypto support. This could increase collision risk and reduce entropy in high-throughput client sessions. This change introduces a deterministic fallback chain that uses the strongest available browser/runtime primitive first, then degrades explicitly. - Add random segment generator with crypto.randomUUID primary path - Fallback to crypto.getRandomValues hex encoding when UUID is unavailable - Keep explicit padded Math.random fallback for legacy runtimes - Add unit tests covering all three entropy-source branches --- frontend/src/lib/utils/session.test.ts | 53 ++++++++++++++++++++++++++ frontend/src/lib/utils/session.ts | 17 ++++++++- 2 files changed, 69 insertions(+), 1 deletion(-) create mode 100644 frontend/src/lib/utils/session.test.ts diff --git a/frontend/src/lib/utils/session.test.ts b/frontend/src/lib/utils/session.test.ts new file mode 100644 index 00000000..10c752b6 --- /dev/null +++ b/frontend/src/lib/utils/session.test.ts @@ -0,0 +1,53 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { generateSessionId } from './session' + +describe('generateSessionId', () => { + afterEach(() => { + vi.useRealTimers() + vi.restoreAllMocks() + vi.unstubAllGlobals() + }) + + it('uses crypto.randomUUID when available', () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-02-09T12:00:00.000Z')) + vi.stubGlobal('crypto', { + randomUUID: () => 'uuid-test-value', + } as unknown as Crypto) + + const sessionId = generateSessionId('chat') + expect(sessionId).toBe('chat-1770638400000-uuid-test-value') + }) + + it('uses crypto.getRandomValues when randomUUID is unavailable', () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-02-09T12:00:00.000Z')) + vi.stubGlobal('crypto', { + getRandomValues: (randomBytes: Uint8Array) => { + randomBytes.fill(15) + return randomBytes + }, + } as unknown as Crypto) + + const sessionId = generateSessionId('chat') + const sessionParts = sessionId.split('-') + const randomSuffix = sessionParts[sessionParts.length - 1] + expect(randomSuffix).toHaveLength(32) + expect(/^[0-9a-f]+$/.test(randomSuffix)).toBe(true) + }) + + it('falls back to padded Math.random output when crypto is unavailable', () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-02-09T12:00:00.000Z')) + vi.stubGlobal('crypto', undefined) + vi.spyOn(Math, 'random').mockReturnValue(0) + + const sessionId = generateSessionId('chat') + const sessionParts = sessionId.split('-') + const randomSuffix = sessionParts[sessionParts.length - 1] + expect(sessionId.endsWith('-')).toBe(false) + expect(randomSuffix).toHaveLength(12) + expect(randomSuffix).toBe('000000000000') + }) +}) + diff --git a/frontend/src/lib/utils/session.ts b/frontend/src/lib/utils/session.ts index 49812c1a..534e54bf 100644 --- a/frontend/src/lib/utils/session.ts +++ b/frontend/src/lib/utils/session.ts @@ -2,6 +2,20 @@ * Session identifier utilities for client-side chat session management. */ +function createSessionRandomPart(): string { + if (typeof crypto !== 'undefined') { + if (typeof crypto.randomUUID === 'function') { + return crypto.randomUUID() + } + if (typeof crypto.getRandomValues === 'function') { + const randomBytes = new Uint8Array(16) + crypto.getRandomValues(randomBytes) + return Array.from(randomBytes, (randomByte) => randomByte.toString(16).padStart(2, '0')).join('') + } + } + return Math.random().toString(36).slice(2, 14).padEnd(12, '0') +} + /** * Generates a unique session identifier with a domain-specific prefix. * @@ -12,5 +26,6 @@ * @returns Unique session ID string in format "{prefix}-{timestamp}-{random}" */ export function generateSessionId(prefix: string): string { - return `${prefix}-${Date.now()}-${Math.random().toString(36).slice(2, 15)}` + const randomPart = createSessionRandomPart() + return `${prefix}-${Date.now()}-${randomPart}` } From 1a2e14183d256e1f341f5d66f42c66fa1ffe7bf6 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:43:31 -0800 Subject: [PATCH 08/19] chore(format): apply spotless formatting leftovers Pre-commit formatting (spotlessApply) reformatted two files during earlier commits, leaving non-functional line-wrap and trailing-line diffs unstaged. These diffs are noise-only and should be isolated so functional commits remain clean and fully focused on behavior changes. This commit captures only formatter output and does not alter runtime behavior or test intent. - Wrap long method-call lines to formatter-compliant style - Normalize chained append formatting in marker payload builder - Remove trailing blank line in test file - Keep logic and assertions unchanged --- .../javachat/service/LocalStoreService.java | 14 ++++++++++---- .../service/ChunkProcessingServiceTest.java | 7 ++++--- 2 files changed, 14 insertions(+), 7 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java b/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java index c7844557..0c7e9a3e 100644 --- a/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java +++ b/src/main/java/com/williamcallahan/javachat/service/LocalStoreService.java @@ -132,7 +132,8 @@ public boolean hasHashMetadataChanged(String hash, String title, String packageN return !normalizedTitle.equals(storedMetadata.title()) || !normalizedPackageName.equals(storedMetadata.packageName()); } catch (IOException markerReadFailure) { - throw new IllegalStateException("Failed to read hash ingestion marker for hash: " + hash, markerReadFailure); + throw new IllegalStateException( + "Failed to read hash ingestion marker for hash: " + hash, markerReadFailure); } } @@ -429,7 +430,8 @@ private String buildHashMarkerPayload(String title, String packageName) { StringBuilder markerPayload = new StringBuilder(); markerPayload.append(HASH_MARKER_INGESTED_FLAG).append('\n'); markerPayload.append(HASH_MARKER_TITLE_PREFIX).append(encodedTitle).append('\n'); - markerPayload.append(HASH_MARKER_PACKAGE_PREFIX) + markerPayload + .append(HASH_MARKER_PACKAGE_PREFIX) .append(encodedPackageName) .append('\n'); return markerPayload.toString(); @@ -442,10 +444,14 @@ private HashMarkerMetadata readHashMarkerMetadata(Path markerPath) throws IOExce for (String markerLine : markerPayload.split("\n")) { String trimmedMarkerLine = markerLine == null ? "" : markerLine.trim(); if (trimmedMarkerLine.startsWith(HASH_MARKER_TITLE_PREFIX)) { - String encodedTitle = trimmedMarkerLine.substring(HASH_MARKER_TITLE_PREFIX.length()).trim(); + String encodedTitle = trimmedMarkerLine + .substring(HASH_MARKER_TITLE_PREFIX.length()) + .trim(); resolvedTitle = decodeHashMetadataField(encodedTitle); } else if (trimmedMarkerLine.startsWith(HASH_MARKER_PACKAGE_PREFIX)) { - String encodedPackageName = trimmedMarkerLine.substring(HASH_MARKER_PACKAGE_PREFIX.length()).trim(); + String encodedPackageName = trimmedMarkerLine + .substring(HASH_MARKER_PACKAGE_PREFIX.length()) + .trim(); resolvedPackageName = decodeHashMetadataField(encodedPackageName); } } diff --git a/src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java b/src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java index 8d92766c..05f6ad0e 100644 --- a/src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java +++ b/src/test/java/com/williamcallahan/javachat/service/ChunkProcessingServiceTest.java @@ -46,7 +46,8 @@ void processAndStoreChunks_reingestsWhenMetadataChangedForExistingHash() throws String chunkHash = contentHasher.generateChunkHash(SOURCE_URL, 0, CHUNK_TEXT); when(localStoreService.isHashIngested(chunkHash)).thenReturn(true); - when(localStoreService.hasHashMetadataChanged(chunkHash, UPDATED_TITLE, PACKAGE_NAME)).thenReturn(true); + when(localStoreService.hasHashMetadataChanged(chunkHash, UPDATED_TITLE, PACKAGE_NAME)) + .thenReturn(true); ChunkProcessingService.ChunkProcessingOutcome processingOutcome = chunkProcessingService.processAndStoreChunks("ignored", SOURCE_URL, UPDATED_TITLE, PACKAGE_NAME); @@ -62,7 +63,8 @@ void processAndStoreChunks_skipsWhenHashIngestedAndMetadataUnchanged() throws IO String chunkHash = contentHasher.generateChunkHash(SOURCE_URL, 0, CHUNK_TEXT); when(localStoreService.isHashIngested(chunkHash)).thenReturn(true); - when(localStoreService.hasHashMetadataChanged(chunkHash, UPDATED_TITLE, PACKAGE_NAME)).thenReturn(false); + when(localStoreService.hasHashMetadataChanged(chunkHash, UPDATED_TITLE, PACKAGE_NAME)) + .thenReturn(false); ChunkProcessingService.ChunkProcessingOutcome processingOutcome = chunkProcessingService.processAndStoreChunks("ignored", SOURCE_URL, UPDATED_TITLE, PACKAGE_NAME); @@ -72,4 +74,3 @@ void processAndStoreChunks_skipsWhenHashIngestedAndMetadataUnchanged() throws IO verify(localStoreService, never()).saveChunkText(SOURCE_URL, 0, CHUNK_TEXT, chunkHash); } } - From bf29e081b3bcfd2e56ea932862e1b0bc11f52d71 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Sun, 8 Feb 2026 22:45:37 -0800 Subject: [PATCH 09/19] test(frontend): deduplicate session test timer setup The session utility tests repeated identical fake-timer initialization in each case. This was introduced by formatter/hook interactions and left as an unstaged follow-up change. This commit centralizes shared timer setup in beforeEach to keep tests clear and reduce repetition without changing assertions or coverage. - Add beforeEach with fixed test clock initialization - Remove duplicated timer setup lines from individual test cases - Keep existing test scenarios and expectations unchanged --- frontend/src/lib/utils/session.test.ts | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/frontend/src/lib/utils/session.test.ts b/frontend/src/lib/utils/session.test.ts index 10c752b6..1217e6fd 100644 --- a/frontend/src/lib/utils/session.test.ts +++ b/frontend/src/lib/utils/session.test.ts @@ -1,7 +1,12 @@ -import { afterEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { generateSessionId } from './session' describe('generateSessionId', () => { + beforeEach(() => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-02-09T12:00:00.000Z')) + }) + afterEach(() => { vi.useRealTimers() vi.restoreAllMocks() @@ -9,8 +14,6 @@ describe('generateSessionId', () => { }) it('uses crypto.randomUUID when available', () => { - vi.useFakeTimers() - vi.setSystemTime(new Date('2026-02-09T12:00:00.000Z')) vi.stubGlobal('crypto', { randomUUID: () => 'uuid-test-value', } as unknown as Crypto) @@ -20,8 +23,6 @@ describe('generateSessionId', () => { }) it('uses crypto.getRandomValues when randomUUID is unavailable', () => { - vi.useFakeTimers() - vi.setSystemTime(new Date('2026-02-09T12:00:00.000Z')) vi.stubGlobal('crypto', { getRandomValues: (randomBytes: Uint8Array) => { randomBytes.fill(15) @@ -37,8 +38,6 @@ describe('generateSessionId', () => { }) it('falls back to padded Math.random output when crypto is unavailable', () => { - vi.useFakeTimers() - vi.setSystemTime(new Date('2026-02-09T12:00:00.000Z')) vi.stubGlobal('crypto', undefined) vi.spyOn(Math, 'random').mockReturnValue(0) From 5cb023cff583cf1bc9afd09fa824ada062c76585 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 13:55:54 -0800 Subject: [PATCH 10/19] feat(analytics): add Simple Analytics via Vite transformIndexHtml plugin MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The site needs privacy-friendly analytics tracking. Rather than hardcoding a script tag in index.html, this uses Vite's built-in transformIndexHtml hook to inject the correct script at build time — latest.dev.js in development (no-op) and latest.js in production (real tracking). - Switch defineConfig from object to function form to access ConfigEnv.mode - Add simple-analytics plugin using transformIndexHtml to return HtmlTagDescriptor - Extract CDN base URL to SIMPLE_ANALYTICS_CDN named constant - Remove stale "Serve from root" comment on base property --- frontend/vite.config.ts | 27 +++++++++++++++++++++++---- 1 file changed, 23 insertions(+), 4 deletions(-) diff --git a/frontend/vite.config.ts b/frontend/vite.config.ts index 07957b43..2239b2f8 100644 --- a/frontend/vite.config.ts +++ b/frontend/vite.config.ts @@ -1,9 +1,28 @@ import { defineConfig } from 'vite' import { svelte } from '@sveltejs/vite-plugin-svelte' -export default defineConfig({ - plugins: [svelte()], - // Serve from root +const SIMPLE_ANALYTICS_CDN = 'https://scripts.simpleanalyticscdn.com' + +export default defineConfig(({ mode }) => ({ + plugins: [ + svelte(), + { + name: 'simple-analytics', + transformIndexHtml() { + const scriptFile = mode === 'development' ? 'latest.dev.js' : 'latest.js' + return [ + { + tag: 'script', + attrs: { + async: true, + src: `${SIMPLE_ANALYTICS_CDN}/${scriptFile}`, + }, + injectTo: 'head', + }, + ] + }, + }, + ], base: '/', server: { port: 5173, @@ -37,4 +56,4 @@ export default defineConfig({ } } } -}) +})) From ae6e4ce92297cf82a89d7c3c009a2ffb0ddb46f7 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:04:57 -0800 Subject: [PATCH 11/19] fix(health): prevent checkInProgress stuck and startup race in ExternalServiceHealth Move URL construction and request spec building inside try blocks so pre-subscription failures reset checkInProgress via markUnhealthy. Add status.reset() before verifyQdrantCollectionsAfterStartup to prevent ApplicationReadyEvent from being blocked by in-flight check. Inline anemic checkQdrantHealthAsync wrapper. --- .../service/ExternalServiceHealth.java | 67 ++++++++++--------- 1 file changed, 35 insertions(+), 32 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java b/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java index 63bfbb51..6e511ac4 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java +++ b/src/main/java/com/williamcallahan/javachat/service/ExternalServiceHealth.java @@ -108,9 +108,17 @@ public void init() { * *

Hybrid collections may be created on {@link ApplicationReadyEvent}; this avoids false negatives * during early startup while still enforcing collection presence after initialization.

+ * + *

The initial connectivity check from {@link #init()} may still be in-flight (reactive subscribe + * is non-blocking). Resetting the service status forces an immediate full collection check regardless + * of whether the connectivity probe has completed.

*/ @EventListener(ApplicationReadyEvent.class) public void verifyQdrantCollectionsAfterStartup() { + ServiceStatus status = serviceStatuses.get(SERVICE_QDRANT); + if (status != null) { + status.reset(); + } checkQdrantHealth(); } @@ -133,9 +141,8 @@ public boolean isHealthy(String serviceName) { // If unhealthy, check if we should retry based on backoff if (status.shouldRetryNow(Instant.now())) { - // Time to retry - trigger async health check if (SERVICE_QDRANT.equals(serviceName)) { - checkQdrantHealthAsync(); + checkQdrantHealth(); } } @@ -193,10 +200,6 @@ public void scheduledQdrantHealthCheck() { isHealthy(SERVICE_QDRANT); } - private void checkQdrantHealthAsync() { - checkQdrantHealth(); - } - private void checkQdrantConnectivity() { ServiceStatus status = serviceStatuses.get(SERVICE_QDRANT); if (status == null) return; @@ -204,21 +207,21 @@ private void checkQdrantConnectivity() { return; } - String base = buildQdrantRestBaseUrl(); - String connectivityPath = qdrantSsl ? "/collections" : "/health"; - String url = base + connectivityPath; + try { + String base = buildQdrantRestBaseUrl(); + String connectivityPath = qdrantSsl ? "/collections" : "/health"; + String connectivityUrl = base + connectivityPath; - var requestSpec = webClient.get().uri(url); - if (qdrantApiKey != null && !qdrantApiKey.isBlank()) { - requestSpec = requestSpec.header("api-key", qdrantApiKey); - } + var requestSpec = webClient.get().uri(connectivityUrl); + if (qdrantApiKey != null && !qdrantApiKey.isBlank()) { + requestSpec = requestSpec.header("api-key", qdrantApiKey); + } - try { requestSpec .retrieve() .toBodilessEntity() .timeout(HEALTH_CHECK_TIMEOUT) - .subscribe(response -> status.markHealthy(), error -> { + .subscribe(ignored -> status.markHealthy(), error -> { status.markUnhealthy(); log.warn( "Qdrant connectivity check failed (exception type: {}) - Will retry in {}", @@ -248,25 +251,25 @@ private void checkQdrantHealth() { return; } - String base = buildQdrantRestBaseUrl(); - List> checks = new ArrayList<>(qdrantCollections.size()); - for (String collection : qdrantCollections) { - if (collection == null || collection.isBlank()) { - continue; - } - String url = base + "/collections/" + collection; - var requestSpec = webClient.get().uri(url); - if (qdrantApiKey != null && !qdrantApiKey.isBlank()) { - requestSpec = requestSpec.header("api-key", qdrantApiKey); + try { + String base = buildQdrantRestBaseUrl(); + List> checks = new ArrayList<>(qdrantCollections.size()); + for (String collection : qdrantCollections) { + if (collection == null || collection.isBlank()) { + continue; + } + String collectionUrl = base + "/collections/" + collection; + var requestSpec = webClient.get().uri(collectionUrl); + if (qdrantApiKey != null && !qdrantApiKey.isBlank()) { + requestSpec = requestSpec.header("api-key", qdrantApiKey); + } + checks.add(requestSpec + .retrieve() + .toBodilessEntity() + .timeout(HEALTH_CHECK_TIMEOUT) + .then()); } - checks.add(requestSpec - .retrieve() - .toBodilessEntity() - .timeout(HEALTH_CHECK_TIMEOUT) - .then()); - } - try { Mono.whenDelayError(checks) .timeout(HEALTH_CHECK_TIMEOUT) .subscribe( From fecd724b6a617644cc457fd7bcd4cd51bf75fafb Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:05:05 -0800 Subject: [PATCH 12/19] fix(web): extract session validation message constants in ChatController Replace inline string literals with SESSION_NOT_FOUND_MESSAGE, SESSION_FOUND_MESSAGE, and SESSION_FOUND_EMPTY_MESSAGE constants to comply with no-magic-literals rule. --- .../com/williamcallahan/javachat/web/ChatController.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/web/ChatController.java b/src/main/java/com/williamcallahan/javachat/web/ChatController.java index 1d55fb28..fa29104a 100644 --- a/src/main/java/com/williamcallahan/javachat/web/ChatController.java +++ b/src/main/java/com/williamcallahan/javachat/web/ChatController.java @@ -45,6 +45,9 @@ public class ChatController extends BaseController { private static final Logger PIPELINE_LOG = LoggerFactory.getLogger("PIPELINE"); private static final AtomicLong REQUEST_SEQUENCE = new AtomicLong(); private static final String SESSION_ID_REQUIRED = "Session ID is required"; + private static final String SESSION_NOT_FOUND_MESSAGE = "Session not found on server"; + private static final String SESSION_FOUND_MESSAGE = "Session found"; + private static final String SESSION_FOUND_EMPTY_MESSAGE = "Session found but empty"; private static final String PIPELINE_LOG_SEPARATOR = "============================================"; private final ChatService chatService; @@ -312,12 +315,12 @@ public ResponseEntity validateSession( } boolean sessionRecognized = chatMemory.hasSession(sessionId); if (!sessionRecognized) { - return ResponseEntity.ok(new SessionValidationResponse(sessionId, 0, false, "Session not found on server")); + return ResponseEntity.ok(new SessionValidationResponse(sessionId, 0, false, SESSION_NOT_FOUND_MESSAGE)); } var turns = chatMemory.getTurns(sessionId); int turnCount = turns.size(); boolean exists = turnCount > 0; - String validationMessage = exists ? "Session found" : "Session found but empty"; + String validationMessage = exists ? SESSION_FOUND_MESSAGE : SESSION_FOUND_EMPTY_MESSAGE; return ResponseEntity.ok(new SessionValidationResponse(sessionId, turnCount, exists, validationMessage)); } From d036370ccaae51eb3706fe0f74da9055e0c4faa3 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:05:11 -0800 Subject: [PATCH 13/19] fix(ingestion): rename banned abbreviation doc to pageDocument in ChunkProcessingService --- .../javachat/service/ChunkProcessingService.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java b/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java index c3634593..e3b53ed1 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java +++ b/src/main/java/com/williamcallahan/javachat/service/ChunkProcessingService.java @@ -185,9 +185,9 @@ public ChunkProcessingOutcome processPdfAndStoreWithPages( allChunkHashes.add(hash); boolean hashAlreadyIngested = hashIngestionLookup.isHashIngested(hash); if (!hashAlreadyIngested || hashIngestionLookup.hasMetadataChanged(hash, title, packageName)) { - Document doc = documentFactory.createDocumentWithPages( + Document pageDocument = documentFactory.createDocumentWithPages( chunkText, url, title, globalIndex, packageName, hash, pageIndex + 1, pageIndex + 1); - pageDocuments.add(doc); + pageDocuments.add(pageDocument); chunkTextStore.saveChunkText(url, globalIndex, chunkText, hash); } else { skipped++; From 36f5c63eb679bc92f6464786a715406fc9c05abe Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:05:27 -0800 Subject: [PATCH 14/19] refactor(ingestion): extract shared metadataText helper to DocumentFactory Move duplicated private metadataText methods from DocsIngestionService and LocalDocsFileIngestionProcessor into DocumentFactory as a static utility to eliminate DRY violation. --- .../javachat/service/DocsIngestionService.java | 12 ++---------- .../javachat/service/DocumentFactory.java | 15 +++++++++++++++ .../LocalDocsFileIngestionProcessor.java | 13 +++---------- 3 files changed, 20 insertions(+), 20 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java b/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java index d9adeed6..90152213 100644 --- a/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java +++ b/src/main/java/com/williamcallahan/javachat/service/DocsIngestionService.java @@ -157,8 +157,8 @@ private void markDocumentsIngested(List discoveredLinks = new java.util.ArrayList<>(); diff --git a/src/main/java/com/williamcallahan/javachat/service/DocumentFactory.java b/src/main/java/com/williamcallahan/javachat/service/DocumentFactory.java index 605a4e8d..99dfe982 100644 --- a/src/main/java/com/williamcallahan/javachat/service/DocumentFactory.java +++ b/src/main/java/com/williamcallahan/javachat/service/DocumentFactory.java @@ -112,6 +112,21 @@ public org.springframework.ai.document.Document createDocumentWithPages( return doc; } + /** + * Extracts a metadata value as a string, returning empty string when absent. + * + * @param document the document to read metadata from + * @param metadataKey the metadata key to look up + * @return the metadata value as a string, or empty string when the key is absent or null + */ + public static String metadataText(org.springframework.ai.document.Document document, String metadataKey) { + Object metadataRaw = document.getMetadata().get(metadataKey); + if (metadataRaw == null) { + return ""; + } + return metadataRaw.toString(); + } + private org.springframework.ai.document.Document createDocumentWithOptionalId(String text, String hash) { if (hash == null || hash.isBlank()) { return new org.springframework.ai.document.Document(text); diff --git a/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java b/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java index 51ed5cce..91998799 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java +++ b/src/main/java/com/williamcallahan/javachat/service/ingestion/LocalDocsFileIngestionProcessor.java @@ -3,6 +3,7 @@ import com.williamcallahan.javachat.config.DocsSourceRegistry; import com.williamcallahan.javachat.domain.ingestion.IngestionLocalFailure; import com.williamcallahan.javachat.service.ChunkProcessingService; +import com.williamcallahan.javachat.service.DocumentFactory; import com.williamcallahan.javachat.service.EmbeddingServiceUnavailableException; import com.williamcallahan.javachat.service.LocalStoreService; import com.williamcallahan.javachat.service.ProgressTracker; @@ -370,8 +371,8 @@ private void markDocumentsIngested(List documents) { if (hashMetadata == null) { continue; } - String title = metadataText(doc, "title"); - String packageName = metadataText(doc, "package"); + String title = DocumentFactory.metadataText(doc, "title"); + String packageName = DocumentFactory.metadataText(doc, "package"); try { localStore.markHashIngested(hashMetadata.toString(), title, packageName); } catch (IOException markHashException) { @@ -380,14 +381,6 @@ private void markDocumentsIngested(List documents) { } } - private String metadataText(Document document, String metadataKey) { - Object metadataRaw = document.getMetadata().get(metadataKey); - if (metadataRaw == null) { - return ""; - } - return metadataRaw.toString(); - } - private void markFileIngested( String url, long fileSizeBytes, From c9678835f25779a5e0943b057164c3d68151dac6 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:06:02 -0800 Subject: [PATCH 15/19] fix(chat): replace unused lambda param with underscore in ChatMemoryService --- .../williamcallahan/javachat/service/ChatMemoryService.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java b/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java index ae4f0c21..1c0207a5 100644 --- a/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java +++ b/src/main/java/com/williamcallahan/javachat/service/ChatMemoryService.java @@ -46,7 +46,7 @@ public List getHistory(String sessionId) { public void addUser(String sessionId, String text) { Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); SessionConversation sessionConversation = - sessionConversations.computeIfAbsent(sessionId, missingSessionId -> new SessionConversation()); + sessionConversations.computeIfAbsent(sessionId, _ -> new SessionConversation()); sessionConversation.addUserMessage(text); } @@ -59,7 +59,7 @@ public void addUser(String sessionId, String text) { public void addAssistant(String sessionId, String text) { Objects.requireNonNull(sessionId, REQUIRE_SESSION_ID); SessionConversation sessionConversation = - sessionConversations.computeIfAbsent(sessionId, missingSessionId -> new SessionConversation()); + sessionConversations.computeIfAbsent(sessionId, _ -> new SessionConversation()); sessionConversation.addAssistantMessage(text); } From 60f0226713aac9f451122c2282299aa29853b493 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:06:11 -0800 Subject: [PATCH 16/19] test(web): extract constants and rename generic variables in session validation test --- .../ChatControllerSessionValidationTest.java | 44 +++++++++++-------- 1 file changed, 25 insertions(+), 19 deletions(-) diff --git a/src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java b/src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java index acd5e243..5912f321 100644 --- a/src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java +++ b/src/test/java/com/williamcallahan/javachat/web/ChatControllerSessionValidationTest.java @@ -16,31 +16,37 @@ */ class ChatControllerSessionValidationTest { + private static final String UNKNOWN_SESSION_ID = "unknown-session-id"; + private static final String RECOGNIZED_SESSION_ID = "recognized-session-id"; + private static final String STORED_MESSAGE_TEXT = "Stored message"; + private static final String SESSION_NOT_FOUND_MESSAGE = "Session not found on server"; + private static final String SESSION_FOUND_MESSAGE = "Session found"; + private static final String SESSION_ID_REQUIRED_MESSAGE = "Session ID is required"; + @Test void validateSession_doesNotCreateUnknownSessionAndReportsRecognizedSessionHistory() { ChatMemoryService chatMemoryService = new ChatMemoryService(); ChatController chatController = new ChatController( null, chatMemoryService, null, null, null, new ExceptionResponseBuilder(), new AppProperties()); - String unknownSessionId = "unknown-session-id"; - assertFalse(chatMemoryService.hasSession(unknownSessionId)); + assertFalse(chatMemoryService.hasSession(UNKNOWN_SESSION_ID)); - ResponseEntity unknownSessionResponse = - chatController.validateSession(unknownSessionId); - assertEquals(HttpStatus.OK, unknownSessionResponse.getStatusCode()); - SessionValidationResponse unknownSessionBody = unknownSessionResponse.getBody(); + ResponseEntity unknownSessionEntity = + chatController.validateSession(UNKNOWN_SESSION_ID); + assertEquals(HttpStatus.OK, unknownSessionEntity.getStatusCode()); + SessionValidationResponse unknownSessionBody = unknownSessionEntity.getBody(); assertNotNull(unknownSessionBody); assertFalse(unknownSessionBody.exists()); - assertEquals("Session not found on server", unknownSessionBody.message()); - assertFalse(chatMemoryService.hasSession(unknownSessionId)); + assertEquals(SESSION_NOT_FOUND_MESSAGE, unknownSessionBody.message()); + assertFalse(chatMemoryService.hasSession(UNKNOWN_SESSION_ID)); - chatMemoryService.addUser("recognized-session-id", "Stored message"); - ResponseEntity recognizedSessionResponse = - chatController.validateSession("recognized-session-id"); - SessionValidationResponse recognizedSessionBody = recognizedSessionResponse.getBody(); + chatMemoryService.addUser(RECOGNIZED_SESSION_ID, STORED_MESSAGE_TEXT); + ResponseEntity recognizedSessionEntity = + chatController.validateSession(RECOGNIZED_SESSION_ID); + SessionValidationResponse recognizedSessionBody = recognizedSessionEntity.getBody(); assertNotNull(recognizedSessionBody); assertTrue(recognizedSessionBody.exists()); - assertEquals("Session found", recognizedSessionBody.message()); + assertEquals(SESSION_FOUND_MESSAGE, recognizedSessionBody.message()); } @Test @@ -49,11 +55,11 @@ void validateSession_returnsBadRequestWhenSessionIdIsBlank() { ChatController chatController = new ChatController( null, chatMemoryService, null, null, null, new ExceptionResponseBuilder(), new AppProperties()); - ResponseEntity response = chatController.validateSession(""); - SessionValidationResponse body = response.getBody(); - assertNotNull(body); - assertEquals(HttpStatus.BAD_REQUEST, response.getStatusCode()); - assertFalse(body.exists()); - assertEquals("Session ID is required", body.message()); + ResponseEntity blankSessionEntity = chatController.validateSession(""); + SessionValidationResponse blankSessionBody = blankSessionEntity.getBody(); + assertNotNull(blankSessionBody); + assertEquals(HttpStatus.BAD_REQUEST, blankSessionEntity.getStatusCode()); + assertFalse(blankSessionBody.exists()); + assertEquals(SESSION_ID_REQUIRED_MESSAGE, blankSessionBody.message()); } } From 8fe786872f123dc8856b48b745e3f5f3c4c6b381 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:06:48 -0800 Subject: [PATCH 17/19] test(health): add retry-gate and backoff progression tests for ExternalServiceHealth --- .../service/ExternalServiceHealthTest.java | 79 ++++++++++++++++++- 1 file changed, 76 insertions(+), 3 deletions(-) diff --git a/src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java b/src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java index f1bef082..608ff7d8 100644 --- a/src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java +++ b/src/test/java/com/williamcallahan/javachat/service/ExternalServiceHealthTest.java @@ -27,9 +27,7 @@ void serviceStatus_marksUnhealthyWithoutDurationOverflow() throws Exception { markUnhealthyMethod.setAccessible(true); markUnhealthyMethod.invoke(serviceStatus); - Field currentBackoffField = serviceStatus.getClass().getDeclaredField("currentBackoff"); - currentBackoffField.setAccessible(true); - Duration currentBackoff = (Duration) currentBackoffField.get(serviceStatus); + Duration currentBackoff = readCurrentBackoff(serviceStatus); assertEquals(Duration.ofDays(1), currentBackoff); } @@ -55,6 +53,75 @@ void serviceStatus_allowsOnlyOneConcurrentCheckUntilCompletion() throws Exceptio assertTrue(checkStartedAfterCompletion); } + @Test + void shouldRetryNow_returnsFalseWhenCheckInProgress() throws Exception { + Object serviceStatus = newServiceStatus(); + Method tryStartCheckMethod = serviceStatus.getClass().getDeclaredMethod("tryStartCheck", Instant.class); + tryStartCheckMethod.setAccessible(true); + Method shouldRetryNowMethod = serviceStatus.getClass().getDeclaredMethod("shouldRetryNow", Instant.class); + shouldRetryNowMethod.setAccessible(true); + + tryStartCheckMethod.invoke(serviceStatus, Instant.now()); + + boolean shouldRetry = (boolean) + shouldRetryNowMethod.invoke(serviceStatus, Instant.now().plusSeconds(3600)); + assertFalse(shouldRetry, "shouldRetryNow must return false while a check is in progress"); + } + + @Test + void shouldRetryNow_returnsFalseBeforeBackoffElapses() throws Exception { + Object serviceStatus = newServiceStatus(); + Method markUnhealthyMethod = serviceStatus.getClass().getDeclaredMethod("markUnhealthy"); + markUnhealthyMethod.setAccessible(true); + Method shouldRetryNowMethod = serviceStatus.getClass().getDeclaredMethod("shouldRetryNow", Instant.class); + shouldRetryNowMethod.setAccessible(true); + + markUnhealthyMethod.invoke(serviceStatus); + + Instant beforeBackoffElapses = Instant.now().plusSeconds(30); + boolean shouldRetry = (boolean) shouldRetryNowMethod.invoke(serviceStatus, beforeBackoffElapses); + assertFalse(shouldRetry, "shouldRetryNow must return false before backoff period elapses"); + } + + @Test + void shouldRetryNow_returnsTrueAfterBackoffElapses() throws Exception { + Object serviceStatus = newServiceStatus(); + Method markUnhealthyMethod = serviceStatus.getClass().getDeclaredMethod("markUnhealthy"); + markUnhealthyMethod.setAccessible(true); + Method shouldRetryNowMethod = serviceStatus.getClass().getDeclaredMethod("shouldRetryNow", Instant.class); + shouldRetryNowMethod.setAccessible(true); + + markUnhealthyMethod.invoke(serviceStatus); + + Instant afterBackoffElapses = Instant.now().plusSeconds(120); + boolean shouldRetry = (boolean) shouldRetryNowMethod.invoke(serviceStatus, afterBackoffElapses); + assertTrue(shouldRetry, "shouldRetryNow must return true after backoff period elapses"); + } + + @Test + void computeBackoffDuration_doublesExponentiallyUntilCap() throws Exception { + Object serviceStatus = newServiceStatus(); + Method markUnhealthyMethod = serviceStatus.getClass().getDeclaredMethod("markUnhealthy"); + markUnhealthyMethod.setAccessible(true); + + // First failure: 1 minute + markUnhealthyMethod.invoke(serviceStatus); + assertEquals(Duration.ofMinutes(1), readCurrentBackoff(serviceStatus)); + + // Second failure: 2 minutes + markUnhealthyMethod.invoke(serviceStatus); + assertEquals(Duration.ofMinutes(2), readCurrentBackoff(serviceStatus)); + + // Third failure: 4 minutes + markUnhealthyMethod.invoke(serviceStatus); + assertEquals(Duration.ofMinutes(4), readCurrentBackoff(serviceStatus)); + + // Fifth failure (after fourth = 8m): 16 minutes + markUnhealthyMethod.invoke(serviceStatus); + markUnhealthyMethod.invoke(serviceStatus); + assertEquals(Duration.ofMinutes(16), readCurrentBackoff(serviceStatus)); + } + private Object newServiceStatus() throws Exception { Class serviceStatusClass = Class.forName("com.williamcallahan.javachat.service.ExternalServiceHealth$ServiceStatus"); @@ -68,4 +135,10 @@ private AtomicInteger readConsecutiveFailures(Object serviceStatus) throws Excep consecutiveFailuresField.setAccessible(true); return (AtomicInteger) consecutiveFailuresField.get(serviceStatus); } + + private Duration readCurrentBackoff(Object serviceStatus) throws Exception { + Field currentBackoffField = serviceStatus.getClass().getDeclaredField("currentBackoff"); + currentBackoffField.setAccessible(true); + return (Duration) currentBackoffField.get(serviceStatus); + } } From 2f990c232a5d27b25b12123ac0993ea2d0c9a50c Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:07:12 -0800 Subject: [PATCH 18/19] test(ingestion): add Base64 metadata round-trip and drift detection tests for LocalStoreService --- .../service/LocalStoreServiceTest.java | 137 ++++++++++++++++++ 1 file changed, 137 insertions(+) create mode 100644 src/test/java/com/williamcallahan/javachat/service/LocalStoreServiceTest.java diff --git a/src/test/java/com/williamcallahan/javachat/service/LocalStoreServiceTest.java b/src/test/java/com/williamcallahan/javachat/service/LocalStoreServiceTest.java new file mode 100644 index 00000000..58b83892 --- /dev/null +++ b/src/test/java/com/williamcallahan/javachat/service/LocalStoreServiceTest.java @@ -0,0 +1,137 @@ +package com.williamcallahan.javachat.service; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +/** + * Verifies hash marker metadata round-trip encoding and metadata drift detection + * in {@link LocalStoreService}. + */ +class LocalStoreServiceTest { + + private static final String SAMPLE_HASH = "abc123def456"; + private static final String SAMPLE_TITLE = "java.util.Optional"; + private static final String SAMPLE_PACKAGE = "java.util"; + private static final String DIFFERENT_TITLE = "java.util.List"; + private static final String DIFFERENT_PACKAGE = "java.util.stream"; + + @TempDir + Path tempDir; + + private LocalStoreService localStoreService; + + @BeforeEach + void setUp() { + String snapshotDir = tempDir.resolve("snapshots").toString(); + String parsedDir = tempDir.resolve("parsed").toString(); + String indexDir = tempDir.resolve("index").toString(); + ProgressTracker progressTracker = new ProgressTracker(parsedDir, indexDir); + localStoreService = new LocalStoreService(snapshotDir, parsedDir, indexDir, progressTracker); + localStoreService.createStoreDirectories(); + } + + @Test + void markHashIngested_roundTripsBase64MetadataCorrectly() throws IOException { + localStoreService.markHashIngested(SAMPLE_HASH, SAMPLE_TITLE, SAMPLE_PACKAGE); + + assertTrue(localStoreService.isHashIngested(SAMPLE_HASH)); + assertFalse( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, SAMPLE_TITLE, SAMPLE_PACKAGE), + "Metadata should match after round-trip"); + } + + @Test + void hasHashMetadataChanged_detectsTitleDrift() throws IOException { + localStoreService.markHashIngested(SAMPLE_HASH, SAMPLE_TITLE, SAMPLE_PACKAGE); + + assertTrue( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, DIFFERENT_TITLE, SAMPLE_PACKAGE), + "Changed title should be detected as metadata drift"); + } + + @Test + void hasHashMetadataChanged_detectsPackageDrift() throws IOException { + localStoreService.markHashIngested(SAMPLE_HASH, SAMPLE_TITLE, SAMPLE_PACKAGE); + + assertTrue( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, SAMPLE_TITLE, DIFFERENT_PACKAGE), + "Changed package should be detected as metadata drift"); + } + + @Test + void hasHashMetadataChanged_returnsFalseWhenMarkerDoesNotExist() { + assertFalse( + localStoreService.hasHashMetadataChanged("nonexistent_hash", SAMPLE_TITLE, SAMPLE_PACKAGE), + "Non-existent marker should report no change"); + } + + @Test + void markHashIngested_updatesMarkerWhenMetadataChanges() throws IOException { + localStoreService.markHashIngested(SAMPLE_HASH, SAMPLE_TITLE, SAMPLE_PACKAGE); + localStoreService.markHashIngested(SAMPLE_HASH, DIFFERENT_TITLE, DIFFERENT_PACKAGE); + + assertFalse( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, DIFFERENT_TITLE, DIFFERENT_PACKAGE), + "Metadata should match updated values after overwrite"); + assertTrue( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, SAMPLE_TITLE, SAMPLE_PACKAGE), + "Original metadata should now be detected as drift"); + } + + @Test + void hasHashMetadataChanged_throwsOnCorruptedBase64Marker() throws IOException { + Path markerPath = tempDir.resolve("index").resolve(SAMPLE_HASH); + String corruptedPayload = "1\ntitleB64=!!!not-valid-base64!!!\npackageB64=alsoBroken\n"; + Files.writeString(markerPath, corruptedPayload, StandardCharsets.UTF_8); + + IllegalStateException thrown = assertThrows( + IllegalStateException.class, + () -> localStoreService.hasHashMetadataChanged(SAMPLE_HASH, SAMPLE_TITLE, SAMPLE_PACKAGE)); + assertEquals("Failed to read hash ingestion marker for hash: " + SAMPLE_HASH, thrown.getMessage()); + } + + @Test + void markHashIngested_handlesNullAndEmptyMetadata() throws IOException { + localStoreService.markHashIngested(SAMPLE_HASH, null, null); + + assertTrue(localStoreService.isHashIngested(SAMPLE_HASH)); + assertFalse( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, "", ""), + "Null metadata normalizes to empty string and should match empty query"); + } + + @Test + void markHashIngested_noArgOverloadWritesEmptyMetadata() throws IOException { + localStoreService.markHashIngested(SAMPLE_HASH); + + assertTrue(localStoreService.isHashIngested(SAMPLE_HASH)); + assertFalse( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, "", ""), + "No-arg overload should write empty metadata"); + } + + @Test + void markHashIngested_roundTripsUnicodeMetadata() throws IOException { + String unicodeTitle = "クラス概要 — java.util.Optional"; + String unicodePackage = "日本語パッケージ"; + + localStoreService.markHashIngested(SAMPLE_HASH, unicodeTitle, unicodePackage); + + assertFalse( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, unicodeTitle, unicodePackage), + "Unicode metadata should survive Base64 round-trip"); + assertTrue( + localStoreService.hasHashMetadataChanged(SAMPLE_HASH, "ASCII title", unicodePackage), + "Changed title should be detected even when package contains Unicode"); + } +} From 24261e31ea09dbffa22526ecada38142cc695a18 Mon Sep 17 00:00:00 2001 From: William Callahan Date: Tue, 10 Feb 2026 17:07:20 -0800 Subject: [PATCH 19/19] chore(spotbugs): suppress CRLF_INJECTION_LOGS for ExternalServiceHealth --- config/spotbugs/spotbugs-exclude.xml | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/config/spotbugs/spotbugs-exclude.xml b/config/spotbugs/spotbugs-exclude.xml index 8d0e4f1c..69999fef 100644 --- a/config/spotbugs/spotbugs-exclude.xml +++ b/config/spotbugs/spotbugs-exclude.xml @@ -170,6 +170,10 @@ + + + +