ahri_tre_lake/
redcap_upload.rs1use crate::{LakeError, ScratchAttempt, ScratchAttemptId, UploadValidation};
3use sha2::{Digest, Sha256};
4use std::io::{Read, Write};
5
6pub struct RedcapUploadRole<'a> {
7 pub length: u64,
8 pub digest: &'a str,
9}
10pub struct RedcapUpload {
11 attempt: ScratchAttempt,
12}
13impl RedcapUpload {
14 pub fn receive(
15 scratch: &ScratchAttempt,
16 source: &mut dyn Read,
17 roles: [RedcapUploadRole<'_>; 3],
18 validation: &UploadValidation,
19 ) -> Result<Self, LakeError> {
20 let attempt = scratch.create_child(ScratchAttemptId::new(
21 &uuid::Uuid::new_v4().simple().to_string(),
22 )?)?;
23 let upload = Self { attempt };
24 let result = (|| {
25 let mut total = 0_u64;
26 for (index, role) in roles.into_iter().enumerate() {
27 total = total.checked_add(role.length).ok_or_else(invalid)?;
28 if role.length == 0
29 || total > validation.budgets.payload_bytes
30 || total > validation.budgets.decoded_file_bytes
31 || (index < 2 && role.length > 8 * 1024 * 1024)
32 {
33 return Err(invalid());
34 }
35 let path = upload.path(index);
36 let mut file = std::fs::OpenOptions::new()
37 .write(true)
38 .create_new(true)
39 .open(path)?;
40 let mut hash = Sha256::new();
41 let mut remaining = role.length;
42 let mut buffer = [0; 16 * 1024];
43 while remaining != 0 {
44 validation.check_deadline()?;
45 let size = buffer.len().min(remaining as usize);
46 source.read_exact(&mut buffer[..size])?;
47 hash.update(&buffer[..size]);
48 file.write_all(&buffer[..size])?;
49 remaining -= size as u64;
50 }
51 let digest = format!(
52 "sha256:{}",
53 hash.finalize()
54 .iter()
55 .map(|b| format!("{b:02x}"))
56 .collect::<String>()
57 );
58 if role.digest != digest {
59 return Err(invalid());
60 }
61 }
62 if source.read(&mut [0])? != 0 {
64 return Err(invalid());
65 }
66 validation.validate(&upload.path(2))?;
67 let connection = duckdb::Connection::open_in_memory()?;
68 connection
69 .execute_batch("SET memory_limit='128MB'; SET threads=1; SET temp_directory=''")?;
70 let path = upload.path(2).to_string_lossy().replace('\'', "''");
71 let mut statement = connection.prepare(&format!("SELECT record, field_name, value FROM read_csv('{path}', header=true, all_varchar=true) LIMIT 0"))?;
72 statement.query([])?.next()?;
73 validation.check_deadline()?;
74 Ok(())
75 })();
76 if let Err(error) = result {
77 upload.cleanup()?;
78 return Err(error);
79 }
80 Ok(upload)
81 }
82 pub fn path(&self, role: usize) -> std::path::PathBuf {
83 self.attempt
84 .path()
85 .join(["project.json", "metadata.json", "records.csv"][role])
86 }
87 pub fn cleanup(self) -> Result<(), LakeError> {
88 self.attempt.cleanup().map_err(Into::into)
89 }
90}
91fn invalid() -> LakeError {
92 std::io::Error::other("Invalid REDCap upload role").into()
93}
94
95pub struct RedcapExecutionLimits {
98 connection: duckdb::Connection,
99 memory: String,
100 threads: u64,
101}
102impl RedcapExecutionLimits {
103 pub fn new(connection: &duckdb::Connection) -> Result<Self, LakeError> {
104 let memory = connection.query_row("SELECT current_setting('memory_limit')", [], |row| {
105 row.get(0)
106 })?;
107 let threads =
108 connection.query_row("SELECT current_setting('threads')", [], |row| row.get(0))?;
109 let limits = Self {
110 connection: connection.try_clone()?,
111 memory,
112 threads,
113 };
114 limits
115 .connection
116 .execute_batch("SET memory_limit='128MB'; SET threads=1")?;
117 Ok(limits)
118 }
119}
120impl Drop for RedcapExecutionLimits {
121 fn drop(&mut self) {
122 let memory = self.memory.replace('\'', "''");
123 let _ = self.connection.execute_batch(&format!(
124 "SET memory_limit='{memory}'; SET threads={}",
125 self.threads
126 ));
127 }
128}