From 334b55ec035c0c84eb00adc3b5cae3456cce44cc Mon Sep 17 00:00:00 2001 From: hubab1 Date: Sun, 19 Jul 2026 11:21:41 +0100 Subject: [PATCH] feat(research): persist ranking observations --- OpenASO.xcodeproj/project.pbxproj | 8 + OpenASO/App/AppServices.swift | 10 + .../AppDetail/AppDetailRefreshService.swift | 57 +- OpenASO/Services/MCP/OpenASOMCPService.swift | 77 +- .../KeywordResearchRankingWorkflow.swift | 259 +++ .../RankingRefreshCoordinator.swift | 805 +++++--- .../Storefront/AppCatalogService.swift | 94 +- .../KeywordResearchProjectStoreTests.swift | 2 + .../KeywordResearchRankingWorkflowTests.swift | 1659 +++++++++++++++++ OpenASOTests/OpenASOMCPServiceTests.swift | 152 +- .../RankingRefreshCoordinatorTests.swift | 632 ++++++- 11 files changed, 3413 insertions(+), 342 deletions(-) create mode 100644 OpenASO/Services/SearchRanking/KeywordResearchRankingWorkflow.swift create mode 100644 OpenASOTests/KeywordResearchRankingWorkflowTests.swift diff --git a/OpenASO.xcodeproj/project.pbxproj b/OpenASO.xcodeproj/project.pbxproj index 455c61d..8efa54d 100644 --- a/OpenASO.xcodeproj/project.pbxproj +++ b/OpenASO.xcodeproj/project.pbxproj @@ -137,6 +137,8 @@ D32C00000000000000000001 /* Persistence/KeywordResearchProjectStore.swift in Sources */ = {isa = PBXBuildFile; fileRef = D32C00000000000000000011 /* Persistence/KeywordResearchProjectStore.swift */; }; D32C00000000000000000002 /* Persistence/KeywordResearchProjectStore.swift in Sources */ = {isa = PBXBuildFile; fileRef = D32C00000000000000000011 /* Persistence/KeywordResearchProjectStore.swift */; }; D32C00000000000000000003 /* KeywordResearchProjectStoreTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = D32C00000000000000000012 /* KeywordResearchProjectStoreTests.swift */; }; + D33A00000000000000000001 /* SearchRanking/KeywordResearchRankingWorkflow.swift in Sources */ = {isa = PBXBuildFile; fileRef = D33A00000000000000000011 /* SearchRanking/KeywordResearchRankingWorkflow.swift */; }; + D33A00000000000000000002 /* KeywordResearchRankingWorkflowTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = D33A00000000000000000012 /* KeywordResearchRankingWorkflowTests.swift */; }; D20B00000000000000000007 /* Persistence/EstimatedKeywordDifficultyStore.swift in Sources */ = {isa = PBXBuildFile; fileRef = D20B00000000000000000014 /* Persistence/EstimatedKeywordDifficultyStore.swift */; }; D20B00000000000000000008 /* Persistence/EstimatedKeywordDifficultyStore.swift in Sources */ = {isa = PBXBuildFile; fileRef = D20B00000000000000000014 /* Persistence/EstimatedKeywordDifficultyStore.swift */; }; D20B00000000000000000009 /* EstimatedKeywordDifficultyStoreTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = D20B00000000000000000015 /* EstimatedKeywordDifficultyStoreTests.swift */; }; @@ -401,6 +403,8 @@ D32A00000000000000000014 /* KeywordResearchPersistenceTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = KeywordResearchPersistenceTests.swift; sourceTree = ""; }; D32C00000000000000000011 /* Persistence/KeywordResearchProjectStore.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Persistence/KeywordResearchProjectStore.swift; sourceTree = ""; }; D32C00000000000000000012 /* KeywordResearchProjectStoreTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = KeywordResearchProjectStoreTests.swift; sourceTree = ""; }; + D33A00000000000000000011 /* SearchRanking/KeywordResearchRankingWorkflow.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = SearchRanking/KeywordResearchRankingWorkflow.swift; sourceTree = ""; }; + D33A00000000000000000012 /* KeywordResearchRankingWorkflowTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = KeywordResearchRankingWorkflowTests.swift; sourceTree = ""; }; D20B00000000000000000014 /* Persistence/EstimatedKeywordDifficultyStore.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Persistence/EstimatedKeywordDifficultyStore.swift; sourceTree = ""; }; D20B00000000000000000015 /* EstimatedKeywordDifficultyStoreTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = EstimatedKeywordDifficultyStoreTests.swift; sourceTree = ""; }; A20000000000000000000011 /* KeywordTableCells.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = KeywordTableCells.swift; sourceTree = ""; }; @@ -548,6 +552,7 @@ D77BC4C40000000000000011 /* ExactV4StoreFixtureTests.swift */, D32A00000000000000000014 /* KeywordResearchPersistenceTests.swift */, D32C00000000000000000012 /* KeywordResearchProjectStoreTests.swift */, + D33A00000000000000000012 /* KeywordResearchRankingWorkflowTests.swift */, D09D00000000000000000007 /* TrackedKeywordRefreshStatusTests.swift */, D20B00000000000000000015 /* EstimatedKeywordDifficultyStoreTests.swift */, D02000000000000000000004 /* v0.3.2-v1 */, @@ -671,6 +676,7 @@ 9838D42ED1A3EA1E6DDDD68C /* OpenASOError.swift */, 035877A6E7443B32A3980E39 /* SearchRanking/RankingMatcher.swift */, 26D4B2C87FDED9CFCFA50D7D /* SearchRanking/RankingRefreshCoordinator.swift */, + D33A00000000000000000011 /* SearchRanking/KeywordResearchRankingWorkflow.swift */, A70000000000000000000002 /* RatingsReviews/ReviewLanguageDetectionService.swift */, 1C1EE58C790084CF7A83CA8A /* SearchRanking/SearchModels.swift */, F1F48387102DE88B407695C3 /* SearchRanking/SearchRankingProvider.swift */, @@ -1178,6 +1184,7 @@ 5FAADABB0DA1EEDB3699E66C /* SearchRanking/RankingMatcher.swift in Sources */, 33B2BE83B4CFC8970D7A3482 /* SearchRanking/RankingRefreshCoordinator.swift in Sources */, D20A00000000000000000001 /* SearchRanking/KeywordDifficultyEstimator.swift in Sources */, + D33A00000000000000000001 /* SearchRanking/KeywordResearchRankingWorkflow.swift in Sources */, 9A6937276AF0E3CAEDA90C36 /* RankingSource.swift in Sources */, A40000000000000000000003 /* RatingsDetail.swift in Sources */, A40000000000000000000001 /* RatingsDashboardSupport.swift in Sources */, @@ -1224,6 +1231,7 @@ D77BC4C40000000000000001 /* ExactV4StoreFixtureTests.swift in Sources */, D32A00000000000000000007 /* KeywordResearchPersistenceTests.swift in Sources */, D32C00000000000000000003 /* KeywordResearchProjectStoreTests.swift in Sources */, + D33A00000000000000000002 /* KeywordResearchRankingWorkflowTests.swift in Sources */, D09D00000000000000000008 /* TrackedKeywordRefreshStatusTests.swift in Sources */, D20B00000000000000000009 /* EstimatedKeywordDifficultyStoreTests.swift in Sources */, D14000000000000000000002 /* AppMetadataRefreshProgressStoreTests.swift in Sources */, diff --git a/OpenASO/App/AppServices.swift b/OpenASO/App/AppServices.swift index 5c1058d..d703b27 100644 --- a/OpenASO/App/AppServices.swift +++ b/OpenASO/App/AppServices.swift @@ -42,6 +42,7 @@ final class AppServices { let refreshProgressStore: AppRefreshProgressStore let mcpServerController: OpenASOMCPServerController let keywordResearchProjectStore: KeywordResearchProjectStore? + let keywordResearchRankingWorkflow: KeywordResearchRankingWorkflow? private(set) var backgroundModelStore: BackgroundModelStore? private(set) var backgroundModelStoreRevision = 0 @@ -55,6 +56,7 @@ final class AppServices { allowsIconNetworkFetches: Bool = true, backgroundModelStore: BackgroundModelStore? = nil, keywordResearchProjectStore: KeywordResearchProjectStore? = nil, + keywordResearchRankingWorkflow: KeywordResearchRankingWorkflow? = nil, refreshObservationClock: RefreshObservationClock = .live, refreshMetricsRecorder: RefreshMetricsRecorder? = nil, providerRequestGateMode: ProviderRequestGateMode? = nil, @@ -241,6 +243,13 @@ final class AppServices { }, metadataEnrichmentHandler: metadataEnrichmentHandler ) + let keywordResearchRankingWorkflow = keywordResearchRankingWorkflow + ?? backgroundModelStore.map { + KeywordResearchRankingWorkflow( + backgroundModelStore: $0, + rankingCoordinator: refreshCoordinator + ) + } let appDetailRefreshService = backgroundModelStore.map { AppDetailRefreshService( backgroundModelStore: $0, @@ -342,6 +351,7 @@ final class AppServices { self.appDetailRefreshService = appDetailRefreshService self.refreshProgressStore = refreshProgressStore self.keywordResearchProjectStore = keywordResearchProjectStore + self.keywordResearchRankingWorkflow = keywordResearchRankingWorkflow self.mcpServerController = OpenASOMCPServerController(portProvider: { settingsStore.mcpServerPort }) { diff --git a/OpenASO/Services/AppDetail/AppDetailRefreshService.swift b/OpenASO/Services/AppDetail/AppDetailRefreshService.swift index fdfaa17..c3902f4 100644 --- a/OpenASO/Services/AppDetail/AppDetailRefreshService.swift +++ b/OpenASO/Services/AppDetail/AppDetailRefreshService.swift @@ -88,7 +88,7 @@ private struct KeywordRefreshResult: Sendable { private struct RankingPersistenceBatchOutcome: Sendable { let outcomes: [KeywordBackgroundRefreshOutcome] let statsRebuildRequests: Set - let successfulPageResults: [RankingRefreshPageResult] + let metadataEnrichmentPageResults: [RankingRefreshPageResult] var failureCount: Int { outcomes.filter { $0.error != nil }.count @@ -576,7 +576,7 @@ final class AppDetailRefreshService: Sendable { outcomes.append(contentsOf: batchOutcome.outcomes) statsRebuildRequests.formUnion(batchOutcome.statsRebuildRequests) failureCount += batchOutcome.failureCount - for pageResult in batchOutcome.successfulPageResults { + for pageResult in batchOutcome.metadataEnrichmentPageResults { try Task.checkCancellation() refreshCoordinator.scheduleTopRankingMetadataEnrichment(for: pageResult) } @@ -682,13 +682,16 @@ final class AppDetailRefreshService: Sendable { do { try await backgroundModelStore.write { modelContext in try Task.checkCancellation() - refreshCoordinator.rebuildDerivedStats(for: requests, in: modelContext) + try refreshCoordinator.rebuildDerivedStats(for: requests, in: modelContext) } try Task.checkCancellation() } catch { if Task.isCancelled { throw CancellationError() } + OpenASOLog.refresh.error( + "Failed to rebuild ranking statistics: \(String(reflecting: error), privacy: .private(mask: .hash))" + ) } } @@ -711,44 +714,30 @@ final class AppDetailRefreshService: Sendable { try Task.checkCancellation() var outcomes: [KeywordBackgroundRefreshOutcome] = [] var statsRebuildRequests = Set() - var successfulPageResults: [RankingRefreshPageResult] = [] + var metadataEnrichmentPageResults: [RankingRefreshPageResult] = [] for pageResult in pageResults { - do { - _ = try refreshCoordinator.persistRankingPage( - pageResult, - in: modelContext, - rebuildDerivedStats: false, - saveChanges: false, - scheduleMetadataEnrichment: false - ) - if let statsRebuildRequest = RankingStatsRebuildRequest(pageRequest: pageResult.request) { - statsRebuildRequests.insert(statsRebuildRequest) - } - successfulPageResults.append(pageResult) - outcomes.append(KeywordBackgroundRefreshOutcome( - trackIdentityKey: pageResult.request.identityKey, - error: nil - )) - } catch { - let mappedError = OpenASOError.map(error) - _ = try? refreshCoordinator.recordRefreshFailure( - identityKey: pageResult.request.identityKey, - error: mappedError, - in: modelContext, - saveChanges: false - ) - outcomes.append(KeywordBackgroundRefreshOutcome( - trackIdentityKey: pageResult.request.identityKey, - error: mappedError - )) + let persistence = try refreshCoordinator.persistRankingPageTransaction( + pageResult, + in: modelContext, + rebuildDerivedStats: false + ) + if let statsRebuildRequest = RankingStatsRebuildRequest(pageRequest: pageResult.request) { + statsRebuildRequests.insert(statsRebuildRequest) + } + if persistence.appliedSharedObservation { + metadataEnrichmentPageResults.append(persistence.canonicalPageResult) } + outcomes.append(KeywordBackgroundRefreshOutcome( + trackIdentityKey: pageResult.request.identityKey, + error: nil + )) } return RankingPersistenceBatchOutcome( outcomes: outcomes, statsRebuildRequests: statsRebuildRequests, - successfulPageResults: successfulPageResults + metadataEnrichmentPageResults: metadataEnrichmentPageResults ) } try Task.checkCancellation() @@ -781,7 +770,7 @@ final class AppDetailRefreshService: Sendable { KeywordBackgroundRefreshOutcome(trackIdentityKey: $0.request.identityKey, error: mappedError) }, statsRebuildRequests: [], - successfulPageResults: [] + metadataEnrichmentPageResults: [] ) } } diff --git a/OpenASO/Services/MCP/OpenASOMCPService.swift b/OpenASO/Services/MCP/OpenASOMCPService.swift index 3dc6eee..0ac04c6 100644 --- a/OpenASO/Services/MCP/OpenASOMCPService.swift +++ b/OpenASO/Services/MCP/OpenASOMCPService.swift @@ -54,6 +54,11 @@ private struct OpenASOMCPHistoryCursor: Codable, Sendable { } } +private struct OpenASOMCPRankingRefreshCommit: Sendable { + let result: OpenASOMCPKeywordRefreshResult + let enrichmentPageResults: [RankingRefreshPageResult] +} + private enum OpenASOMCPHistoryCursorCodec { private static let version = 2 private static let maximumEncodedCursorLength = 4_096 @@ -1989,7 +1994,7 @@ final class OpenASOMCPService: Sendable { default: ResponseLimits.defaultKeywordRefreshTrackLimit, maximum: ResponseLimits.maximumKeywordRefreshTrackLimit ) - let resultLimit = ResponseLimits.defaultRankingAppLimit + let resultLimit = SearchRankingCrawl.fullKeywordRankingLimit let reconciliationPlan = await rankingRefreshScheduler.makeReconciliationPlan() let candidateBatch: ( @@ -2140,8 +2145,9 @@ final class OpenASOMCPService: Sendable { try Task.checkCancellation() let fetchedResults = fetched - let result = try await backgroundModelStore.write { modelContext in + let commit = try await backgroundModelStore.write { modelContext in var outcomes: [OpenASOMCPKeywordRefreshOutcome] = [] + var enrichmentPageResults: [RankingRefreshPageResult] = [] for item in fetchedResults { if let page = item.page, let observedAt = item.fetchedAt { let track: TrackedAppKeyword @@ -2155,13 +2161,15 @@ final class OpenASOMCPService: Sendable { winningCount: 0, confidence: nil ) - track = try rankingRefreshCoordinator.persistRankingPage( + let persistence = try rankingRefreshCoordinator.persistRankingPageTransaction( pageResult, in: modelContext, - rebuildDerivedStats: true, - saveChanges: false, - scheduleMetadataEnrichment: true - ).keywordTrack + rebuildDerivedStats: true + ) + track = persistence.snapshot.keywordTrack + if persistence.appliedSharedObservation { + enrichmentPageResults.append(persistence.canonicalPageResult) + } } else { track = try Self.persistRankingPage( page, @@ -2220,23 +2228,31 @@ final class OpenASOMCPService: Sendable { } } let failures = outcomes.filter { $0.error != nil }.count - return OpenASOMCPKeywordRefreshResult( - summary: OpenASOMCPMutationSummary( - inserted: 0, - updated: 0, - skipped: max(0, candidateBatch.totalCount - reservation.requests.count), - refreshed: outcomes.count - failures, - failed: failures + return OpenASOMCPRankingRefreshCommit( + result: OpenASOMCPKeywordRefreshResult( + summary: OpenASOMCPMutationSummary( + inserted: 0, + updated: 0, + skipped: max(0, candidateBatch.totalCount - reservation.requests.count), + refreshed: outcomes.count - failures, + failed: failures + ), + outcomes: outcomes, + notes: Self.keywordRankingRefreshNotes( + requestedLimit: requestedLimit, + appliedLimit: trackLimit + ) ), - outcomes: outcomes, - notes: Self.keywordRankingRefreshNotes( - requestedLimit: requestedLimit, - appliedLimit: trackLimit - ) + enrichmentPageResults: enrichmentPageResults ) } + if let rankingRefreshCoordinator { + for pageResult in commit.enrichmentPageResults { + rankingRefreshCoordinator.scheduleTopRankingMetadataEnrichment(for: pageResult) + } + } await rankingRefreshScheduler.release(reservation) - return result + return commit.result } catch { await rankingRefreshScheduler.release(reservation) throw error @@ -3278,7 +3294,9 @@ extension OpenASOMCPService { keyword: seed.keyword, storefrontCode: storefront, platform: platform, - limit: ResponseLimits.defaultRankingAppLimit + limit: persistResults + ? SearchRankingCrawl.fullKeywordRankingLimit + : ResponseLimits.defaultRankingAppLimit ) } } catch { @@ -3318,7 +3336,7 @@ extension OpenASOMCPService { storefront: storefront, platform: platform ) - try await backgroundModelStore.write { modelContext in + let enrichmentPageResults = try await backgroundModelStore.write { modelContext in let observedAt = now() if let rankingRefreshCoordinator { _ = try Self.ensureTrackedKeyword( @@ -3335,13 +3353,14 @@ extension OpenASOMCPService { winningCount: 0, confidence: nil ) - _ = try rankingRefreshCoordinator.persistRankingPage( + let persistence = try rankingRefreshCoordinator.persistRankingPageTransaction( pageResult, in: modelContext, - rebuildDerivedStats: true, - saveChanges: false, - scheduleMetadataEnrichment: true + rebuildDerivedStats: true ) + return persistence.appliedSharedObservation + ? [persistence.canonicalPageResult] + : [] } else { _ = try Self.persistRankingPage( page, @@ -3351,6 +3370,12 @@ extension OpenASOMCPService { appCatalogService: appCatalogService, in: modelContext ) + return [] + } + } + if let rankingRefreshCoordinator { + for pageResult in enrichmentPageResults { + rankingRefreshCoordinator.scheduleTopRankingMetadataEnrichment(for: pageResult) } } isTracked = true diff --git a/OpenASO/Services/SearchRanking/KeywordResearchRankingWorkflow.swift b/OpenASO/Services/SearchRanking/KeywordResearchRankingWorkflow.swift new file mode 100644 index 0000000..6e790be --- /dev/null +++ b/OpenASO/Services/SearchRanking/KeywordResearchRankingWorkflow.swift @@ -0,0 +1,259 @@ +import Foundation +import SwiftData + +struct KeywordResearchRankingItemSnapshot: Equatable, Hashable, Identifiable, Sendable { + let id: String + let position: Int + let appStoreID: Int64 + let bundleID: String? + let name: String + let subtitle: String? + let sellerName: String? +} + +struct KeywordResearchRankingObservationSnapshot: Equatable, Hashable, Identifiable, Sendable { + let id: String + let projectGeneration: KeywordResearchProjectGeneration + let keywordGeneration: KeywordResearchKeywordGeneration + let queryKey: String + let term: String + let storefront: String + let platform: AppPlatform + let observedAt: Date + let observedHour: Int + let source: RankingSource + let resultCount: Int + let submissionCount: Int + let winningCount: Int + let confidence: String? + let items: [KeywordResearchRankingItemSnapshot] +} + +/// App-only ranking workflow for pre-live research memberships. +/// +/// SwiftData models remain inside `BackgroundModelStore` operations. The +/// provider request runs between a read preflight and a single revalidated +/// write, so a deleted, replaced, or retargeted generation cannot receive a +/// late observation. +actor KeywordResearchRankingWorkflow { + /// Persisted crawls always request the full supported result window. + /// The shared crawl schema has no completeness dimension, so allowing a + /// caller-selected prefix could replace and prune a fuller same-day crawl. + static let persistedResultLimit = SearchRankingCrawl.fullKeywordRankingLimit + + private let modelStore: BackgroundModelStore + private let rankingCoordinator: RankingRefreshCoordinator + + init( + backgroundModelStore: BackgroundModelStore, + rankingCoordinator: RankingRefreshCoordinator + ) { + self.modelStore = backgroundModelStore + self.rankingCoordinator = rankingCoordinator + } + + func refresh( + projectGeneration: KeywordResearchProjectGeneration, + keywordGeneration: KeywordResearchKeywordGeneration + ) async throws -> KeywordResearchRankingObservationSnapshot { + let target = try await modelStore.read { modelContext in + try Self.requireTarget( + projectGeneration: projectGeneration, + keywordGeneration: keywordGeneration, + in: modelContext + ) + } + try Task.checkCancellation() + + let request = RankingRefreshRequest( + identityKey: keywordGeneration.id.uuidString.lowercased(), + queryKey: target.queryKey, + term: target.term, + storefront: target.storefront, + platform: target.platform + ) + let pageResult: RankingRefreshPageResult + do { + pageResult = try await rankingCoordinator.fetchRankingPage( + for: request, + limit: Self.persistedResultLimit + ) + } catch is CancellationError { + throw CancellationError() + } catch let error as URLError where error.code == .cancelled { + throw CancellationError() + } catch { + throw OpenASOError.map(error) + } + try Task.checkCancellation() + + let commit = try await modelStore.write { modelContext in + try Task.checkCancellation() + let currentTarget = try Self.requireTarget( + projectGeneration: projectGeneration, + keywordGeneration: keywordGeneration, + in: modelContext + ) + guard currentTarget == target, + pageResult.request.queryKey == currentTarget.queryKey + else { + throw KeywordResearchProjectStoreError.staleKeywordRevision( + keywordGeneration.id + ) + } + + let query = try Self.requireQuery(for: currentTarget, in: modelContext) + let persisted = try rankingCoordinator.persistSharedRankingObservation( + pageResult, + query: query, + in: modelContext + ) + try rankingCoordinator.rebuildDerivedStats( + forQueryKey: currentTarget.queryKey, + in: modelContext + ) + let snapshot = Self.snapshot( + persisted.observation, + projectGeneration: projectGeneration, + keywordGeneration: keywordGeneration + ) + try Task.checkCancellation() + return CommitResult( + snapshot: snapshot, + shouldScheduleMetadataEnrichment: persisted.appliedIncomingPage + ) + } + + // A successful write return is the commit point. Do not report + // cancellation after data is durable: doing so would suppress the + // post-commit enrichment and make an equal retry a permanent no-op. + if commit.shouldScheduleMetadataEnrichment { + rankingCoordinator.scheduleTopRankingMetadataEnrichment(for: pageResult) + } + return commit.snapshot + } +} + +private extension KeywordResearchRankingWorkflow { + struct Target: Equatable, Sendable { + let queryKey: String + let term: String + let storefront: String + let platform: AppPlatform + } + + struct CommitResult: Sendable { + let snapshot: KeywordResearchRankingObservationSnapshot + let shouldScheduleMetadataEnrichment: Bool + } + + static func requireTarget( + projectGeneration: KeywordResearchProjectGeneration, + keywordGeneration: KeywordResearchKeywordGeneration, + in modelContext: ModelContext + ) throws -> Target { + let projectID = projectGeneration.id + var projectDescriptor = FetchDescriptor( + predicate: #Predicate { project in + project.id == projectID + } + ) + projectDescriptor.fetchLimit = 1 + guard let project = try modelContext.fetch(projectDescriptor).first else { + throw KeywordResearchProjectStoreError.projectNotFound(projectID) + } + guard project.incarnationID == projectGeneration.incarnationID else { + throw KeywordResearchProjectStoreError.staleProjectRevision(projectID) + } + + let keywordID = keywordGeneration.id + var keywordDescriptor = FetchDescriptor( + predicate: #Predicate { keyword in + keyword.id == keywordID + } + ) + keywordDescriptor.fetchLimit = 1 + guard let keyword = try modelContext.fetch(keywordDescriptor).first else { + throw KeywordResearchProjectStoreError.keywordNotFound(keywordID) + } + guard keyword.incarnationID == keywordGeneration.incarnationID else { + throw KeywordResearchProjectStoreError.staleKeywordRevision(keywordID) + } + guard keyword.projectID == project.id else { + throw KeywordResearchProjectStoreError.keywordNotFound(keywordID) + } + + return Target( + queryKey: keyword.queryKey, + term: keyword.term, + storefront: keyword.storefront, + platform: keyword.platform + ) + } + + static func requireQuery( + for target: Target, + in modelContext: ModelContext + ) throws -> KeywordQuery { + let queryKey = target.queryKey + var descriptor = FetchDescriptor( + predicate: #Predicate { query in + query.queryKey == queryKey + } + ) + descriptor.fetchLimit = 1 + guard let query = try modelContext.fetch(descriptor).first, + query.term == target.term, + query.storefront == target.storefront, + query.platform == target.platform + else { + throw OpenASOError.unexpectedResponse + } + return query + } + + static func snapshot( + _ observation: KeywordRankingCrawl, + projectGeneration: KeywordResearchProjectGeneration, + keywordGeneration: KeywordResearchKeywordGeneration + ) -> KeywordResearchRankingObservationSnapshot { + let items = observation.items + .map { + KeywordResearchRankingItemSnapshot( + id: $0.itemKey, + position: $0.position, + appStoreID: $0.appStoreID, + bundleID: $0.bundleID, + name: $0.name, + subtitle: $0.subtitle, + sellerName: $0.sellerName + ) + } + .sorted { + if $0.position != $1.position { + return $0.position < $1.position + } + if $0.appStoreID != $1.appStoreID { + return $0.appStoreID < $1.appStoreID + } + return $0.id < $1.id + } + return KeywordResearchRankingObservationSnapshot( + id: observation.observationKey, + projectGeneration: projectGeneration, + keywordGeneration: keywordGeneration, + queryKey: observation.queryKey, + term: observation.keyword, + storefront: observation.storefront, + platform: observation.platform, + observedAt: observation.observedAt, + observedHour: observation.observedHour, + source: observation.source, + resultCount: observation.resultCount, + submissionCount: observation.submissionCount, + winningCount: observation.winningCount, + confidence: observation.confidenceRaw, + items: items + ) + } +} diff --git a/OpenASO/Services/SearchRanking/RankingRefreshCoordinator.swift b/OpenASO/Services/SearchRanking/RankingRefreshCoordinator.swift index 3560a6b..26b321d 100644 --- a/OpenASO/Services/SearchRanking/RankingRefreshCoordinator.swift +++ b/OpenASO/Services/SearchRanking/RankingRefreshCoordinator.swift @@ -115,6 +115,22 @@ struct RankingRefreshPageResult: Sendable { let confidence: String? } +/// Context-bound result from the synchronous shared-observation transaction. +/// It must never cross an actor or `ModelContext` boundary. +struct RankingObservationPersistenceResult { + let observation: KeywordRankingCrawl + let appliedIncomingPage: Bool +} + +/// Context-bound result from a tracked ranking persistence transaction. +/// The SwiftData snapshot must remain on the owning model executor; the +/// canonical page result is Sendable and may be scheduled after commit. +struct RankingPagePersistenceResult { + let snapshot: TrackedKeywordDailyRanking + let canonicalPageResult: RankingRefreshPageResult + let appliedSharedObservation: Bool +} + struct RankingMetadataEnrichmentRequest: Hashable, Sendable { let appStoreID: Int64 let storefront: String @@ -150,20 +166,31 @@ struct RankingStatsRebuildRequest: Hashable, Sendable { } } +private struct RankingModelContextHasPendingChangesError: LocalizedError, Sendable { + var errorDescription: String? { + "Save or discard pending edits before refreshing keyword rankings." + } +} + final class RankingRefreshCoordinator: Sendable { private let rankingProvider: any SearchRankingProvider private let appCatalogService: AppCatalogService private let analyticsService: AnalyticsService + private let now: @Sendable () -> Date private let refreshTriggerRecorder: (@Sendable (Date) async -> Void)? - private let metadataEnrichmentHandler: (@Sendable ([RankingMetadataEnrichmentRequest]) async -> Void)? + private let metadataEnrichmentScheduler: (@Sendable ([RankingMetadataEnrichmentRequest]) -> Void)? + private let persistenceMutationCheckpoint: (@Sendable () throws -> Void)? @MainActor init( rankingProvider: any SearchRankingProvider, appCatalogService: AppCatalogService, analyticsService: AnalyticsService? = nil, + now: @escaping @Sendable () -> Date = { Date() }, refreshTriggerRecorder: (@Sendable (Date) async -> Void)? = nil, - metadataEnrichmentHandler: (@Sendable ([RankingMetadataEnrichmentRequest]) async -> Void)? = nil + metadataEnrichmentHandler: (@Sendable ([RankingMetadataEnrichmentRequest]) async -> Void)? = nil, + metadataEnrichmentScheduler: (@Sendable ([RankingMetadataEnrichmentRequest]) -> Void)? = nil, + persistenceMutationCheckpoint: (@Sendable () throws -> Void)? = nil ) { self.rankingProvider = rankingProvider self.appCatalogService = appCatalogService @@ -171,20 +198,32 @@ final class RankingRefreshCoordinator: Sendable { settingsStore: AppSettingsStore(defaults: UserDefaults(suiteName: "com.openaso.analytics.noop") ?? .standard), client: NoOpAnalyticsClient() ) + self.now = now self.refreshTriggerRecorder = refreshTriggerRecorder - self.metadataEnrichmentHandler = metadataEnrichmentHandler + self.persistenceMutationCheckpoint = persistenceMutationCheckpoint + if let metadataEnrichmentScheduler { + self.metadataEnrichmentScheduler = metadataEnrichmentScheduler + } else if let metadataEnrichmentHandler { + self.metadataEnrichmentScheduler = { requests in + Task { + await RefreshObservationScope.$runID.withValue(nil) { + await metadataEnrichmentHandler(requests) + } + } + } + } else { + self.metadataEnrichmentScheduler = nil + } } @MainActor func refresh( track: TrackedAppKeyword, - in modelContext: ModelContext, - limit: Int = SearchRankingCrawl.fullKeywordRankingLimit + in modelContext: ModelContext ) async -> Result { await refresh( track: track, in: modelContext, - limit: limit, recordsTrigger: true ) } @@ -193,13 +232,16 @@ final class RankingRefreshCoordinator: Sendable { private func refresh( track: TrackedAppKeyword, in modelContext: ModelContext, - limit: Int, recordsTrigger: Bool, rebuildDerivedStats: Bool = true ) async -> Result { let request = RankingRefreshRequest(track: track) - let pageResult = await refreshPage(for: request, limit: limit, recordsTrigger: recordsTrigger) + let pageResult = await refreshPage( + for: request, + limit: SearchRankingCrawl.fullKeywordRankingLimit, + recordsTrigger: recordsTrigger + ) switch pageResult { case .success(let pageResult): do { @@ -226,47 +268,46 @@ final class RankingRefreshCoordinator: Sendable { await recordRefreshTriggered() } do { - let page = try await rankingProvider.search( - keyword: request.term, - storefrontCode: request.storefront, - platform: request.platform, - limit: limit - ) - return .success(RankingRefreshPageResult( - request: request, - page: page, - searchedAt: .now, - observedHour: nil, - submissionCount: 1, - winningCount: 1, - confidence: "single_source" - )) + return .success(try await fetchRankingPage(for: request, limit: limit)) } catch { return .failure(OpenASOError.map(error)) } } - @MainActor - func makeRankingPageFetcher( + /// Performs only the provider await and produces an immutable page result. + /// Callers that need cancellation semantics should use this throwing API + /// rather than the legacy `Result` wrapper above. + func fetchRankingPage( + for request: RankingRefreshRequest, limit: Int = SearchRankingCrawl.fullKeywordRankingLimit - ) -> @Sendable (RankingRefreshRequest) async -> Result { - let rankingProvider = rankingProvider + ) async throws -> RankingRefreshPageResult { + try Task.checkCancellation() + let providerPage = try await rankingProvider.search( + keyword: request.term, + storefrontCode: request.storefront, + platform: request.platform, + limit: limit + ) + try Task.checkCancellation() + return RankingRefreshPageResult( + request: request, + page: providerPage.canonicalized(limit: limit), + searchedAt: now(), + observedHour: nil, + submissionCount: 1, + winningCount: 1, + confidence: "single_source" + ) + } + + @MainActor + func makeRankingPageFetcher() -> @Sendable (RankingRefreshRequest) async -> Result { + let coordinator = self return { request in do { - let page = try await rankingProvider.search( - keyword: request.term, - storefrontCode: request.storefront, - platform: request.platform, - limit: limit - ) - return .success(RankingRefreshPageResult( - request: request, - page: page, - searchedAt: .now, - observedHour: nil, - submissionCount: 1, - winningCount: 1, - confidence: "single_source" + return .success(try await coordinator.fetchRankingPage( + for: request, + limit: SearchRankingCrawl.fullKeywordRankingLimit )) } catch { return .failure(OpenASOError.map(error)) @@ -282,23 +323,77 @@ final class RankingRefreshCoordinator: Sendable { saveChanges: Bool = true, scheduleMetadataEnrichment: Bool = true ) throws -> TrackedKeywordDailyRanking { + if saveChanges { + guard !modelContext.hasChanges else { + throw RankingModelContextHasPendingChangesError() + } + let committedResult: RankingPagePersistenceResult + do { + committedResult = try persistRankingPageTransaction( + pageResult, + in: modelContext, + rebuildDerivedStats: rebuildDerivedStats + ) + try modelContext.save() + } catch { + // SwiftData's transaction block does not restore a context's + // in-memory graph when its closure throws. Explicit rollback + // prevents a later save from committing a partial crawl. + modelContext.rollback() + throw error + } + if scheduleMetadataEnrichment, committedResult.appliedSharedObservation { + scheduleTopRankingMetadataEnrichment(for: committedResult.canonicalPageResult) + } + return committedResult.snapshot + } + return try persistRankingPageTransaction( + pageResult, + in: modelContext, + rebuildDerivedStats: rebuildDerivedStats + ).snapshot + } + + /// Mutates only the supplied model context. It never saves or starts + /// metadata work, allowing an outer actor-owned transaction to commit + /// before any enrichment side effect is scheduled. + func persistRankingPageTransaction( + _ pageResult: RankingRefreshPageResult, + in modelContext: ModelContext, + rebuildDerivedStats: Bool = true + ) throws -> RankingPagePersistenceResult { guard let track = try fetchTrackedAppKeyword(identityKey: pageResult.request.identityKey, in: modelContext) else { throw OpenASOError.appNotFound } + let requestTerm = pageResult.request.term + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + let trackTerm = track.term + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + let requestStorefront = pageResult.request.storefront + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + let trackStorefront = track.storefront + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + guard pageResult.request.queryKey == track.queryKey, + requestTerm == trackTerm, + requestStorefront == trackStorefront, + pageResult.request.platform == track.platform + else { + throw OpenASOError.unexpectedResponse + } + let canonicalPageResult = pageResult.canonicalized( + limit: SearchRankingCrawl.fullKeywordRankingLimit + ) return try persistRankingPage( - pageResult.page, - searchedAt: pageResult.searchedAt, - observedHour: pageResult.observedHour, - submissionCount: pageResult.submissionCount, - winningCount: pageResult.winningCount, - confidence: pageResult.confidence, + canonicalPageResult, track: track, trackedApp: track.trackedApp, in: modelContext, - rebuildDerivedStats: rebuildDerivedStats, - saveChanges: saveChanges, - scheduleMetadataEnrichment: scheduleMetadataEnrichment + rebuildDerivedStats: rebuildDerivedStats ) } @@ -309,91 +404,74 @@ final class RankingRefreshCoordinator: Sendable { in modelContext: ModelContext, saveChanges: Bool = true ) throws -> PersistentIdentifier? { - guard let track = try fetchTrackedAppKeyword(identityKey: identityKey, in: modelContext) else { - return nil + if saveChanges, modelContext.hasChanges { + throw RankingModelContextHasPendingChangesError() } + do { + guard let track = try fetchTrackedAppKeyword(identityKey: identityKey, in: modelContext) else { + return nil + } - try TrackedKeywordRefreshStatusStore.set( - "Ranking failed to refresh. \(error.localizedDescription)", - domain: .ranking, - for: track, - in: modelContext - ) - if saveChanges { - try modelContext.save() + try TrackedKeywordRefreshStatusStore.set( + "Ranking failed to refresh. \(error.localizedDescription)", + domain: .ranking, + for: track, + in: modelContext + ) + if saveChanges { + try modelContext.save() + } + return track.persistentModelID + } catch { + if saveChanges { + modelContext.rollback() + } + throw error } - return track.persistentModelID } private func persistRankingPage( - _ page: SearchRankingPage, - searchedAt: Date, - observedHour: Int?, - submissionCount: Int, - winningCount: Int, - confidence: String?, + _ pageResult: RankingRefreshPageResult, track: TrackedAppKeyword, trackedApp: TrackedApp, in modelContext: ModelContext, - rebuildDerivedStats: Bool, - saveChanges: Bool, - scheduleMetadataEnrichment: Bool - ) throws -> TrackedKeywordDailyRanking { - let snapshotKey = TrackedKeywordDailyRanking.makeSnapshotKey( - trackIdentityKey: track.identityKey, - searchedAt: searchedAt, - source: page.source - ) - let observationKey = KeywordRankingCrawl.makeObservationKey( - queryKey: track.queryKey, - observedAt: searchedAt, - source: page.source - ) - - let snapshot = try fetchTrackedKeywordDailyRanking( - snapshotKey: snapshotKey, - track: track, - searchedAt: searchedAt, - source: page.source, - in: modelContext - ) ?? TrackedKeywordDailyRanking( + rebuildDerivedStats: Bool + ) throws -> RankingPagePersistenceResult { + let page = pageResult.page + let searchedAt = pageResult.searchedAt + let incomingSnapshotKey = TrackedKeywordDailyRanking.makeSnapshotKey( + trackIdentityKey: track.identityKey, + searchedAt: searchedAt, + source: page.source + ) + let existingSnapshot = try fetchTrackedKeywordDailyRanking( + snapshotKey: incomingSnapshotKey, + track: track, + searchedAt: searchedAt, + source: page.source, + in: modelContext + ) + let sharedResult = try persistSharedRankingObservation( + pageResult, + query: track.query, + in: modelContext + ) + let incomingWinsTrackedSnapshot = existingSnapshot.map { $0.searchedAt < searchedAt } ?? true + let snapshot: TrackedKeywordDailyRanking + if incomingWinsTrackedSnapshot { + snapshot = existingSnapshot ?? TrackedKeywordDailyRanking( rank: RankingMatcher.rank(for: trackedApp, in: page.items), searchedAt: searchedAt, source: page.source, resultCount: page.resultCount, keywordTrack: track ) - let observation = try fetchKeywordRankingCrawl( - observationKey: observationKey, - queryKey: track.queryKey, - observedAt: searchedAt, - source: page.source, - in: modelContext - ) ?? KeywordRankingCrawl( - keyword: track.term, - storefront: track.storefront, - platform: track.platform, - observedAt: searchedAt, - source: page.source, - resultCount: page.resultCount, - query: track.query, - observedHour: observedHour, - submissionCount: submissionCount, - winningCount: winningCount, - confidence: confidence - ) - let isNewSnapshot = snapshot.modelContext == nil - let isNewObservation = observation.modelContext == nil - if isNewSnapshot { modelContext.insert(snapshot) } - if isNewObservation { - modelContext.insert(observation) - } - snapshot.snapshotKey = snapshotKey + snapshot.snapshotKey = incomingSnapshotKey snapshot.trackIdentityKey = track.identityKey snapshot.rank = RankingMatcher.rank(for: trackedApp, in: page.items) snapshot.searchedAt = searchedAt @@ -402,95 +480,185 @@ final class RankingRefreshCoordinator: Sendable { snapshot.errorMessage = nil snapshot.keywordTrack = track - observation.observationKey = observationKey - observation.queryKey = track.queryKey - observation.query = track.query - observation.keyword = track.term.trimmingCharacters(in: .whitespacesAndNewlines) - observation.storefront = track.storefront.lowercased() - observation.platform = track.platform - observation.observedAt = searchedAt - observation.observedHour = observedHour ?? KeywordRankingCrawl.utcHourBucket(for: searchedAt) - observation.source = page.source - observation.resultCount = page.resultCount - observation.submissionCount = submissionCount - observation.winningCount = winningCount - observation.confidenceRaw = confidence - - var catalogCache = try appCatalogService.makeSearchRankingPageCache( - items: page.items, - storefrontCode: track.storefront, - in: modelContext - ) - var ratingCache = try makeRatingPageCache( - items: page.items, - storefront: track.storefront, - observedAt: snapshot.searchedAt, - in: modelContext - ) - for item in page.items { - let storeApp = try appCatalogService.upsertStoreApp( - from: item, - storefrontCode: track.storefront, - rankingSource: page.source, - fetchedAt: searchedAt, - requestedPlatform: track.platform, - in: modelContext, - cache: &catalogCache - ) - upsertStorefrontRating( - from: item, - storefront: track.storefront, - observedAt: snapshot.searchedAt, - source: storefrontRatingSource(for: page.source), - storeApp: storeApp, - in: modelContext, - cache: &ratingCache - ) - upsertRankedResult( from: item, snapshot: snapshot, - snapshotKey: snapshotKey, + snapshotKey: incomingSnapshotKey, in: modelContext ) - upsertObservationItem( - from: item, - observation: observation, - in: modelContext - ) - } - pruneRankedResults(for: snapshot, keeping: page.items.map(\.appStoreID), in: modelContext) - pruneObservationItems(for: observation, keeping: page.items.map(\.appStoreID), in: modelContext) - - if rebuildDerivedStats { - self.rebuildDerivedStats(for: [RankingStatsRebuildRequest(track: track)], in: modelContext) } - - try TrackedKeywordRefreshStatusStore.set( - nil, - domain: .ranking, - for: track, - updatedAt: snapshot.searchedAt, + pruneRankedResults( + for: snapshot, + keeping: page.items.map(\.appStoreID), in: modelContext ) - track.lastRefreshAt = snapshot.searchedAt track.rankingAppCount = page.resultCount if isNewSnapshot { track.snapshots.append(snapshot) } + } else { + snapshot = existingSnapshot! + } - if saveChanges { - try modelContext.save() - } - if scheduleMetadataEnrichment { - scheduleTopRankingMetadataEnrichment( - items: page.items, - storefront: track.storefront, - platform: track.platform - ) - } - return snapshot + try persistenceMutationCheckpoint?() + + if rebuildDerivedStats { + try self.rebuildDerivedStats( + for: [RankingStatsRebuildRequest(track: track)], + in: modelContext + ) + } + + try TrackedKeywordRefreshStatusStore.set( + nil, + domain: .ranking, + for: track, + updatedAt: searchedAt, + in: modelContext + ) + track.lastRefreshAt = max(track.lastRefreshAt ?? .distantPast, searchedAt) + + return RankingPagePersistenceResult( + snapshot: snapshot, + canonicalPageResult: pageResult, + appliedSharedObservation: sharedResult.appliedIncomingPage + ) + } + + /// Upserts the app-independent ranking observation and its shared catalog + /// and rating data. The caller owns the surrounding transaction and save. + /// No tracked-app model is read or written by this primitive. + @discardableResult + func persistSharedRankingObservation( + _ pageResult: RankingRefreshPageResult, + query: KeywordQuery, + in modelContext: ModelContext + ) throws -> RankingObservationPersistenceResult { + let pageResult = pageResult.canonicalized( + limit: SearchRankingCrawl.fullKeywordRankingLimit + ) + let requestTerm = pageResult.request.term + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + let queryTerm = query.term + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + let requestStorefront = pageResult.request.storefront + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + let queryStorefront = query.storefront + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + guard query.queryKey == pageResult.request.queryKey, + requestTerm == queryTerm, + requestStorefront == queryStorefront, + pageResult.request.platform == query.platform + else { + throw OpenASOError.unexpectedResponse + } + + let observationKey = KeywordRankingCrawl.makeObservationKey( + queryKey: pageResult.request.queryKey, + observedAt: pageResult.searchedAt, + source: pageResult.page.source + ) + let existingObservation = try fetchKeywordRankingCrawl( + observationKey: observationKey, + queryKey: pageResult.request.queryKey, + observedAt: pageResult.searchedAt, + source: pageResult.page.source, + in: modelContext + ) + if let existingObservation, + existingObservation.observedAt >= pageResult.searchedAt { + return RankingObservationPersistenceResult( + observation: existingObservation, + appliedIncomingPage: false + ) + } + + let observation = existingObservation ?? KeywordRankingCrawl( + keyword: pageResult.request.term, + storefront: pageResult.request.storefront, + platform: pageResult.request.platform, + observedAt: pageResult.searchedAt, + source: pageResult.page.source, + resultCount: pageResult.page.resultCount, + query: query, + observedHour: pageResult.observedHour, + submissionCount: pageResult.submissionCount, + winningCount: pageResult.winningCount, + confidence: pageResult.confidence + ) + if observation.modelContext == nil { + modelContext.insert(observation) + } + + observation.observationKey = observationKey + observation.queryKey = pageResult.request.queryKey + observation.query = query + observation.keyword = pageResult.request.term + .trimmingCharacters(in: .whitespacesAndNewlines) + observation.storefront = pageResult.request.storefront + .trimmingCharacters(in: .whitespacesAndNewlines) + .lowercased() + observation.platform = pageResult.request.platform + observation.observedAt = pageResult.searchedAt + observation.observedHour = pageResult.observedHour + ?? KeywordRankingCrawl.utcHourBucket(for: pageResult.searchedAt) + observation.source = pageResult.page.source + observation.resultCount = pageResult.page.resultCount + observation.submissionCount = pageResult.submissionCount + observation.winningCount = pageResult.winningCount + observation.confidenceRaw = pageResult.confidence + + var catalogCache = try appCatalogService.makeSearchRankingPageCache( + items: pageResult.page.items, + storefrontCode: pageResult.request.storefront, + in: modelContext + ) + var ratingCache = try makeRatingPageCache( + items: pageResult.page.items, + storefront: pageResult.request.storefront, + observedAt: pageResult.searchedAt, + in: modelContext + ) + for item in pageResult.page.items { + let storeApp = try appCatalogService.upsertStoreApp( + from: item, + storefrontCode: pageResult.request.storefront, + rankingSource: pageResult.page.source, + fetchedAt: pageResult.searchedAt, + requestedPlatform: pageResult.request.platform, + in: modelContext, + cache: &catalogCache + ) + upsertStorefrontRating( + from: item, + storefront: pageResult.request.storefront, + observedAt: pageResult.searchedAt, + source: storefrontRatingSource(for: pageResult.page.source), + storeApp: storeApp, + in: modelContext, + cache: &ratingCache + ) + upsertObservationItem( + from: item, + observation: observation, + in: modelContext + ) + } + pruneObservationItems( + for: observation, + keeping: pageResult.page.items.map(\.appStoreID), + in: modelContext + ) + + return RankingObservationPersistenceResult( + observation: observation, + appliedIncomingPage: true + ) } func scheduleTopRankingMetadataEnrichment(for pageResult: RankingRefreshPageResult) { @@ -506,7 +674,7 @@ final class RankingRefreshCoordinator: Sendable { storefront: String, platform: AppPlatform ) { - guard let metadataEnrichmentHandler else { return } + guard let metadataEnrichmentScheduler else { return } let requests = Self.topRankingEnrichmentRequests( items: items, storefront: storefront, @@ -514,11 +682,7 @@ final class RankingRefreshCoordinator: Sendable { ) guard !requests.isEmpty else { return } - Task { - await RefreshObservationScope.$runID.withValue(nil) { - await metadataEnrichmentHandler(requests) - } - } + metadataEnrichmentScheduler(requests) } static let metadataEnrichmentTopResultLimit = 20 @@ -768,7 +932,9 @@ final class RankingRefreshCoordinator: Sendable { storefront: normalizedStorefront, ratingDate: ratingDate ) - let snapshot = cache.snapshotsByIdentityKey[snapshotKey] ?? AppDailyRating( + let existingSnapshot = cache.snapshotsByIdentityKey[snapshotKey] + let isNewSnapshot = existingSnapshot == nil + let snapshot = existingSnapshot ?? AppDailyRating( appStoreID: item.appStoreID, storefront: normalizedStorefront, ratingCount: item.ratingCount, @@ -781,24 +947,24 @@ final class RankingRefreshCoordinator: Sendable { source: source, storeApp: storeApp ) - if snapshot.modelContext != nil, observedAt < snapshot.observedAt { + if !isNewSnapshot, observedAt <= snapshot.observedAt { return } - if snapshot.modelContext == nil { + if isNewSnapshot { modelContext.insert(snapshot) cache.snapshotsByIdentityKey[snapshotKey] = snapshot } if snapshot.storeApp !== storeApp { snapshot.storeApp = storeApp } - let snapshotChanged = snapshot.modelContext == nil + let snapshotChanged = isNewSnapshot + || snapshot.observedAt != observedAt || snapshot.ratingCount != item.ratingCount || snapshot.averageRating != item.averageRating || snapshot.ratingDate != ratingDate || snapshot.submissionCount != 1 || snapshot.winningCount != 1 || snapshot.confidenceRaw != "single_source" - || snapshot.observedAt != observedAt || snapshot.source != source if snapshotChanged { snapshot.ratingCount = item.ratingCount @@ -815,7 +981,9 @@ final class RankingRefreshCoordinator: Sendable { appStoreID: item.appStoreID, storefront: normalizedStorefront ) - let latest = cache.latestByIdentityKey[latestKey] ?? LatestAppRating( + let existingLatest = cache.latestByIdentityKey[latestKey] + let isNewLatest = existingLatest == nil + let latest = existingLatest ?? LatestAppRating( appStoreID: item.appStoreID, storefront: normalizedStorefront, ratingCount: item.ratingCount, @@ -828,24 +996,24 @@ final class RankingRefreshCoordinator: Sendable { source: source, storeApp: storeApp ) - if latest.modelContext != nil, observedAt < latest.observedAt { + if !isNewLatest, observedAt <= latest.observedAt { return } - if latest.modelContext == nil { + if isNewLatest { modelContext.insert(latest) cache.latestByIdentityKey[latestKey] = latest } if latest.storeApp !== storeApp { latest.storeApp = storeApp } - let latestChanged = latest.modelContext == nil + let latestChanged = isNewLatest + || latest.observedAt != observedAt || latest.ratingCount != item.ratingCount || latest.averageRating != item.averageRating || latest.ratingDate != ratingDate || latest.submissionCount != 1 || latest.winningCount != 1 || latest.confidenceRaw != "single_source" - || latest.observedAt != observedAt || latest.source != source if latestChanged { latest.ratingCount = item.ratingCount @@ -920,7 +1088,6 @@ final class RankingRefreshCoordinator: Sendable { func refresh( tracks: [TrackedAppKeyword], in modelContext: ModelContext, - limit: Int = SearchRankingCrawl.fullKeywordRankingLimit, analyticsTrigger: String? = nil, progress: (@Sendable (_ completed: Int, _ total: Int, _ failureCount: Int) async -> Void)? = nil ) async -> [RefreshOutcome] { @@ -947,7 +1114,7 @@ final class RankingRefreshCoordinator: Sendable { for requestGroup in requestGroups { let result = await refreshPage( for: requestGroup.providerRequest, - limit: limit, + limit: SearchRankingCrawl.fullKeywordRankingLimit, recordsTrigger: false ) switch result { @@ -976,18 +1143,25 @@ final class RankingRefreshCoordinator: Sendable { )) } catch { let mappedError = OpenASOError.map(error) - do { - try TrackedKeywordRefreshStatusStore.set( - "Ranking failed to refresh. \(mappedError.localizedDescription)", - domain: .ranking, - for: track, - in: modelContext - ) - try modelContext.save() - } catch { + if modelContext.hasChanges { OpenASOLog.refresh.error( - "Failed to persist ranking refresh status: \(String(reflecting: error), privacy: .private(mask: .hash))" + "Skipped ranking failure status because the model context has pending edits." ) + } else { + do { + try TrackedKeywordRefreshStatusStore.set( + "Ranking failed to refresh. \(mappedError.localizedDescription)", + domain: .ranking, + for: track, + in: modelContext + ) + try modelContext.save() + } catch { + modelContext.rollback() + OpenASOLog.refresh.error( + "Failed to persist ranking refresh status: \(String(reflecting: error), privacy: .private(mask: .hash))" + ) + } } outcomes.append(RefreshOutcome( trackID: track.persistentModelID, @@ -1010,18 +1184,25 @@ final class RankingRefreshCoordinator: Sendable { continue } - do { - try TrackedKeywordRefreshStatusStore.set( - "Ranking failed to refresh. \(error.localizedDescription)", - domain: .ranking, - for: track, - in: modelContext - ) - try modelContext.save() - } catch { + if modelContext.hasChanges { OpenASOLog.refresh.error( - "Failed to persist ranking refresh status: \(String(reflecting: error), privacy: .private(mask: .hash))" + "Skipped ranking failure status because the model context has pending edits." ) + } else { + do { + try TrackedKeywordRefreshStatusStore.set( + "Ranking failed to refresh. \(error.localizedDescription)", + domain: .ranking, + for: track, + in: modelContext + ) + try modelContext.save() + } catch { + modelContext.rollback() + OpenASOLog.refresh.error( + "Failed to persist ranking refresh status: \(String(reflecting: error), privacy: .private(mask: .hash))" + ) + } } outcomes.append(RefreshOutcome( trackID: track.persistentModelID, @@ -1037,8 +1218,21 @@ final class RankingRefreshCoordinator: Sendable { } } if !statsRebuildRequests.isEmpty { - rebuildDerivedStats(for: statsRebuildRequests, in: modelContext) - try? modelContext.save() + if modelContext.hasChanges { + OpenASOLog.refresh.error( + "Skipped ranking statistics rebuild because the model context has pending edits." + ) + } else { + do { + try rebuildDerivedStats(for: statsRebuildRequests, in: modelContext) + try modelContext.save() + } catch { + modelContext.rollback() + OpenASOLog.refresh.error( + "Failed to rebuild ranking statistics: \(String(reflecting: error), privacy: .private(mask: .hash))" + ) + } + } } if let analyticsTrigger { await captureKeywordRefreshCompleted( @@ -1052,8 +1246,7 @@ final class RankingRefreshCoordinator: Sendable { @MainActor func refreshStaleTracks( - in modelContext: ModelContext, - limit: Int = SearchRankingCrawl.fullKeywordRankingLimit + in modelContext: ModelContext ) async -> [RefreshOutcome] { let descriptor = FetchDescriptor() let tracks = (try? modelContext.fetch(descriptor)) ?? [] @@ -1064,7 +1257,6 @@ final class RankingRefreshCoordinator: Sendable { return await refresh( tracks: staleTracks, in: modelContext, - limit: limit, analyticsTrigger: "daily_refresh", progress: nil ) @@ -1092,17 +1284,27 @@ final class RankingRefreshCoordinator: Sendable { func rebuildDerivedStats( for requests: some Sequence, in modelContext: ModelContext - ) { - let requests = Set(requests) - for request in requests { - rebuildAppKeywordStats(queryKey: request.queryKey, in: modelContext) + ) throws { + let queryKeys = Set(requests.map(\.queryKey)) + for queryKey in queryKeys { + try rebuildAppKeywordStats(queryKey: queryKey, in: modelContext) } } - private func rebuildAppKeywordStats(queryKey: String, in modelContext: ModelContext) { - let metrics = try? fetchKeywordMetrics(queryKey: queryKey, in: modelContext) - let observations = (try? fetchKeywordRankingCrawls(queryKey: queryKey, in: modelContext)) ?? [] - let existingStats = (try? fetchAppKeywordStats(queryKey: queryKey, in: modelContext)) ?? [] + func rebuildDerivedStats( + forQueryKey queryKey: String, + in modelContext: ModelContext + ) throws { + try rebuildAppKeywordStats(queryKey: queryKey, in: modelContext) + } + + private func rebuildAppKeywordStats( + queryKey: String, + in modelContext: ModelContext + ) throws { + let metrics = try fetchKeywordMetrics(queryKey: queryKey, in: modelContext) + let observations = try fetchKeywordRankingCrawls(queryKey: queryKey, in: modelContext) + let existingStats = try fetchAppKeywordStats(queryKey: queryKey, in: modelContext) struct KeywordAggregate { var appStoreID: Int64 @@ -1118,7 +1320,13 @@ final class RankingRefreshCoordinator: Sendable { } var aggregates: [Int64: KeywordAggregate] = [:] - for observation in observations.sorted(by: { $0.observedAt < $1.observedAt }) { + let orderedObservations = observations.sorted { lhs, rhs in + if lhs.observedAt != rhs.observedAt { + return lhs.observedAt < rhs.observedAt + } + return lhs.observationKey < rhs.observationKey + } + for observation in orderedObservations { for item in observation.items { if var aggregate = aggregates[item.appStoreID] { aggregate.bestRank = min(aggregate.bestRank, item.position) @@ -1306,3 +1514,116 @@ struct RefreshOutcome { let searchedAt: Date? let error: OpenASOError? } + +extension RankingRefreshPageResult { + func canonicalized(limit: Int) -> RankingRefreshPageResult { + RankingRefreshPageResult( + request: request, + page: page.canonicalized(limit: limit), + searchedAt: searchedAt, + observedHour: observedHour, + submissionCount: submissionCount, + winningCount: winningCount, + confidence: confidence + ) + } +} + +extension SearchRankingPage { + func canonicalized(limit: Int) -> SearchRankingPage { + let boundedLimit = max(0, limit) + var seenAppStoreIDs = Set() + let canonicalItems = items + .sorted { lhs, rhs in + if lhs.position != rhs.position { + return lhs.position < rhs.position + } + if lhs.appStoreID != rhs.appStoreID { + return lhs.appStoreID < rhs.appStoreID + } + return Self.canonicalTiePrecedes(lhs, rhs) + } + .compactMap { item in + seenAppStoreIDs.insert(item.appStoreID).inserted + ? item + : nil + } + .prefix(boundedLimit) + return SearchRankingPage(items: Array(canonicalItems), source: source) + } + + private static func canonicalTiePrecedes( + _ lhs: SearchRankingItem, + _ rhs: SearchRankingItem + ) -> Bool { + let comparisons: [ComparisonResult] = [ + compare(lhs.bundleID, rhs.bundleID), + compare(lhs.name, rhs.name), + compare(lhs.subtitle, rhs.subtitle), + compare(lhs.sellerName, rhs.sellerName), + compare(lhs.iconURLString, rhs.iconURLString), + compare( + lhs.releaseDate.map { $0.timeIntervalSinceReferenceDate.bitPattern }, + rhs.releaseDate.map { $0.timeIntervalSinceReferenceDate.bitPattern } + ), + compare( + lhs.currentVersionReleaseDate.map { $0.timeIntervalSinceReferenceDate.bitPattern }, + rhs.currentVersionReleaseDate.map { $0.timeIntervalSinceReferenceDate.bitPattern } + ), + compare(lhs.version, rhs.version), + compare(lhs.primaryGenreID, rhs.primaryGenreID), + compare(lhs.primaryGenreName, rhs.primaryGenreName), + compare(lhs.descriptionText, rhs.descriptionText), + compare(lhs.releaseNotes, rhs.releaseNotes), + compare(lhs.supportedLanguageCodes, rhs.supportedLanguageCodes), + compare(lhs.screenshotURLs, rhs.screenshotURLs), + compare(lhs.ipadScreenshotURLs, rhs.ipadScreenshotURLs), + compare(lhs.appletvScreenshotURLs, rhs.appletvScreenshotURLs), + compare(lhs.ratingCount, rhs.ratingCount), + compare( + lhs.averageRating.map(\.bitPattern), + rhs.averageRating.map(\.bitPattern) + ), + compare(lhs.platform.rawValue, rhs.platform.rawValue), + ] + return comparisons.first { $0 != .orderedSame } == .orderedAscending + } + + private static func compare( + _ lhs: Value, + _ rhs: Value + ) -> ComparisonResult { + if lhs < rhs { return .orderedAscending } + if lhs > rhs { return .orderedDescending } + return .orderedSame + } + + private static func compare( + _ lhs: Value?, + _ rhs: Value? + ) -> ComparisonResult { + switch (lhs, rhs) { + case (.none, .none): + return .orderedSame + case (.none, .some): + return .orderedAscending + case (.some, .none): + return .orderedDescending + case (.some(let lhs), .some(let rhs)): + return compare(lhs, rhs) + } + } + + private static func compare( + _ lhs: [Value], + _ rhs: [Value] + ) -> ComparisonResult { + for (lhsValue, rhsValue) in zip(lhs, rhs) { + let comparison = compare(lhsValue, rhsValue) + if comparison != .orderedSame { + return comparison + } + } + return compare(lhs.count, rhs.count) + } +} diff --git a/OpenASO/Services/Storefront/AppCatalogService.swift b/OpenASO/Services/Storefront/AppCatalogService.swift index ee83ea6..2b7fd81 100644 --- a/OpenASO/Services/Storefront/AppCatalogService.swift +++ b/OpenASO/Services/Storefront/AppCatalogService.swift @@ -160,6 +160,13 @@ final class AppCatalogService: Sendable { cache: inout SearchRankingPageCache ) throws -> StoreApp { let normalizedStorefront = normalizedStorefrontCode(storefrontCode) + let storefrontMetadataIdentityKey = AppStorefrontMetadata.makeIdentityKey( + appStoreID: item.appStoreID, + storefront: normalizedStorefront + ) + let storefrontMetadataFetchedAt = cache.storefrontMetadataByIdentityKey[ + storefrontMetadataIdentityKey + ]?.lastFetchedAt let metadataSource = metadataSource(for: rankingSource) let storeApp: StoreApp let isExisting: Bool @@ -192,27 +199,27 @@ final class AppCatalogService: Sendable { isExisting = false } - if !isExisting || fetchedAt >= storeApp.lastMetadataRefreshAt { - update( - storeApp, - storefront: normalizedStorefront, - bundleID: item.bundleID, - name: item.name, - subtitle: item.subtitle, - sellerName: item.sellerName, - iconURLString: item.iconURLString, - supportedLanguageCodes: item.supportedLanguageCodes, - supportedLanguageCodesSource: metadataSource, - supportedLanguageCodesFetchedAt: fetchedAt, - releaseDate: item.releaseDate, - currentVersionReleaseDate: item.currentVersionReleaseDate, - version: item.version, - primaryGenreID: item.primaryGenreID, - primaryGenreName: item.primaryGenreName, - defaultPlatform: isExisting ? nil : requestedPlatform, - fetchedAt: fetchedAt - ) - } + update( + storeApp, + storefront: normalizedStorefront, + bundleID: item.bundleID, + name: item.name, + subtitle: item.subtitle, + sellerName: item.sellerName, + iconURLString: item.iconURLString, + supportedLanguageCodes: item.supportedLanguageCodes, + supportedLanguageCodesSource: metadataSource, + supportedLanguageCodesFetchedAt: fetchedAt, + releaseDate: item.releaseDate, + currentVersionReleaseDate: item.currentVersionReleaseDate, + version: item.version, + primaryGenreID: item.primaryGenreID, + primaryGenreName: item.primaryGenreName, + defaultPlatform: isExisting ? nil : requestedPlatform, + fetchedAt: fetchedAt, + storefrontMetadataFetchedAt: storefrontMetadataFetchedAt, + allowsNoncanonicalGlobalEvidence: true + ) try upsertStorefrontMetadata( from: item, storefront: normalizedStorefront, @@ -464,40 +471,56 @@ final class AppCatalogService: Sendable { primaryGenreName: String?, defaultPlatform: AppPlatform?, allowsNoncanonicalFallback: Bool = true, - fetchedAt: Date = .now + fetchedAt: Date = .now, + storefrontMetadataFetchedAt: Date? = nil, + allowsNoncanonicalGlobalEvidence: Bool = false ) { var changed = false + let acceptsGlobalEvidence = fetchedAt >= storeApp.lastMetadataRefreshAt + let acceptsStorefrontEvidence = fetchedAt >= (storefrontMetadataFetchedAt ?? .distantPast) + let hasNewerCrossSourceGlobalEvidence = fetchedAt < storeApp.lastMetadataRefreshAt + && storeApp.supportedLanguageCodesSource != supportedLanguageCodesSource + let acceptsCanonicalEvidence = acceptsStorefrontEvidence + && !hasNewerCrossSourceGlobalEvidence let isCanonicalStorefront = isCanonicalStorefront( storeApp: storeApp, storefront: storefront ) + let acceptsGlobalFields = acceptsGlobalEvidence + && (isCanonicalStorefront || allowsNoncanonicalGlobalEvidence) - if isCanonicalStorefront, let bundleID, !bundleID.isEmpty { + if acceptsGlobalFields, let bundleID, !bundleID.isEmpty { changed = assignIfChanged(storeApp, \.bundleID, bundleID) || changed } let canonicalNameIsMissing = nonEmpty(storeApp.name) == nil - if isCanonicalStorefront || (allowsNoncanonicalFallback && canonicalNameIsMissing) { + if (isCanonicalStorefront && acceptsCanonicalEvidence) + || (allowsNoncanonicalFallback && canonicalNameIsMissing) + { changed = assignIfChanged(storeApp, \.name, name) || changed } let canonicalSubtitleIsMissing = nonEmpty(storeApp.subtitle) == nil - if (isCanonicalStorefront || (allowsNoncanonicalFallback && canonicalSubtitleIsMissing)), + if ((isCanonicalStorefront && acceptsCanonicalEvidence) + || (allowsNoncanonicalFallback && canonicalSubtitleIsMissing)), let subtitle, !subtitle.isEmpty { changed = assignIfChanged(storeApp, \.subtitle, subtitle) || changed } - if isCanonicalStorefront, let sellerName, !sellerName.isEmpty { + if acceptsGlobalFields, let sellerName, !sellerName.isEmpty { changed = assignIfChanged(storeApp, \.sellerName, sellerName) || changed } - if isCanonicalStorefront, let iconURLString, !iconURLString.isEmpty { + if isCanonicalStorefront, + acceptsCanonicalEvidence, + let iconURLString, + !iconURLString.isEmpty { changed = assignIfChanged(storeApp, \.iconURLString, iconURLString) || changed } let normalizedLanguages = normalizedLanguageCodes(supportedLanguageCodes) - if isCanonicalStorefront, !normalizedLanguages.isEmpty { + if acceptsGlobalFields, !normalizedLanguages.isEmpty { let codesChanged = assignIfChanged( storeApp, \.supportedLanguageCodes, @@ -518,30 +541,30 @@ final class AppCatalogService: Sendable { } } - if isCanonicalStorefront, let releaseDate { + if acceptsGlobalFields, let releaseDate { changed = assignIfChanged(storeApp, \.releaseDate, releaseDate) || changed } - if isCanonicalStorefront, let currentVersionReleaseDate { + if acceptsGlobalFields, let currentVersionReleaseDate { changed = assignIfChanged(storeApp, \.currentVersionReleaseDate, currentVersionReleaseDate) || changed } - if isCanonicalStorefront, let version, !version.isEmpty { + if acceptsGlobalFields, let version, !version.isEmpty { changed = assignIfChanged(storeApp, \.version, version) || changed } - if isCanonicalStorefront, let primaryGenreID { + if acceptsGlobalFields, let primaryGenreID { changed = assignIfChanged(storeApp, \.primaryGenreID, primaryGenreID) || changed } - if isCanonicalStorefront, let primaryGenreName, !primaryGenreName.isEmpty { + if acceptsGlobalFields, let primaryGenreName, !primaryGenreName.isEmpty { changed = assignIfChanged(storeApp, \.primaryGenreName, primaryGenreName) || changed } - if isCanonicalStorefront, let defaultPlatform { + if acceptsGlobalFields, let defaultPlatform { changed = assignIfChanged(storeApp, \.defaultPlatformRaw, defaultPlatform.rawValue) || changed } - if isCanonicalStorefront && (changed || storeApp.lastMetadataRefreshAt != fetchedAt) { + if acceptsGlobalFields && (changed || storeApp.lastMetadataRefreshAt != fetchedAt) { storeApp.lastMetadataRefreshAt = fetchedAt } } @@ -593,7 +616,6 @@ final class AppCatalogService: Sendable { cache.storefrontMetadataByIdentityKey[identityKey] = metadata isNewMetadata = true } - guard fetchedAt >= metadata.lastFetchedAt else { return } // App Store Web search rows are intentionally sparse. Once a detail diff --git a/OpenASOTests/KeywordResearchProjectStoreTests.swift b/OpenASOTests/KeywordResearchProjectStoreTests.swift index e6528ee..95217de 100644 --- a/OpenASOTests/KeywordResearchProjectStoreTests.swift +++ b/OpenASOTests/KeywordResearchProjectStoreTests.swift @@ -692,6 +692,7 @@ struct KeywordResearchProjectStoreTests { ) #expect(appServices.keywordResearchProjectStore === researchStore) + #expect(appServices.keywordResearchRankingWorkflow != nil) let sharedProject = try await researchStore.createProject(name: "Shared MCP reader") let mcpReader = try #require(mcpService.keywordResearchProjectStore) #expect(try await mcpReader.listProjects(offset: 0, limit: 50) == [sharedProject]) @@ -708,6 +709,7 @@ struct KeywordResearchProjectStoreTests { providerRequestGateMode: .disabled ) let autoStore = try #require(autoAppServices.keywordResearchProjectStore) + #expect(autoAppServices.keywordResearchRankingWorkflow != nil) let autoProject = try await autoStore.createProject(name: "Auto-wired app store") #expect(try await autoStore.listProjects() == [autoProject]) diff --git a/OpenASOTests/KeywordResearchRankingWorkflowTests.swift b/OpenASOTests/KeywordResearchRankingWorkflowTests.swift new file mode 100644 index 0000000..e9c96b7 --- /dev/null +++ b/OpenASOTests/KeywordResearchRankingWorkflowTests.swift @@ -0,0 +1,1659 @@ +import Foundation +import SwiftData +import Synchronization +import Testing +@testable import OpenASO + +@MainActor +struct KeywordResearchRankingWorkflowTests { + @Test + func refreshPersistsCanonicalSharedObservationWithoutCreatingTrackedState() async throws { + let searchedAt = utcDate(year: 2026, month: 7, day: 1, hour: 13) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage( + items: [ + rankingItem(position: 2, appStoreID: 20, name: "Zulu", ratingCount: 20), + rankingItem(position: 1, appStoreID: 10, name: "First", ratingCount: 10), + rankingItem(position: 2, appStoreID: 20, name: "Alpha", ratingCount: 21), + rankingItem(position: 9, appStoreID: 10, name: "Late duplicate", ratingCount: 99), + ], + source: .iTunesFallback + )), + ]) + let fixture = try await makeFixture(provider: provider, dates: [searchedAt]) + + let snapshot = try await fixture.workflow.refresh( + projectGeneration: fixture.project.generation, + keywordGeneration: fixture.keyword.generation + ) + + #expect(snapshot.projectGeneration == fixture.project.generation) + #expect(snapshot.keywordGeneration == fixture.keyword.generation) + #expect(snapshot.queryKey == fixture.keyword.queryKey) + #expect(snapshot.term == "launch::planner") + #expect(snapshot.storefront == "gb") + #expect(snapshot.platform == .ipad) + #expect(snapshot.observedAt == searchedAt) + #expect(snapshot.resultCount == 2) + #expect(snapshot.items.map(\.appStoreID) == [10, 20]) + #expect(snapshot.items.map(\.name) == ["First", "Alpha"]) + + let calls = await provider.recordedCalls() + #expect(calls == [ProviderCall( + keyword: "launch::planner", + storefront: "gb", + platform: .ipad, + limit: 200 + )]) + + let state = try await databaseState(in: fixture.backgroundStore) + #expect(state.crawlCount == 1) + #expect(state.observationItemCount == 2) + #expect(state.storeAppIDs == [10, 20]) + #expect(state.metadataCount == 2) + #expect(state.screenshotCount == 2) + #expect(state.latestRatingIDs == [10, 20]) + #expect(state.dailyRatingIDs == [10, 20]) + #expect(state.statsAppStoreIDs == [10, 20]) + #expect(state.trackedCount == 0) + } + + @Test + func canonicalizationIsPermutationStableForDuplicateProviderRows() async throws { + let lower = rankingItem(position: 2, appStoreID: 20, name: "Alpha", ratingCount: 1) + let higher = rankingItem(position: 2, appStoreID: 20, name: "Zulu", ratingCount: 2) + let first = rankingItem(position: 1, appStoreID: 10, name: "First") + let request = RankingRefreshRequest( + identityKey: "research-membership", + queryKey: "opaque::query::key::v1", + term: "opaque::term", + storefront: "gb", + platform: .ipad + ) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage(items: [higher, first, lower], source: .iTunesFallback)), + .page(SearchRankingPage(items: [lower, higher, first], source: .iTunesFallback)), + ]) + let coordinator = RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: AppCatalogService(appResolver: NoOpAppResolver()) + ) + + let forward = try await coordinator.fetchRankingPage(for: request, limit: 200) + let reversed = try await coordinator.fetchRankingPage(for: request, limit: 200) + + #expect(forward.page.items == reversed.page.items) + #expect(forward.page.items.map(\.appStoreID) == [10, 20]) + #expect(forward.page.items.map(\.name) == ["First", "Alpha"]) + } + + @Test + func delayedRetargetRejectsOldMembershipGenerationWithoutLateWrites() async throws { + let provider = GatedResearchRankingProvider() + let fixture = try await makeFixture( + provider: provider, + dates: [utcDate(year: 2026, month: 7, day: 2, hour: 13)] + ) + let task = Task { + try await fixture.workflow.refresh( + projectGeneration: fixture.project.generation, + keywordGeneration: fixture.keyword.generation + ) + } + await provider.waitUntilStarted() + + let projectAfterRemoval = try await fixture.projectStore.removeKeyword( + revision: fixture.keyword.revision, + from: fixture.project.revision + ) + _ = try await fixture.projectStore.addKeyword( + id: fixture.keyword.id, + to: projectAfterRemoval.revision, + term: "retargeted keyword", + storefront: "us", + platform: .iphone + ) + await provider.succeed(SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: 10)], + source: .iTunesFallback + )) + + await #expect(throws: KeywordResearchProjectStoreError.staleKeywordRevision( + fixture.keyword.id + )) { + _ = try await task.value + } + #expect(try await sharedWriteCount(in: fixture.backgroundStore) == 0) + } + + @Test + func delayedProjectReplacementRejectsOldGenerationWithoutLateWrites() async throws { + let provider = GatedResearchRankingProvider() + let fixture = try await makeFixture( + provider: provider, + dates: [utcDate(year: 2026, month: 7, day: 3, hour: 13)] + ) + let task = Task { + try await fixture.workflow.refresh( + projectGeneration: fixture.project.generation, + keywordGeneration: fixture.keyword.generation + ) + } + await provider.waitUntilStarted() + + try await fixture.projectStore.deleteProject(revision: fixture.project.revision) + _ = try await fixture.projectStore.createProject( + id: fixture.project.id, + name: "Replacement project", + defaultStorefront: "gb", + defaultPlatform: .ipad + ) + await provider.succeed(SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: 10)], + source: .iTunesFallback + )) + + await #expect(throws: KeywordResearchProjectStoreError.staleProjectRevision( + fixture.project.id + )) { + _ = try await task.value + } + #expect(try await sharedWriteCount(in: fixture.backgroundStore) == 0) + } + + @Test + func cancellationAfterProviderStartsIsPreservedAndRollsBackAllWrites() async throws { + let provider = GatedResearchRankingProvider() + let metadataRecorder = MetadataEnrichmentRecorder() + let fixture = try await makeFixture( + provider: provider, + dates: [utcDate(year: 2026, month: 7, day: 4, hour: 13)], + metadataEnrichmentScheduler: metadataRecorder.record + ) + let task = Task { + try await fixture.workflow.refresh( + projectGeneration: fixture.project.generation, + keywordGeneration: fixture.keyword.generation + ) + } + await provider.waitUntilStarted() + + task.cancel() + await provider.succeed(SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: 10)], + source: .iTunesFallback + )) + + await #expect(throws: CancellationError.self) { + _ = try await task.value + } + #expect(try await sharedWriteCount(in: fixture.backgroundStore) == 0) + #expect(metadataRecorder.recordedBatches().isEmpty) + } + + @Test + func providerFailureIsReturnedWithoutPersistenceSideEffects() async throws { + let provider = ScriptedResearchRankingProvider(steps: [.failure(.networkUnavailable)]) + let metadataRecorder = MetadataEnrichmentRecorder() + let fixture = try await makeFixture( + provider: provider, + dates: [], + metadataEnrichmentScheduler: metadataRecorder.record + ) + + await #expect(throws: OpenASOError.networkUnavailable) { + _ = try await fixture.workflow.refresh( + projectGeneration: fixture.project.generation, + keywordGeneration: fixture.keyword.generation + ) + } + + #expect(try await sharedWriteCount(in: fixture.backgroundStore) == 0) + #expect(metadataRecorder.recordedBatches().isEmpty) + } + + @Test + func newerSameDayPageReplacesObservationWhileOlderAndEqualPagesAreCompleteNoOps() async throws { + let firstDate = utcDate(year: 2026, month: 7, day: 5, hour: 12) + let newestDate = utcDate(year: 2026, month: 7, day: 5, hour: 13) + let olderDate = utcDate(year: 2026, month: 7, day: 5, hour: 12, minute: 30) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 1, ratingCount: 10), + rankingItem(position: 2, appStoreID: 2, ratingCount: 20), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 2, ratingCount: 21), + rankingItem(position: 2, appStoreID: 3, ratingCount: 30), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 4, ratingCount: 40), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 5, ratingCount: 50), + ], source: .iTunesFallback)), + ]) + let metadataRecorder = MetadataEnrichmentRecorder() + let fixture = try await makeFixture( + provider: provider, + dates: [firstDate, newestDate, olderDate, newestDate], + metadataEnrichmentScheduler: metadataRecorder.record + ) + + _ = try await refresh(fixture) + #expect(metadataRecorder.recordedBatches() == [[ + RankingMetadataEnrichmentRequest(appStoreID: 1, storefront: "gb", platform: .ipad), + RankingMetadataEnrichmentRequest(appStoreID: 2, storefront: "gb", platform: .ipad), + ]]) + let newest = try await refresh(fixture) + #expect(metadataRecorder.recordedBatches() == [ + [ + RankingMetadataEnrichmentRequest(appStoreID: 1, storefront: "gb", platform: .ipad), + RankingMetadataEnrichmentRequest(appStoreID: 2, storefront: "gb", platform: .ipad), + ], + [ + RankingMetadataEnrichmentRequest(appStoreID: 2, storefront: "gb", platform: .ipad), + RankingMetadataEnrichmentRequest(appStoreID: 3, storefront: "gb", platform: .ipad), + ], + ]) + let afterOlder = try await refresh(fixture) + #expect(metadataRecorder.recordedBatches().count == 2) + let afterEqualConflict = try await refresh(fixture) + #expect(metadataRecorder.recordedBatches().count == 2) + + #expect(afterOlder == newest) + #expect(afterEqualConflict == newest) + let observations = try await observationRecords(in: fixture.backgroundStore) + #expect(observations == [ObservationRecord( + observedAt: newestDate, + source: .iTunesFallback, + appStoreIDs: [2, 3], + positions: [1, 2] + )]) + let state = try await databaseState(in: fixture.backgroundStore) + #expect(state.storeAppIDs == [1, 2, 3]) + #expect(state.latestRatingIDs == [1, 2, 3]) + #expect(state.statsAppStoreIDs == [2, 3]) + #expect(!state.storeAppIDs.contains(4)) + #expect(!state.storeAppIDs.contains(5)) + } + + @Test + func differentDayAndSourceCreateDistinctCrawlsAndOlderRatingsCannotReplaceNewerValues() async throws { + let newestFirstDay = utcDate(year: 2026, month: 7, day: 6, hour: 14) + let olderOtherSource = utcDate(year: 2026, month: 7, day: 6, hour: 13) + let secondDay = utcDate(year: 2026, month: 7, day: 7, hour: 14) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 1, ratingCount: 100, averageRating: 4.0), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem(position: 2, appStoreID: 1, ratingCount: 50, averageRating: 2.0), + ], source: .appStoreWeb)), + .page(SearchRankingPage(items: [ + rankingItem(position: 3, appStoreID: 1, ratingCount: 200, averageRating: 4.5), + ], source: .iTunesFallback)), + ]) + let fixture = try await makeFixture( + provider: provider, + dates: [newestFirstDay, olderOtherSource, secondDay] + ) + + _ = try await refresh(fixture) + _ = try await refresh(fixture) + let firstDayRating = try await latestRatingRecord(in: fixture.backgroundStore) + #expect(firstDayRating == RatingRecord( + observedAt: newestFirstDay, + ratingCount: 100, + averageRating: 4.0 + )) + _ = try await refresh(fixture) + + let observations = try await observationRecords(in: fixture.backgroundStore) + #expect(observations.count == 3) + #expect(Set(observations.map(\.source)) == [.iTunesFallback, .appStoreWeb]) + #expect(try await latestRatingRecord(in: fixture.backgroundStore) == RatingRecord( + observedAt: secondDay, + ratingCount: 200, + averageRating: 4.5 + )) + let stats = try await statsRecords(in: fixture.backgroundStore) + #expect(stats == [StatsRecord( + appStoreID: 1, + bestRank: 1, + latestRank: 3, + observationCount: 3 + )]) + } + + @Test + func delayedSourceUpdateCannotRegressNewerCrossSourceCatalogAndRatingWatermarks() async throws { + let firstITunesObservation = utcDate(year: 2026, month: 7, day: 8, hour: 12) + let newerWebObservation = utcDate(year: 2026, month: 7, day: 8, hour: 13) + let delayedITunesObservation = utcDate(year: 2026, month: 7, day: 8, hour: 12, minute: 30) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 1, name: "Fresh", ratingCount: 100), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 1, name: "Fresh", ratingCount: 100), + ], source: .appStoreWeb)), + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 1, name: "Stale", ratingCount: 50), + ], source: .iTunesFallback)), + ]) + let fixture = try await makeFixture( + provider: provider, + dates: [firstITunesObservation, newerWebObservation, delayedITunesObservation] + ) + + _ = try await refresh(fixture) + _ = try await refresh(fixture) + _ = try await refresh(fixture) + + let record = try await crossSourceWatermarkRecord( + queryKey: fixture.keyword.queryKey, + appStoreID: 1, + storefront: "gb", + in: fixture.backgroundStore + ) + #expect(record == CrossSourceWatermarkRecord( + latestRatingCount: 100, + latestRatingObservedAt: newerWebObservation, + latestRatingSource: .appStorePage, + dailyRatingCount: 100, + dailyRatingObservedAt: newerWebObservation, + dailyRatingSource: .appStorePage, + storeAppName: "Fresh", + storeAppMetadataObservedAt: newerWebObservation, + storeAppLanguageSource: .appStoreWebSearch, + storefrontMetadataName: "Stale", + storefrontMetadataObservedAt: delayedITunesObservation, + storefrontMetadataSource: .iTunesSearch, + iTunesCrawlName: "Stale", + iTunesCrawlObservedAt: delayedITunesObservation + )) + } + + @Test + func defaultStorefrontCanonicalMetadataUsesItsOwnWatermarkWithoutRegressingGlobalFields() async throws { + let defaultStorefrontObservation = utcDate(year: 2026, month: 7, day: 9, hour: 12) + let newerNonDefaultObservation = utcDate(year: 2026, month: 7, day: 9, hour: 13) + let newerDefaultStorefrontObservation = utcDate( + year: 2026, + month: 7, + day: 9, + hour: 12, + minute: 30 + ) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage(items: [ + rankingItem( + position: 1, + appStoreID: 77, + name: "Default t12", + bundleID: "example.default.t12", + sellerName: "Default Seller t12", + iconURLString: "https://example.com/default-t12.png", + version: "1.0", + supportedLanguageCodes: ["EN"] + ), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem( + position: 1, + appStoreID: 77, + name: "Non-default t13", + bundleID: "example.global.t13", + sellerName: "Global Seller t13", + iconURLString: "https://example.com/non-default-t13.png", + version: "2.0", + supportedLanguageCodes: ["EN"] + ), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem( + position: 1, + appStoreID: 77, + name: "Default t12.5", + bundleID: "example.stale.t12-5", + sellerName: "Stale Seller t12.5", + iconURLString: "https://example.com/default-t12-5.png", + version: "1.5", + supportedLanguageCodes: ["FR"] + ), + ], source: .iTunesFallback)), + ]) + let fixture = try await makeFixture( + provider: provider, + dates: [ + defaultStorefrontObservation, + newerNonDefaultObservation, + newerDefaultStorefrontObservation, + ], + term: "catalog watermark" + ) + let nonDefaultAddition = try await fixture.projectStore.addKeyword( + to: fixture.project.revision, + term: fixture.keyword.term, + storefront: "us", + platform: .ipad + ) + let projectGeneration = nonDefaultAddition.project.generation + + _ = try await fixture.workflow.refresh( + projectGeneration: projectGeneration, + keywordGeneration: fixture.keyword.generation + ) + _ = try await fixture.workflow.refresh( + projectGeneration: projectGeneration, + keywordGeneration: nonDefaultAddition.keyword.generation + ) + _ = try await fixture.workflow.refresh( + projectGeneration: projectGeneration, + keywordGeneration: fixture.keyword.generation + ) + + let record = try await catalogWatermarkRecord( + appStoreID: 77, + defaultStorefront: "gb", + nonDefaultStorefront: "us", + in: fixture.backgroundStore + ) + #expect(record == CatalogWatermarkRecord( + defaultStorefront: "gb", + name: "Default t12.5", + iconURLString: "https://example.com/default-t12-5.png", + bundleID: "example.global.t13", + sellerName: "Global Seller t13", + version: "2.0", + globalMetadataObservedAt: newerNonDefaultObservation, + supportedLanguageCodes: ["EN"], + supportedLanguageCodesObservedAt: newerNonDefaultObservation, + defaultMetadataName: "Default t12.5", + defaultMetadataIconURLString: "https://example.com/default-t12-5.png", + defaultMetadataObservedAt: newerDefaultStorefrontObservation, + nonDefaultMetadataName: "Non-default t13", + nonDefaultMetadataIconURLString: "https://example.com/non-default-t13.png", + nonDefaultMetadataObservedAt: newerNonDefaultObservation + )) + } + + @Test + func persistedResearchAlwaysRequestsAndCapsTheFullResultWindow() async throws { + let searchedAt = utcDate(year: 2026, month: 7, day: 8, hour: 13) + let oversizedItems = (1...201).map { + rankingItem(position: $0, appStoreID: Int64($0)) + } + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage(items: oversizedItems, source: .iTunesFallback)), + ]) + let fixture = try await makeFixture( + provider: provider, + dates: [searchedAt] + ) + + let snapshot = try await refresh(fixture) + + #expect(snapshot.items.count == 200) + #expect(snapshot.items.last?.appStoreID == 200) + #expect(await provider.recordedCalls().map(\.limit) == [200]) + } + + @Test + func corruptedOpaqueQueryIsRejectedWithoutRewritingOrSharedWrites() async throws { + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: 1)], + source: .iTunesFallback + )), + ]) + let fixture = try await makeFixture( + provider: provider, + dates: [utcDate(year: 2026, month: 7, day: 10, hour: 13)] + ) + try await fixture.backgroundStore.write { modelContext in + let queryKey = fixture.keyword.queryKey + var descriptor = FetchDescriptor( + predicate: #Predicate { query in query.queryKey == queryKey } + ) + descriptor.fetchLimit = 1 + let query = try #require(modelContext.fetch(descriptor).first) + query.term = "corrupted term" + } + + await #expect(throws: OpenASOError.unexpectedResponse) { + _ = try await refresh(fixture) + } + + #expect(try await sharedWriteCount(in: fixture.backgroundStore) == 0) + let queryTerms = try await fixture.backgroundStore.read { modelContext in + try modelContext.fetch(FetchDescriptor()).map(\.term) + } + #expect(queryTerms == ["corrupted term"]) + } + + @Test + func newerResearchCrawlCannotDriveDelayedTrackedPageState() async throws { + let trackedPageAt = utcDate(year: 2026, month: 7, day: 11, hour: 12) + let failureAt = utcDate(year: 2026, month: 7, day: 11, hour: 12, minute: 30) + let researchPageAt = utcDate(year: 2026, month: 7, day: 11, hour: 13) + let previousRefreshAt = utcDate(year: 2026, month: 7, day: 10, hour: 12) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: 1_300)], + source: .iTunesFallback + )), + ]) + let fixture = try await makeFixture(provider: provider, dates: [researchPageAt]) + let trackIdentityKey = try await fixture.backgroundStore.write { modelContext in + let queryKey = fixture.keyword.queryKey + var descriptor = FetchDescriptor( + predicate: #Predicate { query in query.queryKey == queryKey } + ) + descriptor.fetchLimit = 1 + guard let query = try modelContext.fetch(descriptor).first else { + throw OpenASOError.unexpectedResponse + } + let trackedApp = TrackedApp( + appStoreID: 99, + bundleID: "example.99", + name: "Tracked", + sellerName: "Tracked Seller", + defaultPlatform: .ipad, + createdAt: utcDate(year: 2026, month: 6, day: 1, hour: 10) + ) + let track = TrackedAppKeyword( + term: query.term, + storefront: query.storefront, + platform: query.platform, + trackedApp: trackedApp, + query: query, + createdAt: utcDate(year: 2026, month: 6, day: 2, hour: 10) + ) + track.lastRefreshAt = previousRefreshAt + track.rankingAppCount = 17 + trackedApp.keywordTracks.append(track) + modelContext.insert(trackedApp) + modelContext.insert(track) + try TrackedKeywordRefreshStatusStore.set( + "t12.5 ranking failure", + domain: .ranking, + for: track, + updatedAt: failureAt, + in: modelContext + ) + return track.identityKey + } + + _ = try await refresh(fixture) + + let trackedPage = RankingRefreshPageResult( + request: RankingRefreshRequest( + identityKey: trackIdentityKey, + queryKey: fixture.keyword.queryKey, + term: fixture.keyword.term, + storefront: fixture.keyword.storefront, + platform: fixture.keyword.platform + ), + page: SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 1_200), + rankingItem(position: 2, appStoreID: 99), + ], source: .iTunesFallback), + searchedAt: trackedPageAt, + observedHour: nil, + submissionCount: 1, + winningCount: 1, + confidence: "tracked_page" + ) + let appliedSharedObservation = try await fixture.backgroundStore.write { modelContext in + try fixture.coordinator.persistRankingPageTransaction( + trackedPage, + in: modelContext + ).appliedSharedObservation + } + + #expect(!appliedSharedObservation) + let state = try await trackedSharedRaceState( + trackIdentityKey: trackIdentityKey, + queryKey: fixture.keyword.queryKey, + in: fixture.backgroundStore + ) + #expect(state.sharedObservedAt == researchPageAt) + #expect(state.sharedAppStoreIDs == [1_300]) + #expect(state.snapshotSearchedAt == trackedPageAt) + #expect(state.snapshotRank == 2) + #expect(state.snapshotResultCount == 2) + #expect(state.snapshotAppStoreIDs == [1_200, 99]) + #expect(state.snapshotPositions == [1, 2]) + #expect(state.trackLastRefreshAt == trackedPageAt) + #expect(state.trackRankingAppCount == 2) + #expect(state.rankingStatusMessage == "t12.5 ranking failure") + #expect(state.rankingStatusUpdatedAt == failureAt) + } + + @Test + func researchRefreshLeavesFullySeededTrackedGraphUnchanged() async throws { + let researchPageAt = utcDate(year: 2026, month: 7, day: 12, hour: 13) + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: 10)], + source: .iTunesFallback + )), + ]) + let fixture = try await makeFixture(provider: provider, dates: [researchPageAt]) + try await seedFullyPopulatedTrackedGraph( + queryKey: fixture.keyword.queryKey, + in: fixture.backgroundStore + ) + let before = try await trackedGraphState(in: fixture.backgroundStore) + + _ = try await refresh(fixture) + + let after = try await trackedGraphState(in: fixture.backgroundStore) + #expect(after == before) + } + + @Test + func trackedPersistenceDefersMetadataUntilExplicitPostCommitScheduling() async throws { + let container = try ModelContainerFactory.makeModelContainer(isStoredInMemoryOnly: true) + let modelContext = ModelContext(container) + let trackedApp = TrackedApp( + appStoreID: 99, + bundleID: "example.99", + name: "Tracked", + sellerName: "Seller", + defaultPlatform: .ipad + ) + let query = try KeywordQuery.fetchOrInsert( + term: "deferred metadata", + storefront: "gb", + platform: .ipad, + in: modelContext + ) + let track = TrackedAppKeyword( + term: query.term, + storefront: query.storefront, + platform: query.platform, + trackedApp: trackedApp, + query: query + ) + trackedApp.keywordTracks.append(track) + modelContext.insert(trackedApp) + modelContext.insert(track) + try modelContext.save() + + let recorder = MetadataEnrichmentRecorder() + let coordinator = RankingRefreshCoordinator( + rankingProvider: ScriptedResearchRankingProvider(steps: []), + appCatalogService: AppCatalogService(appResolver: NoOpAppResolver()), + metadataEnrichmentScheduler: recorder.record + ) + let pageResult = RankingRefreshPageResult( + request: RankingRefreshRequest(track: track), + page: SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: 10)], + source: .iTunesFallback + ), + searchedAt: utcDate(year: 2026, month: 7, day: 13, hour: 13), + observedHour: nil, + submissionCount: 1, + winningCount: 1, + confidence: "single_source" + ) + + let persistence = try coordinator.persistRankingPageTransaction( + pageResult, + in: modelContext + ) + #expect(recorder.recordedBatches().isEmpty) + + try modelContext.save() + #expect(recorder.recordedBatches().isEmpty) + + coordinator.scheduleTopRankingMetadataEnrichment( + for: persistence.canonicalPageResult + ) + let requests = try #require(recorder.recordedBatches().first) + #expect(requests == [RankingMetadataEnrichmentRequest( + appStoreID: 10, + storefront: "gb", + platform: .ipad + )]) + } + + @Test + func equalTimestampTrackedRetryCannotMutateSharedOrTrackedRows() async throws { + let searchedAt = utcDate(year: 2026, month: 7, day: 11, hour: 13) + let container = try ModelContainerFactory.makeModelContainer(isStoredInMemoryOnly: true) + let modelContext = ModelContext(container) + let trackedApp = TrackedApp( + appStoreID: 99, + bundleID: "example.99", + name: "Tracked", + sellerName: "Seller", + defaultPlatform: .iphone + ) + let query = try KeywordQuery.fetchOrInsert( + term: "tracked term", + storefront: "us", + platform: .iphone, + in: modelContext + ) + let track = TrackedAppKeyword( + term: query.term, + storefront: query.storefront, + platform: query.platform, + trackedApp: trackedApp, + query: query + ) + trackedApp.keywordTracks.append(track) + modelContext.insert(trackedApp) + modelContext.insert(track) + try modelContext.save() + let provider = ScriptedResearchRankingProvider(steps: [ + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 99), + rankingItem(position: 2, appStoreID: 1), + ], source: .iTunesFallback)), + .page(SearchRankingPage(items: [ + rankingItem(position: 1, appStoreID: 2), + ], source: .iTunesFallback)), + ]) + let coordinator = RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: AppCatalogService(appResolver: NoOpAppResolver()), + now: DateSequence([searchedAt, searchedAt]).next + ) + + let first = await coordinator.refresh(track: track, in: modelContext) + let second = await coordinator.refresh(track: track, in: modelContext) + + guard case .success(let firstSnapshot) = first, + case .success(let secondSnapshot) = second + else { + Issue.record("Expected both tracked refreshes to succeed") + return + } + #expect(firstSnapshot.persistentModelID == secondSnapshot.persistentModelID) + #expect(secondSnapshot.rank == 1) + #expect(secondSnapshot.topResults.map(\.appStoreID).sorted() == [1, 99]) + #expect(try modelContext.fetch(FetchDescriptor()).count == 1) + #expect(try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted() == [1, 99]) + #expect(try modelContext.fetch(FetchDescriptor()).count == 1) + #expect(try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted() == [1, 99]) + #expect(try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted() == [1, 99]) + #expect(try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted() == [1, 99]) + } +} + +private extension KeywordResearchRankingWorkflowTests { + func makeFixture( + provider: any SearchRankingProvider, + dates: [Date], + term: String = "launch::planner", + metadataEnrichmentScheduler: (@Sendable ([RankingMetadataEnrichmentRequest]) -> Void)? = nil + ) async throws -> WorkflowFixture { + let container = try ModelContainerFactory.makeModelContainer(isStoredInMemoryOnly: true) + let backgroundStore = BackgroundModelStore(modelContainer: container) + let createdAt = utcDate(year: 2026, month: 6, day: 1, hour: 12) + let projectStore = KeywordResearchProjectStore( + backgroundModelStore: backgroundStore, + now: { createdAt } + ) + let project = try await projectStore.createProject( + name: "Pre-live research", + defaultStorefront: "gb", + defaultPlatform: .ipad + ) + let addition = try await projectStore.addKeyword( + to: project.revision, + term: term, + storefront: "gb", + platform: .ipad + ) + let dateSequence = DateSequence(dates) + let coordinator = RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: AppCatalogService(appResolver: NoOpAppResolver()), + now: dateSequence.next, + metadataEnrichmentScheduler: metadataEnrichmentScheduler + ) + return WorkflowFixture( + container: container, + backgroundStore: backgroundStore, + projectStore: projectStore, + coordinator: coordinator, + workflow: KeywordResearchRankingWorkflow( + backgroundModelStore: backgroundStore, + rankingCoordinator: coordinator + ), + project: addition.project, + keyword: addition.keyword + ) + } + + func refresh( + _ fixture: WorkflowFixture + ) async throws -> KeywordResearchRankingObservationSnapshot { + try await fixture.workflow.refresh( + projectGeneration: fixture.project.generation, + keywordGeneration: fixture.keyword.generation + ) + } +} + +private struct WorkflowFixture { + let container: ModelContainer + let backgroundStore: BackgroundModelStore + let projectStore: KeywordResearchProjectStore + let coordinator: RankingRefreshCoordinator + let workflow: KeywordResearchRankingWorkflow + let project: KeywordResearchProjectSnapshot + let keyword: KeywordResearchKeywordSnapshot +} + +private struct ProviderCall: Equatable, Sendable { + let keyword: String + let storefront: String + let platform: AppPlatform + let limit: Int +} + +private actor ScriptedResearchRankingProvider: SearchRankingProvider { + enum Step: Sendable { + case page(SearchRankingPage) + case failure(OpenASOError) + } + + private var steps: [Step] + private var calls: [ProviderCall] = [] + + init(steps: [Step]) { + self.steps = steps + } + + func search( + keyword: String, + storefrontCode: String, + platform: AppPlatform, + limit: Int + ) async throws -> SearchRankingPage { + calls.append(ProviderCall( + keyword: keyword, + storefront: storefrontCode, + platform: platform, + limit: limit + )) + guard !steps.isEmpty else { + throw OpenASOError.unexpectedResponse + } + switch steps.removeFirst() { + case .page(let page): + return page + case .failure(let error): + throw error + } + } + + func recordedCalls() -> [ProviderCall] { + calls + } +} + +private actor GatedResearchRankingProvider: SearchRankingProvider { + private var continuation: CheckedContinuation? + private var pendingPage: SearchRankingPage? + private var didStart = false + private var startWaiters: [CheckedContinuation] = [] + + func search( + keyword: String, + storefrontCode: String, + platform: AppPlatform, + limit: Int + ) async throws -> SearchRankingPage { + didStart = true + let waiters = startWaiters + startWaiters.removeAll() + for waiter in waiters { + waiter.resume() + } + if let pendingPage { + self.pendingPage = nil + return pendingPage + } + return await withCheckedContinuation { continuation in + self.continuation = continuation + } + } + + func waitUntilStarted() async { + guard !didStart else { return } + await withCheckedContinuation { continuation in + startWaiters.append(continuation) + } + } + + func succeed(_ page: SearchRankingPage) { + guard let continuation else { + pendingPage = page + return + } + self.continuation = nil + continuation.resume(returning: page) + } +} + +private final class MetadataEnrichmentRecorder: Sendable { + private let batches = Mutex<[[RankingMetadataEnrichmentRequest]]>([]) + + func record(_ requests: [RankingMetadataEnrichmentRequest]) { + batches.withLock { batches in + batches.append(requests) + } + } + + func recordedBatches() -> [[RankingMetadataEnrichmentRequest]] { + batches.withLock { batches in + batches + } + } +} + +private final class DateSequence: Sendable { + private let dates: Mutex<[Date]> + + init(_ dates: [Date]) { + self.dates = Mutex(dates) + } + + func next() -> Date { + dates.withLock { dates in + precondition(!dates.isEmpty, "A test provider returned more pages than expected") + return dates.removeFirst() + } + } +} + +private struct NoOpAppResolver: AppResolver { + func resolve(appStoreID: Int64, storefrontCode: String) async throws -> ResolvedApp { + throw OpenASOError.appNotFound + } + + func searchApps( + named query: String, + storefrontCode: String, + limit: Int + ) async throws -> [ResolvedApp] { + [] + } +} + +private struct DatabaseState: Sendable { + let crawlCount: Int + let observationItemCount: Int + let storeAppIDs: [Int64] + let metadataCount: Int + let screenshotCount: Int + let latestRatingIDs: [Int64] + let dailyRatingIDs: [Int64] + let statsAppStoreIDs: [Int64] + let trackedCount: Int +} + +private struct TrackedSharedRaceState: Sendable { + let sharedObservedAt: Date + let sharedAppStoreIDs: [Int64] + let snapshotSearchedAt: Date + let snapshotRank: Int? + let snapshotResultCount: Int + let snapshotAppStoreIDs: [Int64] + let snapshotPositions: [Int] + let trackLastRefreshAt: Date? + let trackRankingAppCount: Int? + let rankingStatusMessage: String? + let rankingStatusUpdatedAt: Date? +} + +private func trackedSharedRaceState( + trackIdentityKey: String, + queryKey: String, + in store: BackgroundModelStore +) async throws -> TrackedSharedRaceState { + try await store.read { modelContext in + let targetTrackIdentityKey = trackIdentityKey + var trackDescriptor = FetchDescriptor( + predicate: #Predicate { track in + track.identityKey == targetTrackIdentityKey + } + ) + trackDescriptor.fetchLimit = 1 + guard let track = try modelContext.fetch(trackDescriptor).first else { + throw OpenASOError.appNotFound + } + + let targetQueryKey = queryKey + var crawlDescriptor = FetchDescriptor( + predicate: #Predicate { crawl in + crawl.queryKey == targetQueryKey + } + ) + crawlDescriptor.fetchLimit = 1 + guard let crawl = try modelContext.fetch(crawlDescriptor).first else { + throw OpenASOError.unexpectedResponse + } + + var snapshotDescriptor = FetchDescriptor( + predicate: #Predicate { snapshot in + snapshot.trackIdentityKey == targetTrackIdentityKey + } + ) + snapshotDescriptor.fetchLimit = 1 + guard let snapshot = try modelContext.fetch(snapshotDescriptor).first else { + throw OpenASOError.unexpectedResponse + } + let rankedResults = snapshot.topResults.sorted { + if $0.position != $1.position { return $0.position < $1.position } + return $0.appStoreID < $1.appStoreID + } + let status = try TrackedKeywordRefreshStatusStore.snapshot( + for: track, + in: modelContext + ) + return TrackedSharedRaceState( + sharedObservedAt: crawl.observedAt, + sharedAppStoreIDs: crawl.items.sorted { $0.position < $1.position }.map(\.appStoreID), + snapshotSearchedAt: snapshot.searchedAt, + snapshotRank: snapshot.rank, + snapshotResultCount: snapshot.resultCount, + snapshotAppStoreIDs: rankedResults.map(\.appStoreID), + snapshotPositions: rankedResults.map(\.position), + trackLastRefreshAt: track.lastRefreshAt, + trackRankingAppCount: track.rankingAppCount, + rankingStatusMessage: status.rankingMessage, + rankingStatusUpdatedAt: status.rankingUpdatedAt + ) + } +} + +private struct TrackedAppStateRecord: Equatable, Sendable { + let appStoreID: Int64 + let bundleID: String? + let name: String + let sellerName: String? + let createdAt: Date + let isPinned: Bool + let sidebarSortOrder: Int +} + +private struct TrackedKeywordStateRecord: Equatable, Sendable { + let identityKey: String + let appStoreID: Int64 + let term: String + let storefront: String + let platform: AppPlatform + let rankingAppCount: Int? + let lastRefreshAt: Date? + let notes: String + let statusMessage: String? + let createdAt: Date +} + +private struct TrackedSnapshotStateRecord: Equatable, Sendable { + let snapshotKey: String + let trackIdentityKey: String + let rank: Int? + let searchedAt: Date + let source: RankingSource + let resultCount: Int + let errorMessage: String? +} + +private struct TrackedRankedResultStateRecord: Equatable, Sendable { + let snapshotKey: String + let position: Int + let appStoreID: Int64 + let bundleID: String? + let name: String + let subtitle: String? + let sellerName: String? +} + +private struct TrackedStatusStateRecord: Equatable, Sendable { + let statusKey: String + let trackIdentityKey: String + let trackCreatedAt: Date + let appStoreID: Int64 + let domainRaw: String + let message: String? + let updatedAt: Date +} + +private struct TrackedAttemptStateRecord: Equatable, Sendable { + let trackIdentityKey: String + let appStoreID: Int64 + let lastRankingRefreshAttemptAt: Date +} + +private struct TrackedGraphState: Equatable, Sendable { + let apps: [TrackedAppStateRecord] + let keywords: [TrackedKeywordStateRecord] + let snapshots: [TrackedSnapshotStateRecord] + let results: [TrackedRankedResultStateRecord] + let statuses: [TrackedStatusStateRecord] + let attempts: [TrackedAttemptStateRecord] +} + +private func seedFullyPopulatedTrackedGraph( + queryKey: String, + in store: BackgroundModelStore +) async throws { + try await store.write { modelContext in + let targetQueryKey = queryKey + var descriptor = FetchDescriptor( + predicate: #Predicate { query in + query.queryKey == targetQueryKey + } + ) + descriptor.fetchLimit = 1 + guard let query = try modelContext.fetch(descriptor).first else { + throw OpenASOError.unexpectedResponse + } + + let appCreatedAt = utcDate(year: 2026, month: 5, day: 1, hour: 9) + let trackCreatedAt = utcDate(year: 2026, month: 5, day: 2, hour: 9) + let snapshotAt = utcDate(year: 2026, month: 5, day: 3, hour: 9) + let statusAt = utcDate(year: 2026, month: 5, day: 4, hour: 9) + let attemptAt = utcDate(year: 2026, month: 5, day: 5, hour: 9) + let trackedApp = TrackedApp( + appStoreID: 900, + bundleID: "example.tracked.900", + name: "Seeded tracked app", + sellerName: "Seeded seller", + defaultPlatform: .ipad, + sidebarSortOrder: 23, + createdAt: appCreatedAt + ) + trackedApp.isPinned = true + let track = TrackedAppKeyword( + term: query.term, + storefront: query.storefront, + platform: query.platform, + trackedApp: trackedApp, + query: query, + createdAt: trackCreatedAt + ) + track.rankingAppCount = 88 + track.lastRefreshAt = snapshotAt + track.notes = "seeded notes" + trackedApp.keywordTracks.append(track) + modelContext.insert(trackedApp) + modelContext.insert(track) + + let snapshot = TrackedKeywordDailyRanking( + rank: 7, + searchedAt: snapshotAt, + source: .appStoreWeb, + resultCount: 88, + errorMessage: "seeded snapshot error", + keywordTrack: track + ) + track.snapshots.append(snapshot) + modelContext.insert(snapshot) + let rankedResult = TrackedKeywordRankedResult( + position: 4, + appStoreID: 901, + bundleID: "example.result.901", + name: "Seeded ranked result", + subtitle: "Seeded subtitle", + sellerName: "Seeded result seller", + snapshot: snapshot + ) + snapshot.topResults.append(rankedResult) + modelContext.insert(rankedResult) + + try TrackedKeywordRefreshStatusStore.set( + "seeded ranking failure", + domain: .ranking, + for: track, + updatedAt: statusAt, + in: modelContext + ) + track.statusMessage = "untouched legacy scalar" + modelContext.insert(TrackedAppKeywordRefreshAttempt( + trackIdentityKey: track.identityKey, + appStoreID: track.appStoreID, + lastRankingRefreshAttemptAt: attemptAt + )) + } +} + +private func trackedGraphState(in store: BackgroundModelStore) async throws -> TrackedGraphState { + try await store.read { modelContext in + let apps = try modelContext.fetch(FetchDescriptor()) + .map { + TrackedAppStateRecord( + appStoreID: $0.appStoreID, + bundleID: $0.bundleID, + name: $0.name, + sellerName: $0.sellerName, + createdAt: $0.createdAt, + isPinned: $0.isPinned, + sidebarSortOrder: $0.sidebarSortOrder + ) + } + .sorted { $0.appStoreID < $1.appStoreID } + let keywords = try modelContext.fetch(FetchDescriptor()) + .map { + TrackedKeywordStateRecord( + identityKey: $0.identityKey, + appStoreID: $0.appStoreID, + term: $0.term, + storefront: $0.storefront, + platform: $0.platform, + rankingAppCount: $0.rankingAppCount, + lastRefreshAt: $0.lastRefreshAt, + notes: $0.notes, + statusMessage: $0.statusMessage, + createdAt: $0.createdAt + ) + } + .sorted { $0.identityKey < $1.identityKey } + let snapshots = try modelContext.fetch(FetchDescriptor()) + .map { + TrackedSnapshotStateRecord( + snapshotKey: $0.snapshotKey, + trackIdentityKey: $0.trackIdentityKey, + rank: $0.rank, + searchedAt: $0.searchedAt, + source: $0.source, + resultCount: $0.resultCount, + errorMessage: $0.errorMessage + ) + } + .sorted { $0.snapshotKey < $1.snapshotKey } + let results = try modelContext.fetch(FetchDescriptor()) + .map { + TrackedRankedResultStateRecord( + snapshotKey: $0.snapshotKey, + position: $0.position, + appStoreID: $0.appStoreID, + bundleID: $0.bundleID, + name: $0.name, + subtitle: $0.subtitle, + sellerName: $0.sellerName + ) + } + .sorted { + if $0.snapshotKey != $1.snapshotKey { return $0.snapshotKey < $1.snapshotKey } + if $0.position != $1.position { return $0.position < $1.position } + return $0.appStoreID < $1.appStoreID + } + let statuses = try modelContext.fetch(FetchDescriptor()) + .map { + TrackedStatusStateRecord( + statusKey: $0.statusKey, + trackIdentityKey: $0.trackIdentityKey, + trackCreatedAt: $0.trackCreatedAt, + appStoreID: $0.appStoreID, + domainRaw: $0.domainRaw, + message: $0.message, + updatedAt: $0.updatedAt + ) + } + .sorted { $0.statusKey < $1.statusKey } + let attempts = try modelContext.fetch(FetchDescriptor()) + .map { + TrackedAttemptStateRecord( + trackIdentityKey: $0.trackIdentityKey, + appStoreID: $0.appStoreID, + lastRankingRefreshAttemptAt: $0.lastRankingRefreshAttemptAt + ) + } + .sorted { $0.trackIdentityKey < $1.trackIdentityKey } + return TrackedGraphState( + apps: apps, + keywords: keywords, + snapshots: snapshots, + results: results, + statuses: statuses, + attempts: attempts + ) + } +} + +private func databaseState(in store: BackgroundModelStore) async throws -> DatabaseState { + try await store.read { modelContext in + let trackedCount = try modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + return DatabaseState( + crawlCount: try modelContext.fetchCount(FetchDescriptor()), + observationItemCount: try modelContext.fetchCount(FetchDescriptor()), + storeAppIDs: try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted(), + metadataCount: try modelContext.fetchCount(FetchDescriptor()), + screenshotCount: try modelContext.fetchCount(FetchDescriptor()), + latestRatingIDs: try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted(), + dailyRatingIDs: try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted(), + statsAppStoreIDs: try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted(), + trackedCount: trackedCount + ) + } +} + +private func sharedWriteCount(in store: BackgroundModelStore) async throws -> Int { + try await store.read { modelContext in + try modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + + modelContext.fetchCount(FetchDescriptor()) + } +} + +private struct ObservationRecord: Equatable, Sendable { + let observedAt: Date + let source: RankingSource + let appStoreIDs: [Int64] + let positions: [Int] +} + +private func observationRecords(in store: BackgroundModelStore) async throws -> [ObservationRecord] { + try await store.read { modelContext in + try modelContext.fetch(FetchDescriptor()) + .map { observation in + let items = observation.items.sorted { + if $0.position != $1.position { return $0.position < $1.position } + return $0.appStoreID < $1.appStoreID + } + return ObservationRecord( + observedAt: observation.observedAt, + source: observation.source, + appStoreIDs: items.map(\.appStoreID), + positions: items.map(\.position) + ) + } + .sorted { + if $0.observedAt != $1.observedAt { return $0.observedAt < $1.observedAt } + return $0.source.rawValue < $1.source.rawValue + } + } +} + +private struct RatingRecord: Equatable, Sendable { + let observedAt: Date + let ratingCount: Int? + let averageRating: Double? +} + +private func latestRatingRecord(in store: BackgroundModelStore) async throws -> RatingRecord? { + try await store.read { modelContext in + try modelContext.fetch(FetchDescriptor()).first.map { + RatingRecord( + observedAt: $0.observedAt, + ratingCount: $0.ratingCount, + averageRating: $0.averageRating + ) + } + } +} + +private struct CrossSourceWatermarkRecord: Equatable, Sendable { + let latestRatingCount: Int? + let latestRatingObservedAt: Date + let latestRatingSource: AppStorefrontSource + let dailyRatingCount: Int? + let dailyRatingObservedAt: Date + let dailyRatingSource: AppStorefrontSource + let storeAppName: String + let storeAppMetadataObservedAt: Date + let storeAppLanguageSource: AppStorefrontMetadataSource? + let storefrontMetadataName: String + let storefrontMetadataObservedAt: Date + let storefrontMetadataSource: AppStorefrontMetadataSource + let iTunesCrawlName: String + let iTunesCrawlObservedAt: Date +} + +private func crossSourceWatermarkRecord( + queryKey: String, + appStoreID: Int64, + storefront: String, + in store: BackgroundModelStore +) async throws -> CrossSourceWatermarkRecord { + try await store.read { modelContext in + let normalizedStorefront = storefront.trimmingCharacters(in: .whitespacesAndNewlines).lowercased() + let latestRatingKey = LatestAppRating.makeIdentityKey( + appStoreID: appStoreID, + storefront: normalizedStorefront + ) + let storefrontMetadataKey = AppStorefrontMetadata.makeIdentityKey( + appStoreID: appStoreID, + storefront: normalizedStorefront + ) + let targetAppStoreID = appStoreID + let targetQueryKey = queryKey + let iTunesSourceRaw = RankingSource.iTunesFallback.rawValue + + var latestRatingDescriptor = FetchDescriptor( + predicate: #Predicate { rating in + rating.identityKey == latestRatingKey + } + ) + latestRatingDescriptor.fetchLimit = 1 + var dailyRatingDescriptor = FetchDescriptor( + predicate: #Predicate { rating in + rating.appStoreID == targetAppStoreID + && rating.storefront == normalizedStorefront + } + ) + dailyRatingDescriptor.fetchLimit = 1 + var storeAppDescriptor = FetchDescriptor( + predicate: #Predicate { app in + app.appStoreID == targetAppStoreID + } + ) + storeAppDescriptor.fetchLimit = 1 + var storefrontMetadataDescriptor = FetchDescriptor( + predicate: #Predicate { metadata in + metadata.identityKey == storefrontMetadataKey + } + ) + storefrontMetadataDescriptor.fetchLimit = 1 + var iTunesCrawlDescriptor = FetchDescriptor( + predicate: #Predicate { crawl in + crawl.queryKey == targetQueryKey + && crawl.sourceRaw == iTunesSourceRaw + } + ) + iTunesCrawlDescriptor.fetchLimit = 1 + + guard let latestRating = try modelContext.fetch(latestRatingDescriptor).first, + let dailyRating = try modelContext.fetch(dailyRatingDescriptor).first, + let storeApp = try modelContext.fetch(storeAppDescriptor).first, + let storefrontMetadata = try modelContext.fetch(storefrontMetadataDescriptor).first, + let iTunesCrawl = try modelContext.fetch(iTunesCrawlDescriptor).first, + let iTunesCrawlItem = iTunesCrawl.items.first(where: { + $0.appStoreID == targetAppStoreID + }) else { + throw OpenASOError.unexpectedResponse + } + + return CrossSourceWatermarkRecord( + latestRatingCount: latestRating.ratingCount, + latestRatingObservedAt: latestRating.observedAt, + latestRatingSource: latestRating.source, + dailyRatingCount: dailyRating.ratingCount, + dailyRatingObservedAt: dailyRating.observedAt, + dailyRatingSource: dailyRating.source, + storeAppName: storeApp.name, + storeAppMetadataObservedAt: storeApp.lastMetadataRefreshAt, + storeAppLanguageSource: storeApp.supportedLanguageCodesSource, + storefrontMetadataName: storefrontMetadata.name, + storefrontMetadataObservedAt: storefrontMetadata.lastFetchedAt, + storefrontMetadataSource: storefrontMetadata.source, + iTunesCrawlName: iTunesCrawlItem.name, + iTunesCrawlObservedAt: iTunesCrawl.observedAt + ) + } +} + +private struct CatalogWatermarkRecord: Equatable, Sendable { + let defaultStorefront: String + let name: String + let iconURLString: String? + let bundleID: String? + let sellerName: String? + let version: String? + let globalMetadataObservedAt: Date + let supportedLanguageCodes: [String] + let supportedLanguageCodesObservedAt: Date? + let defaultMetadataName: String + let defaultMetadataIconURLString: String? + let defaultMetadataObservedAt: Date + let nonDefaultMetadataName: String + let nonDefaultMetadataIconURLString: String? + let nonDefaultMetadataObservedAt: Date +} + +private func catalogWatermarkRecord( + appStoreID: Int64, + defaultStorefront: String, + nonDefaultStorefront: String, + in store: BackgroundModelStore +) async throws -> CatalogWatermarkRecord { + try await store.read { modelContext in + let targetAppStoreID = appStoreID + let defaultMetadataKey = AppStorefrontMetadata.makeIdentityKey( + appStoreID: appStoreID, + storefront: defaultStorefront + ) + let nonDefaultMetadataKey = AppStorefrontMetadata.makeIdentityKey( + appStoreID: appStoreID, + storefront: nonDefaultStorefront + ) + var storeAppDescriptor = FetchDescriptor( + predicate: #Predicate { app in + app.appStoreID == targetAppStoreID + } + ) + storeAppDescriptor.fetchLimit = 1 + var defaultMetadataDescriptor = FetchDescriptor( + predicate: #Predicate { metadata in + metadata.identityKey == defaultMetadataKey + } + ) + defaultMetadataDescriptor.fetchLimit = 1 + var nonDefaultMetadataDescriptor = FetchDescriptor( + predicate: #Predicate { metadata in + metadata.identityKey == nonDefaultMetadataKey + } + ) + nonDefaultMetadataDescriptor.fetchLimit = 1 + + guard let storeApp = try modelContext.fetch(storeAppDescriptor).first, + let defaultMetadata = try modelContext.fetch(defaultMetadataDescriptor).first, + let nonDefaultMetadata = try modelContext.fetch(nonDefaultMetadataDescriptor).first else { + throw OpenASOError.unexpectedResponse + } + + return CatalogWatermarkRecord( + defaultStorefront: storeApp.defaultStorefront, + name: storeApp.name, + iconURLString: storeApp.iconURLString, + bundleID: storeApp.bundleID, + sellerName: storeApp.sellerName, + version: storeApp.version, + globalMetadataObservedAt: storeApp.lastMetadataRefreshAt, + supportedLanguageCodes: storeApp.supportedLanguageCodes, + supportedLanguageCodesObservedAt: storeApp.supportedLanguageCodesFetchedAt, + defaultMetadataName: defaultMetadata.name, + defaultMetadataIconURLString: defaultMetadata.iconURLString, + defaultMetadataObservedAt: defaultMetadata.lastFetchedAt, + nonDefaultMetadataName: nonDefaultMetadata.name, + nonDefaultMetadataIconURLString: nonDefaultMetadata.iconURLString, + nonDefaultMetadataObservedAt: nonDefaultMetadata.lastFetchedAt + ) + } +} + +private struct StatsRecord: Equatable, Sendable { + let appStoreID: Int64 + let bestRank: Int? + let latestRank: Int? + let observationCount: Int +} + +private func statsRecords(in store: BackgroundModelStore) async throws -> [StatsRecord] { + try await store.read { modelContext in + try modelContext.fetch(FetchDescriptor()) + .map { + StatsRecord( + appStoreID: $0.appStoreID, + bestRank: $0.bestRank, + latestRank: $0.latestRank, + observationCount: $0.observationCount + ) + } + .sorted { $0.appStoreID < $1.appStoreID } + } +} + +private func rankingItem( + position: Int, + appStoreID: Int64, + name: String? = nil, + bundleID: String? = nil, + sellerName: String? = nil, + iconURLString: String? = nil, + version: String? = nil, + supportedLanguageCodes: [String]? = nil, + ratingCount: Int? = nil, + averageRating: Double? = nil +) -> SearchRankingItem { + SearchRankingItem( + position: position, + appStoreID: appStoreID, + bundleID: bundleID ?? "example.\(appStoreID)", + name: name ?? "App \(appStoreID)", + subtitle: "Subtitle \(appStoreID)", + sellerName: sellerName ?? "Seller \(appStoreID)", + iconURLString: iconURLString ?? "https://example.com/\(appStoreID).png", + version: version ?? "1.0", + primaryGenreID: 6002, + primaryGenreName: "Utilities", + descriptionText: "Description \(appStoreID)", + releaseNotes: "Notes \(appStoreID)", + supportedLanguageCodes: supportedLanguageCodes ?? ["EN", "FR"], + screenshotURLs: ["https://example.com/\(appStoreID)-iphone.png"], + ratingCount: ratingCount, + averageRating: averageRating, + platform: .ipad + ) +} + +private func utcDate( + year: Int, + month: Int, + day: Int, + hour: Int, + minute: Int = 0 +) -> Date { + var calendar = Calendar(identifier: .gregorian) + calendar.timeZone = TimeZone(secondsFromGMT: 0)! + return calendar.date(from: DateComponents( + timeZone: calendar.timeZone, + year: year, + month: month, + day: day, + hour: hour, + minute: minute + ))! +} diff --git a/OpenASOTests/OpenASOMCPServiceTests.swift b/OpenASOTests/OpenASOMCPServiceTests.swift index 54e591e..33f8403 100644 --- a/OpenASOTests/OpenASOMCPServiceTests.swift +++ b/OpenASOTests/OpenASOMCPServiceTests.swift @@ -830,6 +830,57 @@ struct OpenASOMCPServiceTests { #expect(competitors.first?.appStoreID == "456") } + @Test + func discoverKeywordLandscapeCoordinatorPathPersistsFullRankingBeforeEnrichment() async throws { + let rankingProvider = StubMCPRankingProvider(pages: [ + "calorie tracker::us::iphone": SearchRankingPage(items: [ + makeRankingItem(position: 1, appStoreID: 456, name: "MyFitnessPal", ratingCount: 1_000), + makeRankingItem(position: 2, appStoreID: 123, name: "Calorie Tracker", ratingCount: 100) + ], source: .iTunesFallback) + ]) + let container = try ModelContainerFactory.makeModelContainer(isStoredInMemoryOnly: true) + let enrichmentRecorder = MCPMetadataEnrichmentRecorder(modelContainer: container) + let context = try MCPTestContext( + modelContainer: container, + rankingProvider: rankingProvider, + useRankingRefreshCoordinator: true, + metadataEnrichmentScheduler: enrichmentRecorder.record + ) + try context.insertTrackedApp(appStoreID: 123, name: "Calorie Tracker") + + let result = try await context.service.discoverKeywordLandscape( + appStoreID: 123, + storefronts: ["us"], + platform: "iphone", + keywordLimit: 1, + competitorLimit: 1, + includeReviews: false + ) + + #expect(result.errors.isEmpty) + #expect(result.verifiedKeywords.count == 1) + #expect(result.verifiedKeywords.first?.keyword == "calorie tracker") + #expect(result.verifiedKeywords.first?.isTracked == true) + #expect(result.verifiedKeywords.first?.targetRank == 2) + #expect(await rankingProvider.searchedLimitsSnapshot() == [200]) + #expect(enrichmentRecorder.snapshot() == [MCPMetadataEnrichmentRecord( + requests: [ + RankingMetadataEnrichmentRequest( + appStoreID: 456, + storefront: "us", + platform: .iphone + ), + RankingMetadataEnrichmentRequest( + appStoreID: 123, + storefront: "us", + platform: .iphone + ) + ], + committedCrawlCount: 1, + committedSnapshotCount: 1 + )]) + } + @Test func suggestKeywordsStopsVerificationBeforeToolTimeoutBudget() async throws { let rankingProvider = StubMCPRankingProvider(pages: [:]) @@ -918,6 +969,57 @@ struct OpenASOMCPServiceTests { #expect(metricsRefresh.outcomes.first?.track.statusMessage?.contains("Connect an Apple Ads") == true) } + @Test + func refreshKeywordRankingsCoordinatorPathRequestsFullPageAndEnrichesAfterCommit() async throws { + let rankingProvider = StubMCPRankingProvider(pages: [ + "calorie tracker::us::iphone": SearchRankingPage(items: [ + makeRankingItem(position: 1, appStoreID: 456, name: "MyFitnessPal", ratingCount: 1_000), + makeRankingItem(position: 2, appStoreID: 123, name: "Cal AI", ratingCount: 100) + ], source: .iTunesFallback) + ]) + let container = try ModelContainerFactory.makeModelContainer(isStoredInMemoryOnly: true) + let enrichmentRecorder = MCPMetadataEnrichmentRecorder(modelContainer: container) + let context = try MCPTestContext( + modelContainer: container, + rankingProvider: rankingProvider, + useRankingRefreshCoordinator: true, + metadataEnrichmentScheduler: enrichmentRecorder.record + ) + try context.insertTrackedApp(appStoreID: 123, name: "Cal AI") + _ = try await context.service.addKeywords( + appStoreID: 123, + keywords: ["calorie tracker"], + storefronts: ["us"], + platform: "iphone" + ) + + let refresh = try await context.service.refreshKeywordRankings( + appStoreID: 123, + storefronts: ["us"], + platform: "iphone" + ) + + #expect(refresh.summary.refreshed == 1) + #expect(refresh.outcomes.first?.track.latestRank == 2) + #expect(await rankingProvider.searchedLimitsSnapshot() == [200]) + #expect(enrichmentRecorder.snapshot() == [MCPMetadataEnrichmentRecord( + requests: [ + RankingMetadataEnrichmentRequest( + appStoreID: 456, + storefront: "us", + platform: .iphone + ), + RankingMetadataEnrichmentRequest( + appStoreID: 123, + storefront: "us", + platform: .iphone + ) + ], + committedCrawlCount: 1, + committedSnapshotCount: 1 + )]) + } + @Test func refreshKeywordRankingsLimitCapsKeywordTracksRefreshed() async throws { let rankingProvider = StubMCPRankingProvider(pages: [ @@ -3189,6 +3291,7 @@ private struct MCPTestContext { resolver: StubMCPAppResolver = StubMCPAppResolver(), rankingProvider: (any SearchRankingProvider)? = nil, useRankingRefreshCoordinator: Bool = false, + metadataEnrichmentScheduler: (@Sendable ([RankingMetadataEnrichmentRequest]) -> Void)? = nil, rankingRefreshScheduler: OpenASOMCPRankingRefreshScheduler = OpenASOMCPRankingRefreshScheduler(), persistRankingRefreshAttempts: OpenASOMCPService.RankingRefreshAttemptsPersistence? = nil, includeReviewService: Bool = false, @@ -3217,7 +3320,11 @@ private struct MCPTestContext { ) let rankingRefreshCoordinator = rankingProvider.flatMap { provider in useRankingRefreshCoordinator - ? RankingRefreshCoordinator(rankingProvider: provider, appCatalogService: appCatalogService) + ? RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: appCatalogService, + metadataEnrichmentScheduler: metadataEnrichmentScheduler + ) : nil } self.service = OpenASOMCPService( @@ -3666,10 +3773,48 @@ private struct StubMCPAppResolver: AppResolver { } } +private struct MCPMetadataEnrichmentRecord: Equatable, Sendable { + let requests: [RankingMetadataEnrichmentRequest] + let committedCrawlCount: Int + let committedSnapshotCount: Int +} + +private final class MCPMetadataEnrichmentRecorder: Sendable { + private let modelContainer: ModelContainer + private let records = Mutex<[MCPMetadataEnrichmentRecord]>([]) + + init(modelContainer: ModelContainer) { + self.modelContainer = modelContainer + } + + func record(_ requests: [RankingMetadataEnrichmentRequest]) { + let modelContext = ModelContext(modelContainer) + let record = MCPMetadataEnrichmentRecord( + requests: requests, + committedCrawlCount: (try? modelContext.fetchCount( + FetchDescriptor() + )) ?? -1, + committedSnapshotCount: (try? modelContext.fetchCount( + FetchDescriptor() + )) ?? -1 + ) + records.withLock { records in + records.append(record) + } + } + + func snapshot() -> [MCPMetadataEnrichmentRecord] { + records.withLock { records in + records + } + } +} + private actor StubMCPRankingProvider: SearchRankingProvider { private var pages: [String: SearchRankingPage] private var failures: [String: OpenASOError] private var searchedKeys: [String] = [] + private var searchedLimits: [Int] = [] init( pages: [String: SearchRankingPage], @@ -3691,6 +3836,10 @@ private actor StubMCPRankingProvider: SearchRankingProvider { searchedKeys } + func searchedLimitsSnapshot() -> [Int] { + searchedLimits + } + func search(keyword: String, storefrontCode: String, platform: AppPlatform, limit: Int) async throws -> SearchRankingPage { let key = TrackedAppKeyword.makeQueryKey( term: keyword, @@ -3698,6 +3847,7 @@ private actor StubMCPRankingProvider: SearchRankingProvider { platform: platform ) searchedKeys.append(key) + searchedLimits.append(limit) if let failure = failures[key] { throw failure } diff --git a/OpenASOTests/RankingRefreshCoordinatorTests.swift b/OpenASOTests/RankingRefreshCoordinatorTests.swift index 2c993fa..2c875ef 100644 --- a/OpenASOTests/RankingRefreshCoordinatorTests.swift +++ b/OpenASOTests/RankingRefreshCoordinatorTests.swift @@ -320,7 +320,7 @@ struct RankingRefreshCoordinatorTests { let catalogService = AppCatalogService(appResolver: resolver) let coordinator = RankingRefreshCoordinator(rankingProvider: provider, appCatalogService: catalogService) - let result = await coordinator.refresh(track: track, in: modelContext, limit: 10) + let result = await coordinator.refresh(track: track, in: modelContext) switch result { case .success(let snapshot): @@ -380,6 +380,210 @@ struct RankingRefreshCoordinatorTests { #expect(ratingSnapshots.first?.averageRating == 4.65041) } + @Test + func persistenceFailureAfterMutationRollsBackTheEntireRefresh() throws { + let container = try makeInMemoryContainer() + let modelContext = ModelContext(container) + modelContext.autosaveEnabled = false + let trackedApp = TrackedApp( + appStoreID: 700, + bundleID: "example.tracked.700", + name: "Tracked", + sellerName: "Example", + defaultPlatform: .iphone + ) + let track = try makeTrackedAppKeyword( + term: "rollback test", + trackedApp: trackedApp, + in: modelContext + ) + trackedApp.keywordTracks.append(track) + modelContext.insert(trackedApp) + modelContext.insert(track) + try modelContext.save() + let baseline = try rankingPersistenceState(in: modelContext) + + let checkpoint = FailingRankingPersistenceCheckpoint(failingOnCall: 1) + let coordinator = RankingRefreshCoordinator( + rankingProvider: StubRankingProvider( + page: SearchRankingPage(items: [], source: .iTunesFallback) + ), + appCatalogService: AppCatalogService(appResolver: StubAppResolver()), + persistenceMutationCheckpoint: checkpoint.call + ) + let pageResult = RankingRefreshPageResult( + request: RankingRefreshRequest(track: track), + page: SearchRankingPage( + items: [SearchRankingItem( + position: 1, + appStoreID: 701, + bundleID: "example.result.701", + name: "Result", + sellerName: "Example", + screenshotURLs: ["https://example.com/result-701.png"], + ratingCount: 42, + averageRating: 4.5, + platform: .iphone + )], + source: .iTunesFallback + ), + searchedAt: date( + year: 2026, + month: 7, + day: 19, + hour: 12, + minute: 0, + calendar: utcCalendar() + ), + observedHour: nil, + submissionCount: 1, + winningCount: 1, + confidence: "single_source" + ) + + #expect(throws: OpenASOError.unexpectedResponse) { + _ = try coordinator.persistRankingPage(pageResult, in: modelContext) + } + + #expect(checkpoint.callCount() == 1) + #expect(!modelContext.hasChanges) + #expect(try rankingPersistenceState(in: modelContext) == baseline) + } + + @Test + func pendingUnrelatedEditArrivingDuringFetchIsNeitherCommittedNorRolledBack() async throws { + let container = try makeInMemoryContainer() + let modelContext = ModelContext(container) + modelContext.autosaveEnabled = false + let trackedApp = TrackedApp( + appStoreID: 710, + bundleID: "example.tracked.710", + name: "Tracked", + sellerName: "Example", + defaultPlatform: .iphone + ) + let track = try makeTrackedAppKeyword( + term: "pending edit test", + trackedApp: trackedApp, + in: modelContext + ) + trackedApp.keywordTracks.append(track) + modelContext.insert(trackedApp) + modelContext.insert(track) + let unrelatedStorefront = Storefront( + code: "zz", + name: "Persisted storefront name", + flagEmoji: "ZZ", + languageCode: "en" + ) + modelContext.insert(unrelatedStorefront) + try modelContext.save() + let baseline = try rankingPersistenceState(in: modelContext) + + let provider = GatedRankingProvider() + let coordinator = RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: AppCatalogService(appResolver: StubAppResolver()) + ) + let refreshTask = Task { @MainActor in + await coordinator.refresh(tracks: [track], in: modelContext) + } + await provider.waitUntilStarted() + + unrelatedStorefront.name = "Pending storefront name" + #expect(modelContext.hasChanges) + await provider.succeed(SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: trackedApp.appStoreID, platform: .iphone)], + source: .iTunesFallback + )) + let outcomes = await refreshTask.value + + #expect(outcomes.count == 1) + #expect(outcomes.first?.error != nil) + #expect(outcomes.first?.snapshotID == nil) + #expect(modelContext.hasChanges) + #expect(unrelatedStorefront.name == "Pending storefront name") + + let verificationContext = ModelContext(container) + verificationContext.autosaveEnabled = false + let durableStorefront = try #require( + verificationContext.fetch(FetchDescriptor()) + .first(where: { $0.code == "zz" }) + ) + #expect(durableStorefront.name == "Persisted storefront name") + #expect(try rankingPersistenceState(in: verificationContext) == baseline) + #expect(try verificationContext.fetchCount(FetchDescriptor()) == 0) + } + + @Test + func pendingEditFromProgressSkipsFinalStatsWithoutCommitOrRollback() async throws { + let container = try makeInMemoryContainer() + let modelContext = ModelContext(container) + modelContext.autosaveEnabled = false + let trackedApp = TrackedApp( + appStoreID: 720, + bundleID: "example.tracked.720", + name: "Tracked", + sellerName: "Example", + defaultPlatform: .iphone + ) + let track = try makeTrackedAppKeyword( + term: "stats reentrancy test", + trackedApp: trackedApp, + in: modelContext + ) + trackedApp.keywordTracks.append(track) + modelContext.insert(trackedApp) + modelContext.insert(track) + let unrelatedStorefront = Storefront( + code: "yy", + name: "Persisted stats storefront name", + flagEmoji: "YY", + languageCode: "en" + ) + modelContext.insert(unrelatedStorefront) + try modelContext.save() + + let provider = StubRankingProvider(page: SearchRankingPage( + items: [rankingItem(position: 1, appStoreID: trackedApp.appStoreID, platform: .iphone)], + source: .iTunesFallback + )) + let coordinator = RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: AppCatalogService(appResolver: StubAppResolver()) + ) + let editInjector = RankingPendingEditInjector( + storefront: unrelatedStorefront, + pendingName: "Pending stats storefront name" + ) + + let outcomes = await coordinator.refresh( + tracks: [track], + in: modelContext, + progress: { completed, _, _ in + await editInjector.injectIfPagePersistenceCompleted(completed) + } + ) + + #expect(outcomes.count == 1) + #expect(outcomes.first?.error == nil) + #expect(outcomes.first?.rank == 1) + #expect(editInjector.injectionCount == 1) + #expect(modelContext.hasChanges) + #expect(unrelatedStorefront.name == "Pending stats storefront name") + + let verificationContext = ModelContext(container) + verificationContext.autosaveEnabled = false + let durableStorefront = try #require( + verificationContext.fetch(FetchDescriptor()) + .first(where: { $0.code == "yy" }) + ) + #expect(durableStorefront.name == "Persisted stats storefront name") + #expect(try verificationContext.fetchCount(FetchDescriptor()) == 1) + #expect(try verificationContext.fetchCount(FetchDescriptor()) == 1) + #expect(try verificationContext.fetchCount(FetchDescriptor()) == 0) + } + @Test func webEnrichmentBackfillsMissingSubtitleForRankingCatalogApp() throws { let container = try makeInMemoryContainer() @@ -1141,7 +1345,7 @@ struct RankingRefreshCoordinatorTests { appCatalogService: AppCatalogService(appResolver: StubAppResolver()) ) - _ = await coordinator.refresh(track: track, in: modelContext, limit: 10) + _ = await coordinator.refresh(track: track, in: modelContext) provider.page = SearchRankingPage( items: [ @@ -1158,7 +1362,7 @@ struct RankingRefreshCoordinatorTests { source: .iTunesFallback ) - let result = await coordinator.refresh(track: track, in: modelContext, limit: 10) + let result = await coordinator.refresh(track: track, in: modelContext) guard case .success(let snapshot) = result else { Issue.record("Expected refresh to succeed") @@ -1517,6 +1721,274 @@ struct AppDetailRefreshServiceQueueTests { appStoreConnectCredentials: AppStoreConnectCredentials(issuerID: "", keyID: "", privateKey: "") ) } + + @Test + func appDetailFanOutSchedulesMetadataOnlyForAppliedSharedObservation() async throws { + let container = try makeInMemoryContainer() + let modelContext = ModelContext(container) + let fixture = try makeDeduplicatedRefreshFixture(in: modelContext) + let searchedAt = isoDate("2026-07-18T13:00:00Z") + let provider = QueryRankingProvider(responses: [ + QueryRankingProvider.key( + term: fixture.firstTrack.term, + storefront: fixture.firstTrack.storefront, + platform: fixture.firstTrack.platform + ): .page(SearchRankingPage( + items: [ + rankingItem(position: 1, appStoreID: fixture.firstApp.appStoreID, platform: .iphone), + rankingItem(position: 2, appStoreID: fixture.secondApp.appStoreID, platform: .iphone), + ], + source: .iTunesFallback + )), + ]) + let metadataRecorder = RankingMetadataEnrichmentRecorder() + let coordinator = RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: AppCatalogService(appResolver: StubAppResolver()), + now: { searchedAt }, + metadataEnrichmentScheduler: metadataRecorder.record + ) + let httpClient = MockHTTPClient { request in + throw OpenASOError.providerUnavailable( + "Unexpected request to \(request.url?.absoluteString ?? "unknown URL")" + ) + } + let defaults = makeDefaults() + let keychain = InMemoryKeychainService() + let service = AppDetailRefreshService( + backgroundModelStore: BackgroundModelStore(modelContainer: container), + refreshCoordinator: coordinator, + keywordMetricsService: KeywordMetricsService( + httpClient: httpClient, + credentialStore: AppleAdsCredentialStore( + defaults: defaults, + keychain: keychain, + loadsEnvironmentCredentials: false + ), + settingsStore: AppSettingsStore(defaults: defaults), + webSessionStore: AppleAdsWebSessionStore( + defaults: defaults, + keychain: keychain + ) + ), + appStorefrontRatingService: AppStorefrontRatingService(httpClient: httpClient), + appStorefrontReviewService: AppStorefrontReviewService(httpClient: httpClient), + appStoreConnectReviewService: AppStoreConnectReviewService( + httpClient: httpClient, + credentialStore: AppStoreConnectCredentialStore( + defaults: defaults, + keychain: keychain + ) + ) + ) + + let result = await service.refresh(AppDetailRefreshRequest( + app: AppDetailRefreshAppSnapshot( + appStoreID: fixture.firstApp.appStoreID, + bundleID: fixture.firstApp.bundleID, + name: fixture.firstApp.name, + subtitle: fixture.firstApp.subtitle, + sellerName: fixture.firstApp.sellerName, + defaultPlatform: fixture.firstApp.defaultPlatform + ), + workspace: .keywords, + storefrontSelection: .storefront(code: "us"), + trackIdentityKeys: [fixture.firstTrack.identityKey, fixture.secondTrack.identityKey], + trigger: "manual", + refreshKeywords: true, + refreshMetrics: false, + refreshRatings: false, + refreshReviews: false, + recordsRatingsReviewsRefresh: false, + popularityContextAppStoreID: nil, + appleAdsWebSession: nil, + appStoreConnectCredentials: AppStoreConnectCredentials( + issuerID: "", + keyID: "", + privateKey: "" + ) + )) + + #expect(result.keywordOutcomes.count == 2) + #expect(result.keywordOutcomes.allSatisfy { $0.error == nil }) + #expect(result.firstError == nil) + #expect(await provider.callCounts() == ["pages::us::iphone": 1]) + #expect(metadataRecorder.recordedBatches() == [[ + RankingMetadataEnrichmentRequest( + appStoreID: fixture.firstApp.appStoreID, + storefront: "us", + platform: .iphone + ), + RankingMetadataEnrichmentRequest( + appStoreID: fixture.secondApp.appStoreID, + storefront: "us", + platform: .iphone + ), + ]]) + + let state = try await BackgroundModelStore(modelContainer: container).read { modelContext in + try rankingPersistenceState(in: modelContext) + } + #expect(state.crawlCount == 1) + #expect(state.snapshotCount == 2) + } + + @Test + func rankingBatchFailureRollsBackEarlierPagesBeforeRecordingFailures() async throws { + let container = try ModelContainerFactory.makeModelContainer(isStoredInMemoryOnly: true) + let modelContext = ModelContext(container) + let trackedApp = TrackedApp( + appStoreID: 800, + bundleID: "example.tracked.800", + name: "Tracked", + sellerName: "Example", + defaultPlatform: .iphone + ) + let firstTrack = try makeTrackedAppKeyword( + term: "first rollback query", + trackedApp: trackedApp, + in: modelContext + ) + let secondTrack = try makeTrackedAppKeyword( + term: "second rollback query", + trackedApp: trackedApp, + in: modelContext + ) + trackedApp.keywordTracks.append(contentsOf: [firstTrack, secondTrack]) + modelContext.insert(trackedApp) + modelContext.insert(firstTrack) + modelContext.insert(secondTrack) + try modelContext.save() + + let backgroundModelStore = BackgroundModelStore(modelContainer: container) + let baseline = try await backgroundModelStore.read { modelContext in + try rankingPersistenceState(in: modelContext) + } + let provider = QueryRankingProvider(responses: [ + QueryRankingProvider.key( + term: firstTrack.term, + storefront: firstTrack.storefront, + platform: firstTrack.platform + ): .page(SearchRankingPage( + items: [SearchRankingItem( + position: 1, + appStoreID: 801, + bundleID: "example.result.801", + name: "First result", + sellerName: "Example", + screenshotURLs: ["https://example.com/result-801.png"], + ratingCount: 81, + averageRating: 4.1, + platform: .iphone + )], + source: .iTunesFallback + )), + QueryRankingProvider.key( + term: secondTrack.term, + storefront: secondTrack.storefront, + platform: secondTrack.platform + ): .page(SearchRankingPage( + items: [SearchRankingItem( + position: 1, + appStoreID: 802, + bundleID: "example.result.802", + name: "Second result", + sellerName: "Example", + screenshotURLs: ["https://example.com/result-802.png"], + ratingCount: 82, + averageRating: 4.2, + platform: .iphone + )], + source: .iTunesFallback + )), + ]) + let checkpoint = FailingRankingPersistenceCheckpoint(failingOnCall: 2) + let coordinator = RankingRefreshCoordinator( + rankingProvider: provider, + appCatalogService: AppCatalogService(appResolver: StubAppResolver()), + persistenceMutationCheckpoint: checkpoint.call + ) + let httpClient = MockHTTPClient { request in + throw OpenASOError.providerUnavailable( + "Unexpected request to \(request.url?.absoluteString ?? "unknown URL")" + ) + } + let defaults = makeDefaults() + let keychain = InMemoryKeychainService() + let service = AppDetailRefreshService( + backgroundModelStore: backgroundModelStore, + refreshCoordinator: coordinator, + keywordMetricsService: KeywordMetricsService( + httpClient: httpClient, + credentialStore: AppleAdsCredentialStore( + defaults: defaults, + keychain: keychain, + loadsEnvironmentCredentials: false + ), + settingsStore: AppSettingsStore(defaults: defaults), + webSessionStore: AppleAdsWebSessionStore( + defaults: defaults, + keychain: keychain + ) + ), + appStorefrontRatingService: AppStorefrontRatingService(httpClient: httpClient), + appStorefrontReviewService: AppStorefrontReviewService(httpClient: httpClient), + appStoreConnectReviewService: AppStoreConnectReviewService( + httpClient: httpClient, + credentialStore: AppStoreConnectCredentialStore( + defaults: defaults, + keychain: keychain + ) + ) + ) + let result = await service.refresh(AppDetailRefreshRequest( + app: AppDetailRefreshAppSnapshot( + appStoreID: trackedApp.appStoreID, + bundleID: trackedApp.bundleID, + name: trackedApp.name, + subtitle: trackedApp.subtitle, + sellerName: trackedApp.sellerName, + defaultPlatform: trackedApp.defaultPlatform + ), + workspace: .keywords, + storefrontSelection: .storefront(code: "us"), + trackIdentityKeys: [firstTrack.identityKey, secondTrack.identityKey], + trigger: "manual", + refreshKeywords: true, + refreshMetrics: false, + refreshRatings: false, + refreshReviews: false, + recordsRatingsReviewsRefresh: false, + popularityContextAppStoreID: nil, + appleAdsWebSession: nil, + appStoreConnectCredentials: AppStoreConnectCredentials( + issuerID: "", + keyID: "", + privateKey: "" + ) + )) + + #expect(checkpoint.callCount() == 2) + #expect(result.keywordOutcomes.count == 2) + #expect(result.keywordOutcomes.allSatisfy { $0.error != nil }) + #expect(result.firstError != nil) + let state = try await backgroundModelStore.read { modelContext in + try rankingPersistenceState(in: modelContext) + } + #expect(state.crawlCount == baseline.crawlCount) + #expect(state.observationItemCount == baseline.observationItemCount) + #expect(state.snapshotCount == baseline.snapshotCount) + #expect(state.rankedResultCount == baseline.rankedResultCount) + #expect(state.storeAppIDs == baseline.storeAppIDs) + #expect(state.storefrontMetadataCount == baseline.storefrontMetadataCount) + #expect(state.screenshotCount == baseline.screenshotCount) + #expect(state.latestRatingCount == baseline.latestRatingCount) + #expect(state.dailyRatingCount == baseline.dailyRatingCount) + #expect(state.statsCount == baseline.statsCount) + #expect(state.trackStates == baseline.trackStates) + #expect(state.statusCount == 2) + } + } @MainActor @@ -1603,6 +2075,98 @@ private struct DeduplicatedRefreshFixture { } } +private struct RankingPersistenceTrackState: Equatable, Sendable { + let identityKey: String + let rankingAppCount: Int? + let lastRefreshAt: Date? + let snapshotCount: Int +} + +private struct RankingPersistenceState: Equatable, Sendable { + let crawlCount: Int + let observationItemCount: Int + let snapshotCount: Int + let rankedResultCount: Int + let storeAppIDs: [Int64] + let storefrontMetadataCount: Int + let screenshotCount: Int + let latestRatingCount: Int + let dailyRatingCount: Int + let statsCount: Int + let statusCount: Int + let trackStates: [RankingPersistenceTrackState] +} + +private func rankingPersistenceState(in modelContext: ModelContext) throws -> RankingPersistenceState { + let tracks = try modelContext.fetch(FetchDescriptor()) + .map { + RankingPersistenceTrackState( + identityKey: $0.identityKey, + rankingAppCount: $0.rankingAppCount, + lastRefreshAt: $0.lastRefreshAt, + snapshotCount: $0.snapshots.count + ) + } + .sorted { $0.identityKey < $1.identityKey } + return RankingPersistenceState( + crawlCount: try modelContext.fetchCount(FetchDescriptor()), + observationItemCount: try modelContext.fetchCount(FetchDescriptor()), + snapshotCount: try modelContext.fetchCount(FetchDescriptor()), + rankedResultCount: try modelContext.fetchCount(FetchDescriptor()), + storeAppIDs: try modelContext.fetch(FetchDescriptor()).map(\.appStoreID).sorted(), + storefrontMetadataCount: try modelContext.fetchCount(FetchDescriptor()), + screenshotCount: try modelContext.fetchCount(FetchDescriptor()), + latestRatingCount: try modelContext.fetchCount(FetchDescriptor()), + dailyRatingCount: try modelContext.fetchCount(FetchDescriptor()), + statsCount: try modelContext.fetchCount(FetchDescriptor()), + statusCount: try modelContext.fetchCount(FetchDescriptor()), + trackStates: tracks + ) +} + +private final class FailingRankingPersistenceCheckpoint: Sendable { + private struct State: Sendable { + var callCount = 0 + } + + private let failingCall: Int + private let state = Mutex(State()) + + init(failingOnCall: Int) { + self.failingCall = failingOnCall + } + + func call() throws { + let shouldFail = state.withLock { state in + state.callCount += 1 + return state.callCount == failingCall + } + if shouldFail { + throw OpenASOError.unexpectedResponse + } + } + + func callCount() -> Int { + state.withLock { $0.callCount } + } +} + +private final class RankingMetadataEnrichmentRecorder: Sendable { + private let batches = Mutex<[[RankingMetadataEnrichmentRequest]]>([]) + + func record(_ requests: [RankingMetadataEnrichmentRequest]) { + batches.withLock { batches in + batches.append(requests) + } + } + + func recordedBatches() -> [[RankingMetadataEnrichmentRequest]] { + batches.withLock { batches in + batches + } + } +} + @MainActor private func makeDeduplicatedRefreshFixture(in modelContext: ModelContext) throws -> DeduplicatedRefreshFixture { let firstApp = TrackedApp( @@ -1749,6 +2313,68 @@ private actor QueryRankingProvider: SearchRankingProvider { } } +private actor GatedRankingProvider: SearchRankingProvider { + private var continuation: CheckedContinuation? + private var pendingPage: SearchRankingPage? + private var didStart = false + private var startWaiters: [CheckedContinuation] = [] + + func search( + keyword: String, + storefrontCode: String, + platform: AppPlatform, + limit: Int + ) async throws -> SearchRankingPage { + didStart = true + let waiters = startWaiters + startWaiters.removeAll() + for waiter in waiters { + waiter.resume() + } + if let pendingPage { + self.pendingPage = nil + return pendingPage + } + return await withCheckedContinuation { continuation in + self.continuation = continuation + } + } + + func waitUntilStarted() async { + guard !didStart else { return } + await withCheckedContinuation { continuation in + startWaiters.append(continuation) + } + } + + func succeed(_ page: SearchRankingPage) { + guard let continuation else { + pendingPage = page + return + } + self.continuation = nil + continuation.resume(returning: page) + } +} + +@MainActor +private final class RankingPendingEditInjector { + private let storefront: Storefront + private let pendingName: String + private(set) var injectionCount = 0 + + init(storefront: Storefront, pendingName: String) { + self.storefront = storefront + self.pendingName = pendingName + } + + func injectIfPagePersistenceCompleted(_ completed: Int) { + guard completed == 1, injectionCount == 0 else { return } + storefront.name = pendingName + injectionCount += 1 + } +} + @MainActor private final class StubRankingProvider: SearchRankingProvider { var page: SearchRankingPage