Commit a5aac34f1fa2ab79d5f4db54d28aa32ba949f961

Authored by 权海
1 parent 60476ea8

feat(ui):添加新数据表和数据上传

import Foundation
import SQLite3
enum HealthRawStressUploadKind {
case hrv
case realtimeStress
case dailyStress
}
enum HealthRawStressSQLiteUploadError: LocalizedError {
case invalidDatabasePath
case openDatabaseFailed(String)
case prepareFailed(String)
case invalidServerURL
case missingAccessToken
case invalidResponse
case requestFailed(path: String, statusCode: Int, body: String?)
var errorDescription: String? {
switch self {
case .invalidDatabasePath:
return "Invalid HealthRawData sqlite file path."
case .openDatabaseFailed(let message):
return "Open HealthRawData sqlite failed: \(message)"
case .prepareFailed(let message):
return "Prepare HealthRawData sqlite query failed: \(message)"
case .invalidServerURL:
return "Invalid server URL."
case .missingAccessToken:
return "Missing access token."
case .invalidResponse:
return "Invalid upload response."
case .requestFailed(let path, let statusCode, let body):
return "Upload \(path) failed status=\(statusCode) body=\(body ?? "nil")"
}
}
}
final class HealthRawStressSQLiteUploader {
static let shared = HealthRawStressSQLiteUploader()
private let session: URLSession
private let batchSize = 200
init(session: URLSession = .shared) {
self.session = session
}
func uploadHrv(sqliteFilePath: String) async throws -> Int64 {
let rows = try queryRows(
sqliteFilePath: sqliteFilePath,
sql: """
SELECT raw_end_time, raw_hrv, result, state, baseline_resting_hr, baseline_hrv,
is_sleep_likely, is_workout, is_workout_recovery, is_suspected_activity
FROM hrv_results
WHERE uploaded != 1
ORDER BY raw_end_time ASC
"""
)
var uploadedUntil: Int64 = 0
for batch in batches(rows) {
let list = batch.map { row in
[
"data_time": row.int64("raw_end_time"),
"raw_hrv": row.double("raw_hrv"),
"trend_hrv": row.double("result"),
"state": row.int("state"),
"hr_baseline": row.double("baseline_resting_hr"),
"hrv_baseline": row.double("baseline_hrv"),
"is_asleep": row.int("is_sleep_likely"),
"is_workout": row.int("is_workout"),
"is_workout_recovery": row.int("is_workout_recovery"),
"is_suspected_activity": row.int("is_suspected_activity"),
] as [String: Any]
}
try await upload(
path: "/client/doublefeel/health/v2/hrv_trend/",
body: ["data_list": list]
)
uploadedUntil = batch.last?.int64("raw_end_time") ?? uploadedUntil
}
return uploadedUntil
}
func uploadRealtimeStress(sqliteFilePath: String) async throws -> Int64 {
let rows = try queryRows(
sqliteFilePath: sqliteFilePath,
sql: """
SELECT raw_end_time, raw_hr, result, is_sleep_likely, is_workout,
is_workout_recovery, is_suspected_activity
FROM realtime_stress_results
WHERE uploaded != 1
ORDER BY raw_end_time ASC
"""
)
var uploadedUntil: Int64 = 0
for batch in batches(rows) {
let list = batch.map { row in
let stressValue = row.double("result")
return [
"data_time": row.int64("raw_end_time"),
"hr_value": row.double("raw_hr"),
"stress_value": stressValue,
"state": stressState(stressValue),
"is_asleep": row.int("is_sleep_likely"),
"is_workout": row.int("is_workout"),
"is_workout_recovery": row.int("is_workout_recovery"),
"is_suspected_activity": row.int("is_suspected_activity"),
] as [String: Any]
}
try await upload(
path: "/client/doublefeel/health/v2/realtime_stress/",
body: ["data_list": list]
)
uploadedUntil = batch.last?.int64("raw_end_time") ?? uploadedUntil
}
return uploadedUntil
}
func uploadDailyStress(sqliteFilePath: String) async throws -> Bool {
let rows = try queryRows(
sqliteFilePath: sqliteFilePath,
sql: """
SELECT stress_value, stress_score, state, data_time
FROM daily_stress_results
WHERE uploaded != 1
ORDER BY date ASC
"""
)
for batch in batches(rows) {
let list = batch.map { row in
[
"stress_value": row.double("stress_value"),
"stress_score": row.int("stress_score"),
"state": row.int("state"),
"data_time": row.int64("data_time"),
] as [String: Any]
}
try await upload(
path: "/client/doublefeel/health/v2/stress_score/",
body: ["data_list": list]
)
}
return true
}
private func queryRows(sqliteFilePath: String, sql: String) throws -> [SQLiteRow] {
guard FileManager.default.fileExists(atPath: sqliteFilePath) else {
throw HealthRawStressSQLiteUploadError.invalidDatabasePath
}
var db: OpaquePointer?
let flags = SQLITE_OPEN_READONLY | SQLITE_OPEN_FULLMUTEX
guard sqlite3_open_v2(sqliteFilePath, &db, flags, nil) == SQLITE_OK,
let db else {
let message = db.map { String(cString: sqlite3_errmsg($0)) } ?? "unknown"
if let db { sqlite3_close(db) }
throw HealthRawStressSQLiteUploadError.openDatabaseFailed(message)
}
defer { sqlite3_close(db) }
var statement: OpaquePointer?
guard sqlite3_prepare_v2(db, sql, -1, &statement, nil) == SQLITE_OK,
let statement else {
throw HealthRawStressSQLiteUploadError.prepareFailed(
String(cString: sqlite3_errmsg(db))
)
}
defer { sqlite3_finalize(statement) }
var rows: [SQLiteRow] = []
while sqlite3_step(statement) == SQLITE_ROW {
var values: [String: SQLiteValue] = [:]
for index in 0..<sqlite3_column_count(statement) {
let name = String(cString: sqlite3_column_name(statement, index))
values[name] = SQLiteValue(statement: statement, index: index)
}
rows.append(SQLiteRow(values: values))
}
return rows
}
private func batches(_ values: [SQLiteRow]) -> [[SQLiteRow]] {
guard batchSize > 0 else { return [values] }
return stride(from: 0, to: values.count, by: batchSize).map {
Array(values[$0..<Swift.min($0 + batchSize, values.count)])
}
}
private func upload(path: String, body: [String: Any]) async throws {
guard let baseURL = URL(string: AppShared.shared.baseUrl),
let url = URL(string: path, relativeTo: baseURL)?.absoluteURL else {
throw HealthRawStressSQLiteUploadError.invalidServerURL
}
guard let accessToken = AppShared.shared.token, !accessToken.isEmpty else {
throw HealthRawStressSQLiteUploadError.missingAccessToken
}
var request = URLRequest(url: url)
request.httpMethod = "POST"
request.timeoutInterval = 60
request.setValue("application/json", forHTTPHeaderField: "Accept")
request.setValue("application/json", forHTTPHeaderField: "Content-Type")
request.setValue(accessToken, forHTTPHeaderField: "access_token")
request.setValue(AppShared.shared.agent.finalUA, forHTTPHeaderField: "User-Agent")
request.httpBody = try JSONSerialization.data(withJSONObject: body)
let (data, response) = try await session.data(for: request)
guard let httpResponse = response as? HTTPURLResponse else {
throw HealthRawStressSQLiteUploadError.invalidResponse
}
guard (200..<300).contains(httpResponse.statusCode) else {
if httpResponse.statusCode == 401 {
await MainActor.run { AppShared.shared.logout() }
}
throw HealthRawStressSQLiteUploadError.requestFailed(
path: path,
statusCode: httpResponse.statusCode,
body: String(data: data, encoding: .utf8)
)
}
}
private func stressState(_ value: Double) -> Int {
if value >= 81 { return 1 }
if value >= 61 { return 2 }
if value >= 21 { return 3 }
return 4
}
}
private enum SQLiteValue {
case integer(Int64)
case double(Double)
case text(String)
case null
init(statement: OpaquePointer, index: Int32) {
switch sqlite3_column_type(statement, index) {
case SQLITE_INTEGER:
self = .integer(sqlite3_column_int64(statement, index))
case SQLITE_FLOAT:
self = .double(sqlite3_column_double(statement, index))
case SQLITE_TEXT:
self = .text(String(cString: sqlite3_column_text(statement, index)))
default:
self = .null
}
}
}
private struct SQLiteRow {
let values: [String: SQLiteValue]
func int64(_ key: String) -> Int64 {
switch values[key] {
case .integer(let value): return value
case .double(let value): return Int64(value)
case .text(let value): return Int64(value) ?? 0
case .null, .none: return 0
}
}
func int(_ key: String) -> Int {
Int(int64(key))
}
func double(_ key: String) -> Double {
switch values[key] {
case .integer(let value): return Double(value)
case .double(let value): return value
case .text(let value): return Double(value) ?? 0
case .null, .none: return 0
}
}
}
... ...
... ... @@ -15,6 +15,7 @@ final class HealthKitRawDataHostApiImpl: HealthKitRawDataHostApi {
}
func performHealthDataUpload(completion: @escaping (Result<Bool, any Error>) -> Void) {
print("trigger HealthDataUpload")
Task {
let summary = await AnchoredHealthDataUploader.shared.uploadAll()
completion(.success(true))
... ... @@ -22,15 +23,39 @@ final class HealthKitRawDataHostApiImpl: HealthKitRawDataHostApi {
}
func performHRDataUpload(sqliteFilePath: String, completion: @escaping (Result<Int64, any Error>) -> Void) {
print("trigger HRDataUpload")
Task {
do {
let uploadedUntil = try await HealthRawStressSQLiteUploader.shared.uploadRealtimeStress(sqliteFilePath: sqliteFilePath)
completion(.success(uploadedUntil))
} catch {
completion(.failure(error))
}
}
}
func performHRVDataUpload(sqliteFilePath: String, completion: @escaping (Result<Int64, any Error>) -> Void) {
print("trigger HRVDataUpload")
Task {
do {
let uploadedUntil = try await HealthRawStressSQLiteUploader.shared.uploadHrv(sqliteFilePath: sqliteFilePath)
completion(.success(uploadedUntil))
} catch {
completion(.failure(error))
}
}
}
func performAvgRealtimeStressDataUpload(sqliteFilePath: String, completion: @escaping (Result<Bool, any Error>) -> Void) {
print("trigger AvgRealtimeStressDataUpload")
Task {
do {
let success = try await HealthRawStressSQLiteUploader.shared.uploadDailyStress(sqliteFilePath: sqliteFilePath)
completion(.success(success))
} catch {
completion(.failure(error))
}
}
}
... ... @@ -81,7 +106,7 @@ final class HealthKitRawDataHostApiImpl: HealthKitRawDataHostApi {
) {
let startDate = Date(timeIntervalSince1970: TimeInterval(startTime))
let endDate = Date(timeIntervalSince1970: TimeInterval(endTime))
print("trigger getHealthKitRawSleepData from: \(startDate) to \(endDate)")
Task {
do {
let intervals = try await service.fetchSleepData(
... ... @@ -93,6 +118,7 @@ final class HealthKitRawDataHostApiImpl: HealthKitRawDataHostApi {
completion(.success([]))
return
}
print("trigger getHealthKitRawSleepData back: \(points)")
completion(.success([
HealthKitRawSleepDataPoint(
dataType: Int64(NativeHealthDataType.sleep.rawValue),
... ...
... ... @@ -25,22 +25,29 @@ class HealthRawDataCoreService {
HealthRawStressLocalStore? localStore,
AppEnvironmentConfig? environmentConfig,
int Function()? userIdProvider,
bool uploadResultsAfterCalculation = true,
}) : _healthApi = healthApi ?? HealthKitHostApi(),
_rawDataApi = rawDataApi ?? HealthKitRawDataHostApi(),
_localStore = localStore ?? HealthRawStressLocalStore(),
_environmentConfig = environmentConfig,
_userIdProvider = userIdProvider;
_userIdProvider = userIdProvider,
_uploadResultsAfterCalculation = uploadResultsAfterCalculation;
final HealthKitHostApi _healthApi;
final HealthKitRawDataHostApi _rawDataApi;
final HealthRawStressLocalStore _localStore;
final AppEnvironmentConfig? _environmentConfig;
final int Function()? _userIdProvider;
final bool _uploadResultsAfterCalculation;
final StreamController<HealthRawDataUpdatedEvent>
_healthDataUpdatedController =
StreamController<HealthRawDataUpdatedEvent>.broadcast();
Future<void>? _healthDataUpdatedCalculation;
Future<HealthRawStressCalculationResult>? _coreCalculation;
List<int> _pendingHealthDataUpdatedTypes = const <int>[];
bool _isUploadingHrvResults = false;
bool _isUploadingRealtimeStressResults = false;
bool _isUploadingDailyStressResults = false;
int get _userId {
final userId = _userIdProvider?.call() ?? 0;
... ... @@ -127,7 +134,11 @@ class HealthRawDataCoreService {
int? endTime,
int readChunkDays = defaultReadChunkDays,
}) async {
return _runWithLog(
final running = _coreCalculation;
if (running != null) {
return running;
}
final task = _runWithLog(
action: 'startCoreCaculate',
body: () => _readCalculateAndStore(
endTime: endTime,
... ... @@ -135,6 +146,14 @@ class HealthRawDataCoreService {
forceStartTime: null,
),
);
_coreCalculation = task;
try {
return await task;
} finally {
if (identical(_coreCalculation, task)) {
_coreCalculation = null;
}
}
}
Future<HealthRawStressCalculationResult> syncAndStore({
... ... @@ -258,11 +277,18 @@ class HealthRawDataCoreService {
endTime: effectiveEndTime,
);
await _localStore.upsertResult(result);
if (_uploadResultsAfterCalculation) {
await _uploadHrvResults();
await _uploadRealtimeStressResults();
}
final dailyStressPoints = await _calculateAndStoreDailyStressPoints(
userId: userId,
realtimePoints: result.realtimeStressPoints,
nowSeconds: effectiveEndTime,
);
if (_uploadResultsAfterCalculation) {
await _uploadDailyStressResults();
}
final storedResult = result.copyWith(dailyStressPoints: dailyStressPoints);
if (isFirstCalculation && (_environmentConfig?.isDebug ?? false)) {
AppToast.show('首次计算完成');
... ... @@ -450,6 +476,78 @@ class HealthRawDataCoreService {
);
}
Future<void> _uploadHrvResults() async {
if (_isUploadingHrvResults) {
return;
}
if (!await _localStore.hasPendingHrvStressUploads(userId: _userId)) {
return;
}
_isUploadingHrvResults = true;
try {
final uploadedUntil = await _rawDataApi.performHRVDataUpload(
sqliteFilePath: await _localStore.dbPath(_userId),
);
if (uploadedUntil > 0) {
await _localStore.markHrvStressUploadedUntil(
userId: _userId,
rawEndTime: uploadedUntil,
);
}
} catch (error, stackTrace) {
_logError('upload hrv results failed', error, stackTrace);
} finally {
_isUploadingHrvResults = false;
}
}
Future<void> _uploadRealtimeStressResults() async {
if (_isUploadingRealtimeStressResults) {
return;
}
if (!await _localStore.hasPendingRealtimeStressUploads(userId: _userId)) {
return;
}
_isUploadingRealtimeStressResults = true;
try {
final uploadedUntil = await _rawDataApi.performHRDataUpload(
sqliteFilePath: await _localStore.dbPath(_userId),
);
if (uploadedUntil > 0) {
await _localStore.markRealtimeStressUploadedUntil(
userId: _userId,
rawEndTime: uploadedUntil,
);
}
} catch (error, stackTrace) {
_logError('upload realtime stress results failed', error, stackTrace);
} finally {
_isUploadingRealtimeStressResults = false;
}
}
Future<void> _uploadDailyStressResults() async {
if (_isUploadingDailyStressResults) {
return;
}
if (!await _localStore.hasPendingDailyStressUploads(userId: _userId)) {
return;
}
_isUploadingDailyStressResults = true;
try {
final success = await _rawDataApi.performAvgRealtimeStressDataUpload(
sqliteFilePath: await _localStore.dbPath(_userId),
);
if (success) {
await _localStore.markDailyStressUploaded(userId: _userId);
}
} catch (error, stackTrace) {
_logError('upload daily stress results failed', error, stackTrace);
} finally {
_isUploadingDailyStressResults = false;
}
}
HealthRawStressCalculationResult calculate({
required List<HealthKitRawDataPoint> hrvPoints,
required List<HealthKitRawDataPoint> heartRatePoints,
... ... @@ -833,6 +931,7 @@ class HealthRawRealtimeStressPoint {
const HealthRawRealtimeStressPoint({
required this.userId,
required this.rawEndTime,
required this.rawHr,
required this.result,
required this.sourceStartTime,
required this.sourceEndTime,
... ... @@ -842,6 +941,7 @@ class HealthRawRealtimeStressPoint {
final int userId;
final int rawEndTime;
final double rawHr;
final double result;
final int sourceStartTime;
final int sourceEndTime;
... ... @@ -858,6 +958,7 @@ class HealthRawRealtimeStressPoint {
return HealthRawRealtimeStressPoint(
userId: row['user_id'] as int,
rawEndTime: row['raw_end_time'] as int,
rawHr: (row['raw_hr'] as num?)?.toDouble() ?? 0,
result: (row['result'] as num).toDouble(),
sourceStartTime: row['source_start_time'] as int,
sourceEndTime: row['source_end_time'] as int,
... ... @@ -1231,6 +1332,18 @@ class HealthRawStressLocalStore {
);
}
Future<void> markHrvStressUploadedUntil({
required int userId,
required int rawEndTime,
}) {
return _markUploadedUntil(
userId: userId,
table: hrvResultsTable,
timeColumn: 'raw_end_time',
time: rawEndTime,
);
}
Future<void> markRealtimeStressUploaded({
required int userId,
required Iterable<int> rawEndTimes,
... ... @@ -1242,6 +1355,40 @@ class HealthRawStressLocalStore {
);
}
Future<void> markRealtimeStressUploadedUntil({
required int userId,
required int rawEndTime,
}) {
return _markUploadedUntil(
userId: userId,
table: realtimeStressResultsTable,
timeColumn: 'raw_end_time',
time: rawEndTime,
);
}
Future<void> markDailyStressUploaded({required int userId}) async {
final db = await _database(userId);
await db.update(
dailyStressResultsTable,
{'uploaded': 1},
where: 'uploaded != 1',
);
}
Future<bool> hasPendingHrvStressUploads({required int userId}) {
return _hasPendingUploads(userId: userId, table: hrvResultsTable);
}
Future<bool> hasPendingRealtimeStressUploads({required int userId}) {
return _hasPendingUploads(
userId: userId, table: realtimeStressResultsTable);
}
Future<bool> hasPendingDailyStressUploads({required int userId}) {
return _hasPendingUploads(userId: userId, table: dailyStressResultsTable);
}
Future<String> dbPath(int userId) async {
final dir = _rootDirectory ?? await getApplicationDocumentsDirectory();
return '${dir.path}/hrv_result_$userId.sqlite';
... ... @@ -1282,7 +1429,7 @@ class HealthRawStressLocalStore {
final db = await factory.openDatabase(
path,
options: OpenDatabaseOptions(
version: 5,
version: 6,
onCreate: (db, version) async {
await _createTables(db);
},
... ... @@ -1300,6 +1447,9 @@ class HealthRawStressLocalStore {
if (oldVersion < 5) {
await _createDailyStressTable(db);
}
if (oldVersion < 6) {
await _addRealtimeRawHrColumn(db);
}
},
),
);
... ... @@ -1332,6 +1482,7 @@ CREATE TABLE IF NOT EXISTS $hrvResultsTable (
CREATE TABLE IF NOT EXISTS $realtimeStressResultsTable (
raw_end_time INTEGER PRIMARY KEY,
user_id INTEGER NOT NULL,
raw_hr REAL NOT NULL DEFAULT 0,
result REAL NOT NULL,
source_start_time INTEGER NOT NULL,
source_end_time INTEGER NOT NULL,
... ... @@ -1408,6 +1559,16 @@ CREATE TABLE IF NOT EXISTS $dailyStressResultsTable (
}
}
Future<void> _addRealtimeRawHrColumn(DatabaseExecutor db) async {
try {
await db.execute(
'ALTER TABLE $realtimeStressResultsTable ADD COLUMN raw_hr REAL NOT NULL DEFAULT 0',
);
} on DatabaseException catch (error) {
if (!error.isDuplicateColumnError()) rethrow;
}
}
Future<int?> _latestSourceStartTime(int userId, String table) async {
final db = await _database(userId);
final rows = await db.query(
... ... @@ -1436,6 +1597,36 @@ CREATE TABLE IF NOT EXISTS $dailyStressResultsTable (
);
}
Future<void> _markUploadedUntil({
required int userId,
required String table,
required String timeColumn,
required int time,
}) async {
if (time <= 0) return;
final db = await _database(userId);
await db.update(
table,
{'uploaded': 1},
where: '$timeColumn <= ? AND uploaded != 1',
whereArgs: [time],
);
}
Future<bool> _hasPendingUploads({
required int userId,
required String table,
}) async {
final db = await _database(userId);
final rows = await db.query(
table,
columns: ['uploaded'],
where: 'uploaded != 1',
limit: 1,
);
return rows.isNotEmpty;
}
Map<String, Object?> _hrvRow(HealthRawHrvStressPoint point) {
return <String, Object?>{
'raw_end_time': point.rawEndTime,
... ... @@ -1458,6 +1649,7 @@ CREATE TABLE IF NOT EXISTS $dailyStressResultsTable (
return <String, Object?>{
'raw_end_time': point.rawEndTime,
'user_id': point.userId,
'raw_hr': point.rawHr,
'result': point.result,
'source_start_time': point.sourceStartTime,
'source_end_time': point.sourceEndTime,
... ... @@ -1492,16 +1684,32 @@ CREATE TABLE IF NOT EXISTS $dailyStressResultsTable (
String table,
Map<String, Object?> row,
) async {
await db.insert(
final existing = await db.query(
table,
where: 'raw_end_time = ?',
whereArgs: [row['raw_end_time']],
limit: 1,
);
if (existing.isEmpty) {
await db.insert(table, row);
return;
}
final existingRow = existing.first;
if (_matchesStoredValues(
existingRow, row, _rawResultStoredValueKeys(row))) {
return;
}
final uploadPayloadChanged = !_matchesStoredValues(
existingRow,
row,
conflictAlgorithm: ConflictAlgorithm.ignore,
_rawResultUploadValueKeys(row),
);
await db.update(
table,
<String, Object?>{
'user_id': row['user_id'],
if (row.containsKey('raw_hrv')) 'raw_hrv': row['raw_hrv'],
if (row.containsKey('raw_hr')) 'raw_hr': row['raw_hr'],
'result': row['result'],
'source_start_time': row['source_start_time'],
'source_end_time': row['source_end_time'],
... ... @@ -1518,7 +1726,7 @@ CREATE TABLE IF NOT EXISTS $dailyStressResultsTable (
'is_workout_recovery': row['is_workout_recovery'],
'is_sleep_likely': row['is_sleep_likely'],
'is_suspected_activity': row['is_suspected_activity'],
'uploaded': 0,
'uploaded': uploadPayloadChanged ? 0 : existingRow['uploaded'],
},
where: 'raw_end_time = ?',
whereArgs: [row['raw_end_time']],
... ... @@ -1529,10 +1737,21 @@ CREATE TABLE IF NOT EXISTS $dailyStressResultsTable (
DatabaseExecutor db,
Map<String, Object?> row,
) async {
await db.insert(
final existing = await db.query(
dailyStressResultsTable,
where: 'date = ?',
whereArgs: [row['date']],
limit: 1,
);
if (existing.isEmpty) {
await db.insert(dailyStressResultsTable, row);
return;
}
final existingRow = existing.first;
final valueChanged = !_matchesStoredValues(
existingRow,
row,
conflictAlgorithm: ConflictAlgorithm.ignore,
const ['user_id', 'stress_value', 'stress_score', 'state'],
);
await db.update(
dailyStressResultsTable,
... ... @@ -1542,10 +1761,68 @@ CREATE TABLE IF NOT EXISTS $dailyStressResultsTable (
'stress_score': row['stress_score'],
'state': row['state'],
'data_time': row['data_time'],
'uploaded': 0,
'uploaded': valueChanged ? 0 : existingRow['uploaded'],
},
where: 'date = ?',
whereArgs: [row['date']],
);
}
List<String> _rawResultStoredValueKeys(Map<String, Object?> row) {
return <String>[
'user_id',
if (row.containsKey('raw_hrv')) 'raw_hrv',
if (row.containsKey('raw_hr')) 'raw_hr',
'result',
'source_start_time',
'source_end_time',
if (row.containsKey('state')) 'state',
if (row.containsKey('baseline_hrv')) 'baseline_hrv',
if (row.containsKey('baseline_awake_hrv')) 'baseline_awake_hrv',
if (row.containsKey('baseline_sleep_hrv')) 'baseline_sleep_hrv',
if (row.containsKey('baseline_resting_hr')) 'baseline_resting_hr',
'is_workout',
'is_workout_recovery',
'is_sleep_likely',
'is_suspected_activity',
];
}
List<String> _rawResultUploadValueKeys(Map<String, Object?> row) {
return <String>[
'user_id',
if (row.containsKey('raw_hrv')) 'raw_hrv',
if (row.containsKey('raw_hr')) 'raw_hr',
'result',
if (row.containsKey('state')) 'state',
if (row.containsKey('baseline_hrv')) 'baseline_hrv',
if (row.containsKey('baseline_awake_hrv')) 'baseline_awake_hrv',
if (row.containsKey('baseline_sleep_hrv')) 'baseline_sleep_hrv',
if (row.containsKey('baseline_resting_hr')) 'baseline_resting_hr',
'is_workout',
'is_workout_recovery',
'is_sleep_likely',
'is_suspected_activity',
];
}
bool _matchesStoredValues(
Map<String, Object?> existing,
Map<String, Object?> next,
Iterable<String> keys,
) {
for (final key in keys) {
if (!_storedValueEquals(existing[key], next[key])) {
return false;
}
}
return true;
}
bool _storedValueEquals(Object? a, Object? b) {
if (a is num && b is num) {
return (a.toDouble() - b.toDouble()).abs() < 0.000001;
}
return a == b;
}
}
... ...
... ... @@ -365,6 +365,7 @@ class HealthRawStressCalculator {
HealthRawRealtimeStressPoint(
userId: userId,
rawEndTime: currentRawPoint.endTime,
rawHr: currentRawPoint.value ?? 0,
result: stress,
sourceStartTime: sourceRange.start,
sourceEndTime: sourceRange.end,
... ...
... ... @@ -9,6 +9,7 @@ import 'package:doublefeel_flutter/data/models/health/sleep/sleep_statistics_dat
import 'package:doublefeel_flutter/pigeon/health_kit_raw_data_api.g.dart';
const _sleepTypeInBed = 0;
const _sleepTypeAsleep = 1;
const _sleepTypeAsleepUnspecified = 1;
const _sleepTypeAwake = 2;
const _sleepTypeAsleepCore = 3;
... ... @@ -499,7 +500,7 @@ class LocalHealthDataConvert {
final sleepEvaluate = _sleepEvaluate(score);
final sleepDuration = nullableScore == null || sleepEvaluate == null
? 0
: windowEnd - windowStart;
: summary.totalAsleepMinutes.round() * 60;
return _SleepDaySummary(
day: day,
... ... @@ -748,11 +749,10 @@ class LocalHealthDataConvert {
remMinutes += minutes;
asleepMinutes += minutes;
break;
case _sleepTypeAsleepUnspecified:
asleepMinutes += minutes;
break;
default:
asleepMinutes += minutes;
if (_isAsleepSleepType(interval.dataType)) {
asleepMinutes += minutes;
}
break;
}
}
... ... @@ -815,6 +815,14 @@ class LocalHealthDataConvert {
return 3;
}
static bool _isAsleepSleepType(int dataType) {
return dataType == _sleepTypeAsleep ||
dataType == _sleepTypeAsleepUnspecified ||
dataType == _sleepTypeAsleepCore ||
dataType == _sleepTypeAsleepDeep ||
dataType == _sleepTypeAsleepRem;
}
static int? _modeState(Iterable<int> states) {
final counts = <int, int>{};
for (final state in states) {
... ...
... ... @@ -42,6 +42,7 @@ void main() {
rawDataApi: api,
localStore: store,
userIdProvider: () => 42,
uploadResultsAfterCalculation: false,
);
final first = await service.startCoreCaculate(
... ... @@ -99,10 +100,10 @@ void main() {
realtime.map((e) => e.rawEndTime),
[base + 60, base + 420, base + 720],
);
expect(hrv.firstWhere((e) => e.rawEndTime == base + 420).uploaded, isFalse);
expect(hrv.firstWhere((e) => e.rawEndTime == base + 420).uploaded, isTrue);
expect(
realtime.firstWhere((e) => e.rawEndTime == base + 420).uploaded,
isFalse,
isTrue,
);
expect(
realtime
... ... @@ -169,6 +170,7 @@ void main() {
rawDataApi: api,
localStore: store,
userIdProvider: () => 42,
uploadResultsAfterCalculation: false,
);
final result = await service.syncAndStore(
... ... @@ -230,6 +232,42 @@ void main() {
expect(daily.single.uploaded, isFalse);
});
test('same calculated rows keep uploaded state', () async {
final store = _MemoryHealthRawStressLocalStore();
const uploadedDaily = HealthRawDailyStressPoint(
userId: 42,
date: 20260720,
stressValue: 66,
stressScore: 70,
state: HealthRawStressState.attention,
dataTime: 1000,
uploaded: true,
);
store.insertDailyStress(uploadedDaily);
await store.upsertDailyStressPoints(
userId: 42,
points: const [
HealthRawDailyStressPoint(
userId: 42,
date: 20260720,
stressValue: 66,
stressScore: 70,
state: HealthRawStressState.attention,
dataTime: 2000,
),
],
);
final daily = await store.queryDailyStressPoints(
userId: 42,
startDate: 20260720,
endDate: 20260720,
);
expect(daily.single.dataTime, 2000);
expect(daily.single.uploaded, isTrue);
});
test('startCoreCaculate skips raw reads when health auth is missing',
() async {
final api = _FakeHealthKitRawDataHostApi();
... ... @@ -238,6 +276,7 @@ void main() {
rawDataApi: api,
localStore: _MemoryHealthRawStressLocalStore(),
userIdProvider: () => 42,
uploadResultsAfterCalculation: false,
);
final result = await service.startCoreCaculate(readChunkDays: 1);
... ... @@ -262,6 +301,7 @@ void main() {
rawDataApi: api,
localStore: _MemoryHealthRawStressLocalStore(),
userIdProvider: () => 42,
uploadResultsAfterCalculation: false,
);
final eventFuture = service.healthDataUpdatedStream.first;
... ... @@ -328,6 +368,18 @@ class _FakeHealthKitRawDataHostApi extends HealthKitRawDataHostApi {
}
@override
Future<int> performHRDataUpload({required String sqliteFilePath}) async => 0;
@override
Future<int> performHRVDataUpload({required String sqliteFilePath}) async => 0;
@override
Future<bool> performAvgRealtimeStressDataUpload({
required String sqliteFilePath,
}) async =>
false;
@override
Future<List<HealthKitRawSleepDataPoint>> getHealthKitRawSleepData(
int startTime,
int endTime,
... ... @@ -397,6 +449,8 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
point.rawEndTime: point,
};
for (final point in result.hrvStressPoints) {
final existing = hrvByTime[point.rawEndTime];
final same = existing != null && _sameHrvStressPoint(existing, point);
hrvByTime[point.rawEndTime] = HealthRawHrvStressPoint(
userId: point.userId,
rawEndTime: point.rawEndTime,
... ... @@ -410,6 +464,7 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
baselineSleepHrv: point.baselineSleepHrv,
baselineRestingHr: point.baselineRestingHr,
flags: point.flags,
uploaded: same ? existing.uploaded : false,
);
}
_hrv[result.userId] = hrvByTime.values.toList()
... ... @@ -421,13 +476,18 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
point.rawEndTime: point,
};
for (final point in result.realtimeStressPoints) {
final existing = realtimeByTime[point.rawEndTime];
final same =
existing != null && _sameRealtimeStressPoint(existing, point);
realtimeByTime[point.rawEndTime] = HealthRawRealtimeStressPoint(
userId: point.userId,
rawEndTime: point.rawEndTime,
rawHr: point.rawHr,
result: point.result,
sourceStartTime: point.sourceStartTime,
sourceEndTime: point.sourceEndTime,
flags: point.flags,
uploaded: same ? existing.uploaded : false,
);
}
_realtime[result.userId] = realtimeByTime.values.toList()
... ... @@ -445,6 +505,12 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
point.date: point,
};
for (final point in points) {
final existing = byDate[point.date];
final same = existing != null &&
existing.userId == point.userId &&
existing.stressValue == point.stressValue &&
existing.stressScore == point.stressScore &&
existing.state == point.state;
byDate[point.date] = HealthRawDailyStressPoint(
userId: point.userId,
date: point.date,
... ... @@ -452,6 +518,7 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
stressScore: point.stressScore,
state: point.state,
dataTime: point.dataTime,
uploaded: same ? existing.uploaded : false,
);
}
_daily[userId] = byDate.values.toList()
... ... @@ -570,6 +637,17 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
}
@override
Future<void> markHrvStressUploadedUntil({
required int userId,
required int rawEndTime,
}) async {
final times = (_hrv[userId] ?? <HealthRawHrvStressPoint>[])
.where((point) => point.rawEndTime <= rawEndTime)
.map((point) => point.rawEndTime);
await markHrvStressUploaded(userId: userId, rawEndTimes: times);
}
@override
Future<void> markRealtimeStressUploaded({
required int userId,
required Iterable<int> rawEndTimes,
... ... @@ -581,6 +659,7 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
return HealthRawRealtimeStressPoint(
userId: point.userId,
rawEndTime: point.rawEndTime,
rawHr: point.rawHr,
result: point.result,
sourceStartTime: point.sourceStartTime,
sourceEndTime: point.sourceEndTime,
... ... @@ -590,6 +669,82 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
}).toList();
}
@override
Future<void> markRealtimeStressUploadedUntil({
required int userId,
required int rawEndTime,
}) async {
final times = (_realtime[userId] ?? <HealthRawRealtimeStressPoint>[])
.where((point) => point.rawEndTime <= rawEndTime)
.map((point) => point.rawEndTime);
await markRealtimeStressUploaded(userId: userId, rawEndTimes: times);
}
@override
Future<void> markDailyStressUploaded({required int userId}) async {
_daily[userId] = (_daily[userId] ?? <HealthRawDailyStressPoint>[])
.map((point) => HealthRawDailyStressPoint(
userId: point.userId,
date: point.date,
stressValue: point.stressValue,
stressScore: point.stressScore,
state: point.state,
dataTime: point.dataTime,
uploaded: true,
))
.toList();
}
@override
Future<bool> hasPendingHrvStressUploads({required int userId}) async {
return (_hrv[userId] ?? <HealthRawHrvStressPoint>[])
.any((point) => !point.uploaded);
}
@override
Future<bool> hasPendingRealtimeStressUploads({required int userId}) async {
return (_realtime[userId] ?? <HealthRawRealtimeStressPoint>[])
.any((point) => !point.uploaded);
}
@override
Future<bool> hasPendingDailyStressUploads({required int userId}) async {
return (_daily[userId] ?? <HealthRawDailyStressPoint>[])
.any((point) => !point.uploaded);
}
bool _sameHrvStressPoint(
HealthRawHrvStressPoint a,
HealthRawHrvStressPoint b,
) {
return a.userId == b.userId &&
a.rawHrv == b.rawHrv &&
a.result == b.result &&
a.state == b.state &&
a.baselineHrv == b.baselineHrv &&
a.baselineAwakeHrv == b.baselineAwakeHrv &&
a.baselineSleepHrv == b.baselineSleepHrv &&
a.baselineRestingHr == b.baselineRestingHr &&
_sameFlags(a.flags, b.flags);
}
bool _sameRealtimeStressPoint(
HealthRawRealtimeStressPoint a,
HealthRawRealtimeStressPoint b,
) {
return a.userId == b.userId &&
a.rawHr == b.rawHr &&
a.result == b.result &&
_sameFlags(a.flags, b.flags);
}
bool _sameFlags(HealthRawPointFlags a, HealthRawPointFlags b) {
return a.isWorkout == b.isWorkout &&
a.isWorkoutRecovery == b.isWorkoutRecovery &&
a.isSleepLikely == b.isSleepLikely &&
a.isSuspectedActivity == b.isSuspectedActivity;
}
Map<String, Object?> _hrvRow(HealthRawHrvStressPoint point) {
return <String, Object?>{
'raw_end_time': point.rawEndTime,
... ... @@ -615,6 +770,7 @@ class _MemoryHealthRawStressLocalStore extends HealthRawStressLocalStore {
return <String, Object?>{
'raw_end_time': point.rawEndTime,
'user_id': point.userId,
'raw_hr': point.rawHr,
'result': point.result,
'source_start_time': point.sourceStartTime,
'source_end_time': point.sourceEndTime,
... ...
... ... @@ -81,6 +81,146 @@ void main() {
expect(statistics.sleepTrendList?.single.sleepEvaluate, isNull);
expect(statistics.avgSleepDuration, isNull);
});
test('sleep duration excludes awake intervals inside sleep range', () {
final day = DateTime(2026, 7, 16);
final sleepStart = DateTime(2026, 7, 15, 23, 56);
final awakeStart = DateTime(2026, 7, 16, 3);
final awakeEnd = DateTime(2026, 7, 16, 3, 4);
final sleepEnd = DateTime(2026, 7, 16, 8, 24);
final intervals = [
HealthKitRawDataPoint(
dataType: 3,
startTime: LocalHealthDataConvert.unixSeconds(sleepStart),
endTime: LocalHealthDataConvert.unixSeconds(awakeStart),
),
HealthKitRawDataPoint(
dataType: 2,
startTime: LocalHealthDataConvert.unixSeconds(awakeStart),
endTime: LocalHealthDataConvert.unixSeconds(awakeEnd),
),
HealthKitRawDataPoint(
dataType: 3,
startTime: LocalHealthDataConvert.unixSeconds(awakeEnd),
endTime: LocalHealthDataConvert.unixSeconds(sleepEnd),
),
];
final statistics = LocalHealthDataConvert.sleepStatistics(
dateRangeType: 3,
days: [day],
previousDays: const [],
sleepIntervals: intervals,
heartRate: const [],
);
expect(
statistics.sleepTrendList?.single.totalTime,
const Duration(hours: 8, minutes: 24).inSeconds,
);
expect(
statistics.avgSleepDuration,
const Duration(hours: 8, minutes: 24).inSeconds,
);
});
test('sleep duration is rounded to displayed Apple Health minutes', () {
final day = DateTime(2026, 7, 16);
final sleepStart = DateTime(2026, 7, 15, 23, 56, 20);
final sleepEnd = DateTime(2026, 7, 16, 8, 20, 50);
final statistics = LocalHealthDataConvert.sleepStatistics(
dateRangeType: 3,
days: [day],
previousDays: const [],
sleepIntervals: [
HealthKitRawDataPoint(
dataType: 3,
startTime: LocalHealthDataConvert.unixSeconds(sleepStart),
endTime: LocalHealthDataConvert.unixSeconds(sleepEnd),
),
],
heartRate: const [],
);
expect(
statistics.sleepTrendList?.single.totalTime,
const Duration(hours: 8, minutes: 25).inSeconds,
);
expect(
statistics.avgSleepDuration,
const Duration(hours: 8, minutes: 25).inSeconds,
);
});
test('sleep duration only counts explicit asleep sleep stages', () {
final day = DateTime(2026, 7, 16);
final start = DateTime(2026, 7, 16, 1);
final statistics = LocalHealthDataConvert.sleepStatistics(
dateRangeType: 3,
days: [day],
previousDays: const [],
sleepIntervals: [
HealthKitRawDataPoint(
dataType: 0,
startTime: LocalHealthDataConvert.unixSeconds(start),
endTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 10)),
),
),
HealthKitRawDataPoint(
dataType: 1, // .asleep / .asleepUnspecified
startTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 10)),
),
endTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 20)),
),
),
HealthKitRawDataPoint(
dataType: 3,
startTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 20)),
),
endTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 30)),
),
),
HealthKitRawDataPoint(
dataType: 4,
startTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 30)),
),
endTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 40)),
),
),
HealthKitRawDataPoint(
dataType: 5,
startTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 40)),
),
endTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 50)),
),
),
HealthKitRawDataPoint(
dataType: 99,
startTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 50)),
),
endTime: LocalHealthDataConvert.unixSeconds(
start.add(const Duration(minutes: 60)),
),
),
],
heartRate: const [],
);
expect(
statistics.sleepTrendList?.single.totalTime,
const Duration(minutes: 40).inSeconds,
);
});
});
group('HealthRawDailyStressCalculator', () {
... ... @@ -109,6 +249,7 @@ void main() {
return HealthRawRealtimeStressPoint(
userId: userId,
rawEndTime: dayStart + offset,
rawHr: 70,
result: value,
sourceStartTime: dayStart + offset,
sourceEndTime: dayStart + offset,
... ... @@ -156,6 +297,7 @@ void main() {
HealthRawRealtimeStressPoint(
userId: userId,
rawEndTime: dayStart - 300,
rawHr: 80,
result: 80,
sourceStartTime: dayStart - 300,
sourceEndTime: dayStart - 300,
... ... @@ -164,6 +306,7 @@ void main() {
HealthRawRealtimeStressPoint(
userId: userId,
rawEndTime: dayStart + 100,
rawHr: 90,
result: 90,
sourceStartTime: dayStart + 100,
sourceEndTime: dayStart + 100,
... ... @@ -171,6 +314,7 @@ void main() {
HealthRawRealtimeStressPoint(
userId: userId,
rawEndTime: dayStart + 1000,
rawHr: 20,
result: 20,
sourceStartTime: dayStart + 1000,
sourceEndTime: dayStart + 1000,
... ...