Skip to main content

dkdc_lake/
files.rs

1use crate::Lake;
2use anyhow::Result;
3use chrono::{DateTime, Utc};
4use dkdc_config::FILES_TABLE_NAME;
5use duckdb::params;
6
7pub struct File {
8    pub filepath: String,
9    pub filename: String,
10    pub filedata: Vec<u8>,
11    pub filesize: i64,
12    pub fileupdated: DateTime<Utc>,
13}
14
15impl Lake {
16    pub fn create_files_table(&self) -> Result<()> {
17        let sql = format!(
18            "CREATE TABLE IF NOT EXISTS {} (
19                filepath VARCHAR,
20                filename VARCHAR,
21                filedata BLOB,
22                filesize BIGINT,
23                fileupdated TIMESTAMP
24            )",
25            FILES_TABLE_NAME
26        );
27        self.execute(&sql)?;
28        Ok(())
29    }
30
31    pub fn add_file(&self, filepath: &str, filename: &str, data: &[u8]) -> Result<()> {
32        self.create_files_table()?;
33
34        let sql = format!(
35            "INSERT INTO {} (filepath, filename, filedata, filesize, fileupdated)
36             VALUES (?, ?, ?, ?, ?)",
37            FILES_TABLE_NAME
38        );
39
40        let mut stmt = self.prepare(&sql)?;
41        stmt.execute(params![
42            filepath,
43            filename,
44            data,
45            data.len() as i64,
46            Utc::now().to_rfc3339(),
47        ])?;
48
49        Ok(())
50    }
51
52    pub fn get_file(&self, filepath: &str, filename: &str) -> Result<Option<File>> {
53        let sql = format!(
54            "SELECT filepath, filename, filedata, filesize, fileupdated
55             FROM {}
56             WHERE filepath = ? AND filename = ?
57             ORDER BY fileupdated DESC
58             LIMIT 1",
59            FILES_TABLE_NAME
60        );
61
62        let mut stmt = self.prepare(&sql)?;
63        let mut rows = stmt.query(params![filepath, filename])?;
64
65        if let Some(row) = rows.next()? {
66            Ok(Some(File {
67                filepath: row.get(0)?,
68                filename: row.get(1)?,
69                filedata: row.get(2)?,
70                filesize: row.get(3)?,
71                fileupdated: {
72                    // DuckDB returns timestamps as microseconds since epoch
73                    let micros: i64 = row.get(4)?;
74                    let secs = micros / 1_000_000;
75                    let nanos = ((micros % 1_000_000) * 1000) as u32;
76                    DateTime::from_timestamp(secs, nanos).unwrap_or_else(Utc::now)
77                },
78            }))
79        } else {
80            Ok(None)
81        }
82    }
83
84    pub fn list_files(&self, filepath: &str) -> Result<Vec<String>> {
85        let sql = format!(
86            "SELECT DISTINCT filename
87             FROM {}
88             WHERE filepath = ?
89             ORDER BY filename",
90            FILES_TABLE_NAME
91        );
92
93        let mut stmt = self.prepare(&sql)?;
94        let mut rows = stmt.query(params![filepath])?;
95
96        let mut files = Vec::new();
97        while let Some(row) = rows.next()? {
98            files.push(row.get(0)?);
99        }
100
101        Ok(files)
102    }
103
104    pub fn delete_file(&self, filepath: &str, filename: &str) -> Result<()> {
105        let sql = format!(
106            "DELETE FROM {} WHERE filepath = ? AND filename = ?",
107            FILES_TABLE_NAME
108        );
109
110        let mut stmt = self.prepare(&sql)?;
111        stmt.execute(params![filepath, filename])?;
112
113        Ok(())
114    }
115}