fix(analytics): prevent database locks during suspension

This commit is contained in:
Rocky
2026-08-21 21:13:13 +08:00
parent 7bd77509b6
commit ac374631ae
11 changed files with 599 additions and 155 deletions
+1
View File
@@ -17,6 +17,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- **Focused skill icons**: remove math, indices, arrows, shapes, commerce, keyboard, media, text-formatting, automotive, device, and variable-rendering categories from the custom-skill symbol picker. / **精简技能图标**:从自定义技能图标选择器中移除数学、索引、箭头、形状、商业、键盘、媒体、文本格式、汽车、设备与可变渲染分类。
### Fixed
- **Suspension-safe analytics**: batch uploads in the foreground or a system-granted refresh task, close SQLite outside authorized execution windows, and keep keyboard extensions recording-only to prevent system termination from database locks. / **安全挂起埋点**:仅在前台或系统授权的刷新任务中批量上传,在授权执行窗口之外关闭 SQLite,同时让键盘扩展仅负责记录,避免数据库锁触发系统终止。
- **Permanent invitation links**: load the account-scoped server invitation profile after sign-in, cache it per account, and keep sharing the exact stable server URL across refreshes and offline failures. / **永久邀请链接**:登录后异步加载账号级服务端邀请资料,按账号缓存,并在刷新或离线失败时持续分享服务端返回的固定链接。
- **Cross-process settings safety**: App Group updates now write only changed fields so a stale main-app or keyboard-extension snapshot cannot overwrite newer unrelated settings. / **跨进程设置安全**:App Group 更新现在只写入发生变化的字段,避免主 App 或键盘扩展的旧快照覆盖其他较新的设置。
- **Account-data freshness**: Settings and Account now share one background-refreshed snapshot, preserve cached content on failure, retry transient read errors, and isolate optional referral outages from core account and credit updates. / **账号数据新鲜度**:设置页与账号页现在共用同一份后台刷新快照,刷新失败时保留缓存内容,对瞬时读取错误自动重试,并避免可选邀请接口故障影响账号与积分更新。
+169 -80
View File
@@ -12,40 +12,6 @@ import OSGKeyboardShared
import OSLog
import UIKit
actor AnalyticsUploadSignal: AnalyticsUploadTriggering {
typealias Action = @Sendable () async -> Void
private var action: Action?
private var isScheduled = false
private let debounce: Duration
init(debounce: Duration = .milliseconds(250)) {
self.debounce = debounce
}
func install(_ action: @escaping Action) {
self.action = action
}
func requestUpload() {
guard !isScheduled, let action else { return }
isScheduled = true
Task {
try? await Task.sleep(for: debounce)
guard !Task.isCancelled else {
uploadFinished()
return
}
await action()
uploadFinished()
}
}
private func uploadFinished() {
isScheduled = false
}
}
final class HostAnalyticsBearerBridge: AnalyticsBearerProviding, @unchecked Sendable {
static let shared = HostAnalyticsBearerBridge()
@@ -112,6 +78,12 @@ struct HostAnalyticsLogger: AnalyticsLogging {
@MainActor
final class AnalyticsHostService: ObservableObject {
private enum RequestedDatabaseAccess {
case uninitialized
case foreground
case suspended
}
static let shared = AnalyticsHostService()
static let backgroundTaskIdentifier = "com.osgkeyboard.ios.analytics-sync"
@@ -127,16 +99,24 @@ final class AnalyticsHostService: ObservableObject {
qos: .utility
)
private var isMonitoringNetwork = false
private var foregroundUploadTask: Task<Void, Never>?
private var requestedDatabaseAccess = RequestedDatabaseAccess.uninitialized
private var databaseTransitionTask: Task<Void, Never>?
private var databaseTransitionGeneration = 0
private var firstOpenAcquisitionChannel = AnalyticsAcquisitionChannel.unknown
private var pendingAuthenticatedAccountID: UUID?
private var hasPendingAccountDeletion = false
private var backgroundTaskIdentifier: UIBackgroundTaskIdentifier = .invalid
init(
bearerProvider: any AnalyticsBearerProviding = HostAnalyticsBearerBridge.shared,
network: any AnalyticsNetworking = URLSessionAnalyticsNetwork()
network: any AnalyticsNetworking = URLSessionAnalyticsNetwork(),
uploadPolicy: AnalyticsUploadPolicy = .mobileDefault
) {
let signal = AnalyticsUploadSignal()
let signal = AnalyticsUploadSignal(policy: uploadPolicy)
uploadSignal = signal
let environment = Self.environment
let runtime = AnalyticsRuntime.mainApp(
environment: Self.environment,
environment: environment,
uploadConfiguration: AnalyticsUploadConfiguration(endpoint: Self.endpoint),
network: network,
bearerProvider: bearerProvider,
@@ -148,53 +128,64 @@ final class AnalyticsHostService: ObservableObject {
Task {
await signal.install {
await runtime.uploadCoordinator.uploadAvailableEvents(maximumBatches: 4)
guard !Task.isCancelled else { return }
await runtime.uploadCoordinator.uploadAvailableEvents(
maximumBatches: uploadPolicy.maximumBatches
)
}
}
}
func prepare(firstOpenAcquisitionChannel: AnalyticsAcquisitionChannel) {
self.firstOpenAcquisitionChannel = firstOpenAcquisitionChannel
startNetworkMonitoringIfNeeded()
Task {
await runtime.prepare(
firstOpenAcquisitionChannel: firstOpenAcquisitionChannel
)
let enabled = await runtime.isEnabled()
await MainActor.run {
self.isEnabled = enabled
}
client.recordSessionActivity()
await runtime.uploadCoordinator.uploadAvailableEvents(maximumBatches: 8)
}
}
func appDidBecomeActive() {
client.recordSessionActivity()
requestImmediateUpload(maximumBatches: 8)
requestForegroundDatabaseAccess()
}
func appWillResignActive() {
requestSuspendedDatabaseAccess()
}
func appDidEnterBackground() {
Self.scheduleBackgroundRefresh()
let identifier = UIApplication.shared.beginBackgroundTask(
withName: "analytics-sync"
)
guard identifier != .invalid else { return }
let task = Task {
await runtime.uploadCoordinator.uploadAvailableEvents(maximumBatches: 2)
await MainActor.run {
UIApplication.shared.endBackgroundTask(identifier)
requestSuspendedDatabaseAccess()
guard backgroundTaskIdentifier == .invalid else { return }
var identifier: UIBackgroundTaskIdentifier = .invalid
identifier = UIApplication.shared.beginBackgroundTask(
withName: "analytics-quiescence"
) { [weak self] in
Task { @MainActor in
self?.finishBackgroundTask(identifier)
}
}
guard identifier != .invalid else { return }
backgroundTaskIdentifier = identifier
let transitionTask = databaseTransitionTask
Task { [weak self] in
await transitionTask?.value
await MainActor.run {
self?.finishBackgroundTask(identifier)
}
}
foregroundUploadTask = task
}
func observeAuthenticatedAccount(_ accountID: UUID) async {
_ = await runtime.observeAccount(stableIdentifier: accountID.uuidString)
await runtime.uploadCoordinator.uploadAvailableEvents(maximumBatches: 8)
pendingAuthenticatedAccountID = accountID
let transitionTask = databaseTransitionTask
await transitionTask?.value
await processDeferredAccountWork()
}
func handleAccountDeletion() async {
await runtime.handleAccountDeletion()
hasPendingAccountDeletion = true
let transitionTask = databaseTransitionTask
await transitionTask?.value
await processDeferredAccountWork()
}
func setEnabled(_ enabled: Bool) {
@@ -205,7 +196,7 @@ final class AnalyticsHostService: ObservableObject {
self.isEnabled = current
}
if current {
await runtime.uploadCoordinator.uploadAvailableEvents(maximumBatches: 4)
await uploadSignal.requestActivationUpload()
}
}
}
@@ -219,15 +210,6 @@ final class AnalyticsHostService: ObservableObject {
}
}
func requestImmediateUpload(maximumBatches: Int = 4) {
foregroundUploadTask?.cancel()
foregroundUploadTask = Task {
await runtime.uploadCoordinator.uploadAvailableEvents(
maximumBatches: maximumBatches
)
}
}
static func registerBackgroundTask() {
BGTaskScheduler.shared.register(
forTaskWithIdentifier: backgroundTaskIdentifier,
@@ -238,10 +220,12 @@ final class AnalyticsHostService: ObservableObject {
return
}
let operation = Task {
await AnalyticsHostService.shared.runtime.uploadCoordinator
.uploadAvailableEvents(maximumBatches: 4)
refreshTask.setTaskCompleted(success: !Task.isCancelled)
scheduleBackgroundRefresh()
let success = await AnalyticsHostService.shared
.performBackgroundRefresh()
refreshTask.setTaskCompleted(success: success)
if success {
scheduleBackgroundRefresh()
}
}
refreshTask.expirationHandler = {
operation.cancel()
@@ -260,13 +244,118 @@ final class AnalyticsHostService: ObservableObject {
isMonitoringNetwork = true
pathMonitor.pathUpdateHandler = { path in
guard path.status == .satisfied else { return }
Task { @MainActor in
AnalyticsHostService.shared.requestImmediateUpload(maximumBatches: 8)
Task {
await AnalyticsHostService.shared.uploadSignal
.requestActivationUpload()
}
}
pathMonitor.start(queue: monitorQueue)
}
private func requestForegroundDatabaseAccess() {
guard requestedDatabaseAccess != .foreground else { return }
requestedDatabaseAccess = .foreground
databaseTransitionGeneration += 1
let generation = databaseTransitionGeneration
let previous = databaseTransitionTask
let runtime = runtime
let signal = uploadSignal
let client = client
let acquisitionChannel = firstOpenAcquisitionChannel
let task = Task { [weak self] in
await previous?.value
await runtime.repository.resumeDatabaseAccess()
await runtime.prepare(
firstOpenAcquisitionChannel: acquisitionChannel
)
let enabled = await runtime.isEnabled()
await self?.processDeferredAccountWork()
await signal.setActive(true)
client.recordSessionActivity()
await signal.requestActivationUpload()
await MainActor.run {
guard let self else { return }
self.isEnabled = enabled
self.finishDatabaseTransition(generation: generation)
}
}
databaseTransitionTask = task
}
private func requestSuspendedDatabaseAccess() {
guard requestedDatabaseAccess != .suspended else { return }
requestedDatabaseAccess = .suspended
databaseTransitionGeneration += 1
let generation = databaseTransitionGeneration
let previous = databaseTransitionTask
let runtime = runtime
let signal = uploadSignal
let task = Task { [weak self] in
await previous?.value
await signal.pauseAndWait()
await runtime.repository.suspendDatabaseAccess()
await MainActor.run {
self?.finishDatabaseTransition(generation: generation)
}
}
databaseTransitionTask = task
}
private func processDeferredAccountWork() async {
guard requestedDatabaseAccess == .foreground else { return }
if hasPendingAccountDeletion {
await runtime.handleAccountDeletion()
hasPendingAccountDeletion = false
}
guard let accountID = pendingAuthenticatedAccountID else { return }
let observation = await runtime.observeAccount(
stableIdentifier: accountID.uuidString
)
guard requestedDatabaseAccess == .foreground else { return }
if case .unavailable = observation {
return
}
pendingAuthenticatedAccountID = nil
await uploadSignal.requestActivationUpload()
}
private func performBackgroundRefresh() async -> Bool {
let transitionTask = databaseTransitionTask
await transitionTask?.value
guard !Task.isCancelled else { return false }
await runtime.repository.resumeDatabaseAccess()
await runtime.uploadCoordinator.uploadAvailableEvents(maximumBatches: 1)
let success = !Task.isCancelled
if requestedDatabaseAccess != .foreground {
await runtime.repository.suspendDatabaseAccess()
// A foreground transition may race the final close while this
// actor is re-entrant. Re-open if it won during the suspension.
if requestedDatabaseAccess == .foreground {
await runtime.repository.resumeDatabaseAccess()
}
}
return success
}
private func finishDatabaseTransition(generation: Int) {
guard databaseTransitionGeneration == generation else { return }
databaseTransitionTask = nil
}
private func finishBackgroundTask(_ identifier: UIBackgroundTaskIdentifier) {
guard identifier != .invalid,
backgroundTaskIdentifier == identifier else {
return
}
UIApplication.shared.endBackgroundTask(identifier)
backgroundTaskIdentifier = .invalid
}
private static let endpoint = URL(
string: "https://account.osglab.com/v1/analytics/events"
)!
@@ -0,0 +1,156 @@
// AnalyticsUploadScheduling.swift
// OSGKeyboard · Main App
//
// Count- and time-based foreground upload scheduling. Recording remains
// fire-and-forget; networking is single-flight and stops when the app resigns.
import Foundation
import OSGKeyboardShared
struct AnalyticsUploadPolicy: Sendable {
static let mobileDefault = Self(
eventThreshold: 20,
flushInterval: .seconds(60),
activationDelay: .seconds(10),
maximumBatches: 2
)
let eventThreshold: Int
let flushInterval: Duration
let activationDelay: Duration
let maximumBatches: Int
init(
eventThreshold: Int,
flushInterval: Duration,
activationDelay: Duration,
maximumBatches: Int
) {
self.eventThreshold = max(1, eventThreshold)
self.flushInterval = flushInterval
self.activationDelay = activationDelay
self.maximumBatches = max(1, maximumBatches)
}
}
actor AnalyticsUploadSignal: AnalyticsUploadTriggering {
typealias Action = @Sendable () async -> Void
private let policy: AnalyticsUploadPolicy
private var action: Action?
private var isActive = false
private var pendingSignals = 0
private var timerTask: Task<Void, Never>?
private var uploadTask: Task<Void, Never>?
init(policy: AnalyticsUploadPolicy = .mobileDefault) {
self.policy = policy
}
func install(_ action: @escaping Action) {
self.action = action
if isActive, pendingSignals > 0, uploadTask == nil {
schedule(after: policy.activationDelay, replacingTimer: true)
}
}
func requestUpload() {
pendingSignals = min(Int.max, pendingSignals + 1)
guard isActive, uploadTask == nil else { return }
if pendingSignals >= policy.eventThreshold {
schedule(after: .zero, replacingTimer: true)
} else {
schedule(after: policy.flushInterval)
}
}
/// Starts one bounded backlog drain after the foreground has settled.
func requestActivationUpload() {
pendingSignals = max(1, pendingSignals)
guard isActive, uploadTask == nil else { return }
schedule(after: policy.activationDelay, replacingTimer: true)
}
func setActive(_ active: Bool) async {
guard active != isActive else {
if active, pendingSignals > 0, uploadTask == nil {
schedule(after: policy.activationDelay)
}
return
}
isActive = active
if active {
if pendingSignals > 0 {
schedule(after: policy.activationDelay, replacingTimer: true)
}
} else {
await cancelAndWait()
}
}
/// Cancels timers and networking, then waits until the upload action has
/// unwound. Callers can keep a short UIKit background assertion while this
/// returns so no SQLite transaction survives process suspension.
func pauseAndWait() async {
isActive = false
await cancelAndWait()
}
private func schedule(
after delay: Duration,
replacingTimer: Bool = false
) {
guard isActive, action != nil, uploadTask == nil else { return }
if replacingTimer {
timerTask?.cancel()
timerTask = nil
}
guard timerTask == nil else { return }
timerTask = Task { [weak self] in
do {
try await Task.sleep(for: delay)
} catch {
return
}
await self?.timerFired()
}
}
private func timerFired() {
timerTask = nil
guard isActive,
uploadTask == nil,
pendingSignals > 0,
let action else {
return
}
pendingSignals = 0
uploadTask = Task { [weak self] in
await action()
await self?.uploadFinished()
}
}
private func uploadFinished() {
uploadTask = nil
guard isActive, pendingSignals > 0 else { return }
let delay: Duration = pendingSignals >= policy.eventThreshold
? .zero
: policy.flushInterval
schedule(after: delay, replacingTimer: true)
}
private func cancelAndWait() async {
timerTask?.cancel()
timerTask = nil
let task = uploadTask
task?.cancel()
await task?.value
uploadTask = nil
}
}
+5
View File
@@ -84,6 +84,9 @@ struct MainAppRoot: View {
analytics.prepare(
firstOpenAcquisitionChannel: firstOpenAcquisitionChannel
)
if scenePhase == .active {
analytics.appDidBecomeActive()
}
OSGDiag.log(
"MainAppRoot.onAppear scene=\(String(describing: scenePhase)) "
+ "onboarding=\(config.hasCompletedOnboarding) \(OSGDiag.memoryTag())",
@@ -152,6 +155,8 @@ struct MainAppRoot: View {
flowManager.handleScenePhase(phase)
if phase == .active {
analytics.appDidBecomeActive()
} else if phase == .inactive {
analytics.appWillResignActive()
} else if phase == .background {
analytics.appDidEnterBackground()
}
@@ -1,97 +1,32 @@
// AnalyticsExtensionService.swift
// OSGKeyboard · Keyboard Extension
//
// The extension records into the shared SQLite queue and only attempts one
// short anonymous batch when Full Access permits network use.
// The extension only records into durable queues. The host app is the sole
// network uploader, avoiding extension-lifecycle and cross-process lock races.
import Foundation
import OSGKeyboardShared
import OSLog
private actor ExtensionAnalyticsUploadSignal: AnalyticsUploadTriggering {
typealias Action = @Sendable () async -> Void
private var action: Action?
private var canUpload = false
private var isUploading = false
func install(_ action: @escaping Action) {
self.action = action
}
func setCanUpload(_ canUpload: Bool) {
self.canUpload = canUpload
}
func requestUpload() {
guard canUpload, !isUploading, let action else { return }
isUploading = true
Task {
await action()
uploadFinished()
}
}
private func uploadFinished() {
isUploading = false
}
}
private struct ExtensionAnalyticsLogger: AnalyticsLogging {
private let logger = Logger(
subsystem: Bundle.main.bundleIdentifier ?? "com.osgkeyboard.ios.keyboard",
category: "analytics"
)
func log(_ entry: AnalyticsUploadLogEntry) {
let statusCode = entry.statusCode ?? 0
let errorCategory = entry.errorCategory?.rawValue ?? "none"
logger.info(
"outcome=\(entry.outcome.rawValue, privacy: .public) count=\(entry.eventCount, privacy: .public) status=\(statusCode, privacy: .public) attempt=\(entry.attempt, privacy: .public) error=\(errorCategory, privacy: .public)"
)
}
}
final class AnalyticsExtensionService: Sendable {
static let shared = AnalyticsExtensionService()
let client: any AnalyticsClient
private let runtime: AnalyticsRuntime
private let uploadSignal: ExtensionAnalyticsUploadSignal
private init() {
let signal = ExtensionAnalyticsUploadSignal()
uploadSignal = signal
let environment = Self.environment
let runtime = AnalyticsRuntime.keyboardExtension(
environment: Self.environment,
uploadConfiguration: AnalyticsUploadConfiguration(endpoint: Self.endpoint),
trigger: signal,
logger: ExtensionAnalyticsLogger()
environment: environment,
uploadConfiguration: AnalyticsUploadConfiguration(endpoint: Self.endpoint)
)
self.runtime = runtime
client = runtime.client
Task {
await signal.install {
await runtime.uploadCoordinator.uploadAvailableEvents(maximumBatches: 1)
}
}
}
func recordPresentation(hasFullAccess: Bool) {
Task {
await uploadSignal.setCanUpload(hasFullAccess)
client.recordSessionActivity()
client.recordKeyboardActivated()
}
func recordPresentation(hasFullAccess _: Bool) {
client.recordSessionActivity()
client.recordKeyboardActivated()
}
func keyboardWillDisappear() {
Task {
await uploadSignal.setCanUpload(false)
}
}
func keyboardWillDisappear() {}
private static let endpoint = URL(
string: "https://account.osglab.com/v1/analytics/events"
@@ -8,6 +8,38 @@ import Foundation
import XCTest
final class AnalyticsExtensionPrivacyTests: XCTestCase {
func testKeyboardRuntimeRecordsWithoutAutomaticallyUploading() async throws {
let network = ExtensionAnalyticsNetwork()
let runtime = AnalyticsRuntime.keyboardExtension(
environment: AnalyticsEnvironment(appVersion: "2.0.0", osVersion: "26.0"),
repositoryConfiguration: AnalyticsRepositoryConfiguration(
databaseURL: try temporaryDatabaseURL()
),
uploadConfiguration: AnalyticsUploadConfiguration(
endpoint: URL(string: "https://analytics.test/v1/events")!
),
network: network,
wallClock: ExtensionAnalyticsClock(),
monotonicClock: ExtensionAnalyticsMonotonicClock(),
uuidGenerator: ExtensionAnalyticsUUIDGenerator(),
random: ExtensionAnalyticsRandomGenerator()
)
runtime.client.recordKeyboardActivated()
for _ in 0..<200 {
if await runtime.repository.debugSnapshot().pendingEvents.count == 1 {
break
}
try await Task.sleep(for: .milliseconds(10))
}
try await Task.sleep(for: .milliseconds(50))
let pendingSnapshot = await runtime.repository.debugSnapshot()
let requests = await network.requests()
XCTAssertEqual(pendingSnapshot.pendingEvents.count, 1)
XCTAssertTrue(requests.isEmpty)
}
func testKeyboardRuntimeUploadsWithoutAuthorizationOrSensitiveInputFields() async throws {
let databaseURL = try temporaryDatabaseURL()
let network = ExtensionAnalyticsNetwork()
@@ -174,7 +174,7 @@ public struct AnalyticsRepositoryConfiguration: Sendable {
maximumEventCount: Int = 10_000,
maximumStoredBytes: Int = 5 * 1_024 * 1_024,
eventRetention: TimeInterval = 34 * 24 * 60 * 60,
busyTimeoutMilliseconds: Int32 = 2_000
busyTimeoutMilliseconds: Int32 = 250
) {
self.databaseURL = databaseURL
self.maximumEventCount = max(1, maximumEventCount)
@@ -96,6 +96,7 @@ public actor AnalyticsRepository {
private let clock: any AnalyticsWallClock
private let uuidGenerator: any AnalyticsUUIDGenerating
private var database: SQLiteDatabase?
private var isDatabaseAccessSuspended = false
public init(
configuration: AnalyticsRepositoryConfiguration = .appGroupDefault(),
@@ -107,6 +108,19 @@ public actor AnalyticsRepository {
self.uuidGenerator = uuidGenerator
}
/// Re-enables process-local access before foreground or BGTask work.
public func resumeDatabaseAccess() {
isDatabaseAccessSuspended = false
}
/// Runs after earlier actor operations, then closes the process-local
/// connection. Later fire-and-forget records become best-effort no-ops
/// until an explicitly authorized execution window resumes access.
public func suspendDatabaseAccess() {
isDatabaseAccessSuspended = true
database = nil
}
/// The host calls this only after cold-start attribution has been resolved.
/// Installation identity, FIRST_OPEN and its marker share one transaction.
public func prepare(
@@ -926,6 +940,7 @@ public actor AnalyticsRepository {
// MARK: - Metadata and setup
private func openDatabaseIfNeeded() -> SQLiteDatabase? {
guard !isDatabaseAccessSuspended else { return nil }
if let database {
return database
}
@@ -452,6 +452,37 @@ final class AnalyticsRepositoryTests: XCTestCase {
eventTypes = await repository.debugSnapshot().pendingEvents.compactMap(\.eventType)
XCTAssertEqual(Set(eventTypes), [.firstOpen, .aiFeatureSucceeded])
}
func testSuspendedDatabaseDropsRecordsUntilExplicitResume() async throws {
let repository = AnalyticsRepository(
configuration: AnalyticsRepositoryConfiguration(
databaseURL: try analyticsTemporaryDatabaseURL()
),
clock: AnalyticsTestWallClock(),
uuidGenerator: AnalyticsTestUUIDGenerator()
)
await repository.record(
eventType: .keyboardActivated,
context: analyticsTestContext
)
await repository.suspendDatabaseAccess()
await repository.record(
eventType: .keyboardActivated,
context: analyticsTestContext
)
let suspendedSnapshot = await repository.debugSnapshot()
XCTAssertFalse(suspendedSnapshot.isAvailable)
await repository.resumeDatabaseAccess()
await repository.record(
eventType: .keyboardActivated,
context: analyticsTestContext
)
let resumedSnapshot = await repository.debugSnapshot()
XCTAssertTrue(resumedSnapshot.isAvailable)
XCTAssertEqual(resumedSnapshot.pendingEvents.count, 2)
}
}
private extension Array {
@@ -0,0 +1,179 @@
// AnalyticsUploadSchedulingTests.swift
// OSGKeyboardTests
//
// Foreground upload policy regression tests for batching and suspension safety.
@testable import OSGKeyboard
import XCTest
final class AnalyticsUploadSchedulingTests: XCTestCase {
func testCountThresholdCoalescesSignalsIntoOneUpload() async throws {
let probe = AnalyticsUploadProbe()
let signal = AnalyticsUploadSignal(
policy: AnalyticsUploadPolicy(
eventThreshold: 3,
flushInterval: .seconds(10),
activationDelay: .zero,
maximumBatches: 1
)
)
await signal.install {
await probe.recordUpload()
}
await signal.setActive(true)
await signal.requestUpload()
await signal.requestUpload()
try await Task.sleep(for: .milliseconds(30))
let countBeforeThreshold = await probe.uploadCount()
XCTAssertEqual(countBeforeThreshold, 0)
await signal.requestUpload()
try await waitUntil { await probe.uploadCount() == 1 }
let countAfterThreshold = await probe.uploadCount()
XCTAssertEqual(countAfterThreshold, 1)
}
func testIntervalFlushesAQueueBelowThreshold() async throws {
let probe = AnalyticsUploadProbe()
let signal = AnalyticsUploadSignal(
policy: AnalyticsUploadPolicy(
eventThreshold: 20,
flushInterval: .milliseconds(20),
activationDelay: .zero,
maximumBatches: 1
)
)
await signal.install {
await probe.recordUpload()
}
await signal.setActive(true)
await signal.requestUpload()
try await waitUntil { await probe.uploadCount() == 1 }
let uploadCount = await probe.uploadCount()
XCTAssertEqual(uploadCount, 1)
}
func testInactiveSignalDefersUntilNextActivation() async throws {
let probe = AnalyticsUploadProbe()
let signal = AnalyticsUploadSignal(
policy: AnalyticsUploadPolicy(
eventThreshold: 2,
flushInterval: .milliseconds(20),
activationDelay: .milliseconds(20),
maximumBatches: 1
)
)
await signal.install {
await probe.recordUpload()
}
await signal.requestUpload()
await signal.requestUpload()
try await Task.sleep(for: .milliseconds(40))
let inactiveUploadCount = await probe.uploadCount()
XCTAssertEqual(inactiveUploadCount, 0)
await signal.setActive(true)
try await waitUntil { await probe.uploadCount() == 1 }
let activeUploadCount = await probe.uploadCount()
XCTAssertEqual(activeUploadCount, 1)
}
func testLateActionInstallationDoesNotLoseActivationUpload() async throws {
let probe = AnalyticsUploadProbe()
let signal = AnalyticsUploadSignal(
policy: AnalyticsUploadPolicy(
eventThreshold: 20,
flushInterval: .seconds(10),
activationDelay: .milliseconds(20),
maximumBatches: 1
)
)
await signal.setActive(true)
await signal.requestActivationUpload()
await signal.install {
await probe.recordUpload()
}
try await waitUntil { await probe.uploadCount() == 1 }
let uploadCount = await probe.uploadCount()
XCTAssertEqual(uploadCount, 1)
}
func testPauseCancelsAndWaitsForInFlightUpload() async throws {
let probe = AnalyticsUploadProbe()
let signal = AnalyticsUploadSignal(
policy: AnalyticsUploadPolicy(
eventThreshold: 1,
flushInterval: .seconds(10),
activationDelay: .zero,
maximumBatches: 1
)
)
await signal.install {
await probe.recordStarted()
do {
try await Task.sleep(for: .seconds(30))
} catch {
await probe.recordCancellation()
}
}
await signal.setActive(true)
await signal.requestUpload()
try await waitUntil { await probe.didStart() }
await signal.pauseAndWait()
let wasCancelled = await probe.wasCancelled()
XCTAssertTrue(wasCancelled)
}
private func waitUntil(
timeout: Duration = .seconds(1),
condition: @escaping @Sendable () async -> Bool
) async throws {
let clock = ContinuousClock()
let deadline = clock.now.advanced(by: timeout)
while clock.now < deadline {
if await condition() {
return
}
try await Task.sleep(for: .milliseconds(10))
}
XCTFail("Timed out waiting for analytics upload state")
}
}
private actor AnalyticsUploadProbe {
private var uploads = 0
private var started = false
private var cancelled = false
func recordUpload() {
uploads += 1
}
func uploadCount() -> Int {
uploads
}
func recordStarted() {
started = true
}
func didStart() -> Bool {
started
}
func recordCancellation() {
cancelled = true
}
func wasCancelled() -> Bool {
cancelled
}
}
+1
View File
@@ -47,6 +47,7 @@
"OSGKeyboardTests/AnalyticsModelTests",
"OSGKeyboardTests/AnalyticsRepositoryTests",
"OSGKeyboardTests/AnalyticsUploadCoordinatorTests",
"OSGKeyboardTests/AnalyticsUploadSchedulingTests",
"OSGKeyboardTests/AnalyticsAIOperationTests",
"OSGKeyboardExtTests/AnalyticsExtensionPrivacyTests"
]