1use std::fs;
4use std::path::{Path, PathBuf};
5
6use serde::{Deserialize, Serialize};
7use sha2::{Digest, Sha256};
8use wellfare_core::companion_sync::CompanionHealthBundle;
9use wellfare_core::models::{HeartRateRecord, SleepRecord, StepRecord, WeightRecord};
10use wellfare_core::parser::{
11 parse_heart_rate_csv, parse_sleep_csv, parse_steps_csv, parse_weight_csv,
12};
13use wellfare_core::record::{EpistemicStatus, EvidenceType, RecordEnvelope, SensitivityClass};
14
15use super::api::WebizenHostApi;
16
17const QAPP_HEALTH: &str = "wellfair-health";
18
19#[derive(Debug, Clone)]
20pub struct EnvelopeWithSummary {
21 pub envelope: RecordEnvelope,
22 pub summary: Option<String>,
23}
24
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
26pub struct SamsungFileReport {
27 pub path: String,
28 pub kind: String,
29 pub records: u32,
30 pub rejected: u32,
31}
32
33#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
34pub struct SamsungImportReport {
35 pub source: String,
36 pub files: Vec<SamsungFileReport>,
37 pub records_committed: usize,
38 pub records_skipped: usize,
39 pub errors: Vec<String>,
40 #[serde(default, skip_serializing_if = "Option::is_none")]
41 pub checkpoint_hash: Option<String>,
42}
43
44fn content_hash_hex(payload: &str) -> String {
45 hex::encode(Sha256::digest(payload.as_bytes()).as_slice())
46}
47
48fn envelope_from_parts(
49 id: &str,
50 owner_did: &str,
51 author_did: &str,
52 asserted_unix: u32,
53 payload_json: &str,
54) -> RecordEnvelope {
55 RecordEnvelope {
56 id: id.to_string(),
57 owner_did: owner_did.to_string(),
58 author_did: author_did.to_string(),
59 proxy_did: None,
60 epistemic_status: EpistemicStatus::Asserted,
61 evidence_type: EvidenceType::DeviceMeasured,
62 sensitivity: SensitivityClass::Restricted,
63 asserted_time_unix: asserted_unix,
64 valid_time_start_unix: Some(asserted_unix),
65 valid_time_end_unix: None,
66 predecessor_id: None,
67 blob_hash: Some(content_hash_hex(payload_json)),
68 tombstone: false,
69 }
70}
71
72fn weight_envelopes(
73 records: &[WeightRecord],
74 owner: &str,
75 author: &str,
76) -> Vec<EnvelopeWithSummary> {
77 records
78 .iter()
79 .filter_map(|r| {
80 let payload = serde_json::to_string(r).ok()?;
81 let id = format!("urn:wellfair:weight:{}", r.uuid);
82 let unix = r.start_datetime.timestamp().max(0) as u32;
83 let summary = serde_json::json!({
84 "weight_kg": r.weight,
85 "bmi": r.bmi,
86 })
87 .to_string();
88 Some(EnvelopeWithSummary {
89 envelope: envelope_from_parts(&id, owner, author, unix, &payload),
90 summary: Some(summary),
91 })
92 })
93 .collect()
94}
95
96fn sleep_envelopes(records: &[SleepRecord], owner: &str, author: &str) -> Vec<EnvelopeWithSummary> {
97 records
98 .iter()
99 .filter_map(|r| {
100 let payload = serde_json::to_string(r).ok()?;
101 let id = format!("urn:wellfair:sleep:{}", r.uuid);
102 let unix = r.start_datetime.timestamp().max(0) as u32;
103 let summary = serde_json::json!({
104 "duration_min": r.sleep_duration,
105 "efficiency": r.efficiency,
106 "deep_min": r.deep_sleep,
107 "rem_min": r.rem_sleep,
108 "light_min": r.light_sleep,
109 })
110 .to_string();
111 Some(EnvelopeWithSummary {
112 envelope: envelope_from_parts(&id, owner, author, unix, &payload),
113 summary: Some(summary),
114 })
115 })
116 .collect()
117}
118
119fn heart_rate_envelopes(
120 records: &[HeartRateRecord],
121 owner: &str,
122 author: &str,
123) -> Vec<EnvelopeWithSummary> {
124 records
125 .iter()
126 .filter_map(|r| {
127 let payload = serde_json::to_string(r).ok()?;
128 let id = format!("urn:wellfair:heart_rate:{}", r.uuid);
129 let unix = r.start_datetime.timestamp().max(0) as u32;
130 let summary = serde_json::json!({
131 "heart_rate": r.heart_rate,
132 "min": r.min,
133 "max": r.max,
134 })
135 .to_string();
136 Some(EnvelopeWithSummary {
137 envelope: envelope_from_parts(&id, owner, author, unix, &payload),
138 summary: Some(summary),
139 })
140 })
141 .collect()
142}
143
144fn steps_envelopes(records: &[StepRecord], owner: &str, author: &str) -> Vec<EnvelopeWithSummary> {
145 records
146 .iter()
147 .filter_map(|r| {
148 let payload = serde_json::to_string(r).ok()?;
149 let id = format!("urn:wellfair:steps:{}", r.uuid);
150 let unix = r.start_datetime.timestamp().max(0) as u32;
151 let summary = serde_json::json!({
152 "steps": r.count,
153 "distance_m": r.distance,
154 })
155 .to_string();
156 Some(EnvelopeWithSummary {
157 envelope: envelope_from_parts(&id, owner, author, unix, &payload),
158 summary: Some(summary),
159 })
160 })
161 .collect()
162}
163
164#[derive(Debug, Clone, Copy, PartialEq, Eq)]
165pub enum SamsungCsvKind {
166 Weight,
167 Sleep,
168 HeartRate,
169 Steps,
170 Unknown,
171}
172
173fn classify_samsung_csv(name: &str) -> SamsungCsvKind {
174 let lower = name.to_ascii_lowercase();
175 if lower.contains("weight") || lower.contains("body_composition") {
176 SamsungCsvKind::Weight
177 } else if lower.contains("sleep") {
178 SamsungCsvKind::Sleep
179 } else if lower.contains("heart") {
180 SamsungCsvKind::HeartRate
181 } else if lower.contains("step") || lower.contains("walk") {
182 SamsungCsvKind::Steps
183 } else {
184 SamsungCsvKind::Unknown
185 }
186}
187
188fn kind_label(kind: SamsungCsvKind) -> &'static str {
189 match kind {
190 SamsungCsvKind::Weight => "weight",
191 SamsungCsvKind::Sleep => "sleep",
192 SamsungCsvKind::HeartRate => "heart_rate",
193 SamsungCsvKind::Steps => "steps",
194 SamsungCsvKind::Unknown => "unknown",
195 }
196}
197
198pub fn parse_csv_named_content(
200 filename: &str,
201 content: &str,
202 owner_did: &str,
203 author_did: &str,
204) -> Result<(SamsungCsvKind, Vec<EnvelopeWithSummary>, u32), String> {
205 let kind = classify_samsung_csv(filename);
206 match kind {
207 SamsungCsvKind::Weight => {
208 let records = parse_weight_csv(content).map_err(|e| e.to_string())?;
209 Ok((kind, weight_envelopes(&records, owner_did, author_did), 0))
210 }
211 SamsungCsvKind::Sleep => {
212 let records = parse_sleep_csv(content).map_err(|e| e.to_string())?;
213 Ok((kind, sleep_envelopes(&records, owner_did, author_did), 0))
214 }
215 SamsungCsvKind::HeartRate => {
216 let records = parse_heart_rate_csv(content).map_err(|e| e.to_string())?;
217 Ok((
218 kind,
219 heart_rate_envelopes(&records, owner_did, author_did),
220 0,
221 ))
222 }
223 SamsungCsvKind::Steps => {
224 let records = parse_steps_csv(content).map_err(|e| e.to_string())?;
225 Ok((kind, steps_envelopes(&records, owner_did, author_did), 0))
226 }
227 SamsungCsvKind::Unknown => Ok((kind, Vec::new(), 0)),
228 }
229}
230
231fn parse_csv_file(
232 path: &Path,
233 owner_did: &str,
234 author_did: &str,
235) -> Result<(SamsungCsvKind, Vec<EnvelopeWithSummary>, u32), String> {
236 let content = fs::read_to_string(path).map_err(|e| format!("read {}: {e}", path.display()))?;
237 let name = path.file_name().and_then(|n| n.to_str()).unwrap_or("");
238 parse_csv_named_content(name, &content, owner_did, author_did)
239}
240
241fn commit_envelopes(
242 host: &mut WebizenHostApi,
243 report: &mut SamsungImportReport,
244 path_label: &str,
245 kind: SamsungCsvKind,
246 envelopes: Vec<EnvelopeWithSummary>,
247 rejected: u32,
248) {
249 let mut committed = 0u32;
250 for item in envelopes {
251 match host.submit_record_with_summary(
252 QAPP_HEALTH,
253 item.envelope,
254 &report.source,
255 item.summary,
256 ) {
257 Ok(n) => {
258 report.records_committed += n;
259 committed += 1;
260 }
261 Err(e) => {
262 report.records_skipped += 1;
263 report.errors.push(e);
264 }
265 }
266 }
267 report.files.push(SamsungFileReport {
268 path: path_label.to_string(),
269 kind: kind_label(kind).to_string(),
270 records: committed,
271 rejected,
272 });
273}
274
275pub fn ingest_companion_health_bundle(
277 host: &mut WebizenHostApi,
278 bundle: &CompanionHealthBundle,
279 owner_did: &str,
280 author_did: &str,
281) -> SamsungImportReport {
282 let mut report = SamsungImportReport {
283 source: format!("companion:{}", bundle.device_id),
284 files: Vec::new(),
285 records_committed: 0,
286 records_skipped: 0,
287 errors: Vec::new(),
288 checkpoint_hash: None,
289 };
290
291 if let Err(e) = bundle.validate() {
292 report.errors.push(e);
293 return report;
294 }
295
296 for file in &bundle.files {
297 match parse_csv_named_content(&file.filename, &file.csv_content, owner_did, author_did) {
298 Ok((kind, envelopes, rejected)) => {
299 commit_envelopes(host, &mut report, &file.filename, kind, envelopes, rejected);
300 }
301 Err(e) => {
302 report.errors.push(e);
303 report.files.push(SamsungFileReport {
304 path: file.filename.clone(),
305 kind: "error".into(),
306 records: 0,
307 rejected: 1,
308 });
309 }
310 }
311 }
312
313 report
314}
315
316pub fn import_samsung_folder(
318 host: &mut WebizenHostApi,
319 folder: &Path,
320 owner_did: &str,
321 author_did: &str,
322) -> SamsungImportReport {
323 let mut report = SamsungImportReport {
324 source: format!("folder:{}", folder.display()),
325 files: Vec::new(),
326 records_committed: 0,
327 records_skipped: 0,
328 errors: Vec::new(),
329 checkpoint_hash: None,
330 };
331
332 if !folder.is_dir() {
333 report
334 .errors
335 .push(format!("Not a directory: {}", folder.display()));
336 return report;
337 }
338
339 let mut csv_paths: Vec<PathBuf> = Vec::new();
340 let entries = match fs::read_dir(folder) {
341 Ok(e) => e,
342 Err(e) => {
343 report.errors.push(format!("read_dir: {e}"));
344 return report;
345 }
346 };
347
348 for entry in entries.filter_map(Result::ok) {
349 let path = entry.path();
350 if path.extension().and_then(|e| e.to_str()) == Some("csv") {
351 csv_paths.push(path);
352 }
353 }
354 csv_paths.sort();
355
356 for path in csv_paths {
357 let label = path.display().to_string();
358 match parse_csv_file(&path, owner_did, author_did) {
359 Ok((kind, envelopes, rejected)) => {
360 commit_envelopes(host, &mut report, &label, kind, envelopes, rejected);
361 }
362 Err(e) => {
363 report.errors.push(e);
364 report.files.push(SamsungFileReport {
365 path: label,
366 kind: "error".into(),
367 records: 0,
368 rejected: 1,
369 });
370 }
371 }
372 }
373
374 report
375}
376
377#[cfg(test)]
378mod tests {
379 use super::*;
380 use std::io::Write;
381 use wellfare_core::companion_sync::{CompanionCsvFile, CompanionHealthBundle};
382
383 #[test]
384 fn classifies_samsung_csv_names() {
385 assert_eq!(
386 classify_samsung_csv("com.samsung.health.weight.20260101.csv"),
387 SamsungCsvKind::Weight
388 );
389 assert_eq!(classify_samsung_csv("sleep.csv"), SamsungCsvKind::Sleep);
390 }
391
392 #[test]
393 fn weight_csv_produces_envelopes() {
394 let dir = tempfile::tempdir().unwrap();
395 let csv = dir.path().join("weight.csv");
396 let mut f = fs::File::create(&csv).unwrap();
397 writeln!(
398 f,
399 "uuid,start_time,end_time,time_offset,weight,body_fat,muscle_mass,body_water,skeletal_muscle,bmi"
400 )
401 .unwrap();
402 writeln!(
403 f,
404 "a1000001-0000-4000-8000-000000000001,1777632000000,1777632060000,60,72.0,18.5,32.1,55.2,30.5,23.1"
405 )
406 .unwrap();
407
408 let (kind, envelopes, _) = parse_csv_file(&csv, "did:wf:owner", "did:wf:owner").unwrap();
409 assert_eq!(kind, SamsungCsvKind::Weight);
410 assert_eq!(envelopes.len(), 1);
411 assert_eq!(
412 envelopes[0].envelope.evidence_type,
413 EvidenceType::DeviceMeasured
414 );
415 }
416
417 #[test]
418 fn companion_bundle_parses_named_content() {
419 let csv = "uuid,start_time,end_time,time_offset,weight,body_fat,muscle_mass,body_water,skeletal_muscle,bmi\n\
420 a1000001-0000-4000-8000-000000000001,1777632000000,1777632060000,60,72.0,18.5,32.1,55.2,30.5,23.1\n";
421 let (kind, envelopes, _) =
422 parse_csv_named_content("weight.csv", csv, "did:wf:owner", "did:wf:owner").unwrap();
423 assert_eq!(kind, SamsungCsvKind::Weight);
424 assert_eq!(envelopes.len(), 1);
425
426 let bundle = CompanionHealthBundle::new(
427 "pixel-test",
428 1_700_000_000,
429 vec![CompanionCsvFile {
430 filename: "weight.csv".into(),
431 csv_content: csv.into(),
432 }],
433 );
434 assert!(bundle.validate().is_ok());
435 }
436}