Skip to main content

ahri_tre_lake/
redcap_upload.rs

1//! Three verified REDCap role files retained only in configured Trusted scratch.
2use 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            // Force the enclosing framed stream to validate its terminal digest.
63            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
95/// Session-exclusive form work uses a finite engine budget; the prior settings
96/// are restored before the Session can serve another workflow.
97pub 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}