Skip to main content

cognite/api/data_ingestion/
raw.rs

1use std::collections::VecDeque;
2
3use futures::stream::{try_unfold, SelectAll};
4use futures::{FutureExt, TryStream, TryStreamExt};
5
6use crate::api::resource::Resource;
7use crate::dto::items::Items;
8use crate::error::Result;
9use crate::{raw::*, CondBoxedStream, CondSend, CursorState, CursorStreamState};
10use crate::{Cursor, ItemsVec, LimitCursorQuery};
11
12/// Raw is a NoSQL JSON store. Each project can have a variable number of databases,
13/// each of which will have a variable number of tables, each of which will have a variable
14/// number of key-value objects. Only queries on key are supported through this API.
15pub type RawResource = Resource<RawRow>;
16
17impl RawResource {
18    /// List Raw databases in the project.
19    ///
20    /// # Arguments
21    ///
22    /// * `limit` - Maximum number of databases to retrieve.
23    /// * `cursor` - Optional cursor for pagination.
24    pub async fn list_databases(
25        &self,
26        limit: Option<i32>,
27        cursor: Option<String>,
28    ) -> Result<ItemsVec<Database, Cursor>> {
29        let query = LimitCursorQuery { limit, cursor };
30        self.api_client
31            .get_with_params("raw/dbs", Some(query))
32            .await
33    }
34
35    /// Create a list of Raw databases.
36    ///
37    /// # Arguments
38    ///
39    /// * `dbs` - Databases to create.
40    pub async fn create_databases(&self, dbs: &[Database]) -> Result<Vec<Database>> {
41        let items = Items::new(dbs);
42        let result: ItemsVec<Database, Cursor> = self.api_client.post("raw/dbs", &items).await?;
43        Ok(result.items)
44    }
45
46    /// Delete a list of raw databases.
47    ///
48    /// # Arguments
49    ///
50    /// * `to_delete` - Request describing which databases to delete and how.
51    pub async fn delete_databases(&self, to_delete: &DeleteDatabasesRequest) -> Result<()> {
52        self.api_client
53            .post::<::serde_json::Value, DeleteDatabasesRequest>("raw/dbs/delete", to_delete)
54            .await?;
55        Ok(())
56    }
57
58    /// List tables in a a raw database.
59    ///
60    /// # Arguments
61    ///
62    /// * `db_name` - Database to list tables in.
63    /// * `limit` - Maximum number of tables to retrieve.
64    /// * `cursor` - Optional cursor for pagination.
65    pub async fn list_tables(
66        &self,
67        db_name: &str,
68        limit: Option<i32>,
69        cursor: Option<String>,
70    ) -> Result<ItemsVec<Table, Cursor>> {
71        let query = LimitCursorQuery { limit, cursor };
72        let path = format!("raw/dbs/{db_name}/tables");
73        self.api_client.get_with_params(&path, Some(query)).await
74    }
75
76    /// Create tables in a raw database.
77    ///
78    /// # Arguments
79    ///
80    /// * `db_name` - Database to create tables in.
81    /// * `ensure_parent` - If this is set to `true`, create database if it doesn't already exist.
82    /// * `tables` - Tables to create.
83    pub async fn create_tables(
84        &self,
85        db_name: &str,
86        ensure_parent: bool,
87        tables: &[Table],
88    ) -> Result<Vec<Table>> {
89        let query = EnsureParentQuery {
90            ensure_parent: Some(ensure_parent),
91        };
92        let path = format!("raw/dbs/{db_name}/tables");
93        let items = Items::new(tables);
94        let result: ItemsVec<Table, Cursor> = self
95            .api_client
96            .post_with_query(&path, &items, Some(query))
97            .await?;
98        Ok(result.items)
99    }
100
101    /// Delete tables in a raw database.
102    ///
103    /// # Arguments
104    ///
105    /// * `db_name` - Database to delete tables from.
106    /// * `to_delete` - Tables to delete.
107    pub async fn delete_tables(&self, db_name: &str, to_delete: &[Table]) -> Result<()> {
108        let path = format!("raw/dbs/{db_name}/tables/delete");
109        let items = Items::new(to_delete);
110        self.api_client
111            .post::<::serde_json::Value, _>(&path, &items)
112            .await?;
113        Ok(())
114    }
115
116    /// Retrieve cursors for parallel reads. This can be used to efficiently download
117    /// large volumes of data from a raw table in parallel.
118    ///
119    /// # Arguments
120    ///
121    /// * `db_name` - Database to retrieve from.
122    /// * `table_name` - Table to retrieve from.
123    /// * `params` - Optional filter parameters.
124    pub async fn retrieve_cursors_for_parallel_reads(
125        &self,
126        db_name: &str,
127        table_name: &str,
128        params: Option<RetrieveCursorsQuery>,
129    ) -> Result<Vec<String>> {
130        let path = format!("raw/dbs/{db_name}/tables/{table_name}/cursors");
131        let result: ItemsVec<String, Cursor> =
132            self.api_client.get_with_params(&path, params).await?;
133        Ok(result.items)
134    }
135
136    /// Retrieve rows from a table, with some basic filtering options.
137    ///
138    /// # Arguments
139    ///
140    /// * `db_name` - Database to retrieve rows from.
141    /// * `table_name` - Table to retrieve rows from.
142    /// * `params` - Optional filter parameters.
143    pub async fn retrieve_rows(
144        &self,
145        db_name: &str,
146        table_name: &str,
147        params: Option<RetrieveRowsQuery>,
148    ) -> Result<ItemsVec<RawRow, Cursor>> {
149        let path = format!("raw/dbs/{db_name}/tables/{table_name}/rows");
150        self.api_client.get_with_params(&path, params).await
151    }
152
153    /// Retrieve all rows from a table, following cursors. This returns a stream, you can abort the stream whenever you
154    /// want and only resources retrieved up to that point will be returned.
155    ///
156    /// Each item in the stream will be a result, after the first error is returned the
157    /// stream will end.
158    ///
159    /// `limit` in the filter only affects how many rows are returned _per request_.
160    ///
161    /// # Arguments
162    ///
163    /// * `db_name` - Database to retrieve rows from.
164    /// * `table_name` - Table to retrieve rows from.
165    /// * `params` - Optional filter parameters. This can set a cursor to start streaming from there.
166    pub fn retrieve_all_rows_stream<'a>(
167        &'a self,
168        db_name: &'a str,
169        table_name: &'a str,
170        params: Option<RetrieveRowsQuery>,
171    ) -> impl TryStream<Ok = RawRow, Error = crate::Error, Item = Result<RawRow>> + CondSend + 'a
172    {
173        let req = params.unwrap_or_default();
174        let initial_state = match &req.cursor {
175            Some(p) => CursorState::Some(p.to_owned()),
176            None => CursorState::Initial,
177        };
178        let state = CursorStreamState {
179            req,
180            responses: VecDeque::new(),
181            next_cursor: initial_state,
182        };
183
184        try_unfold(state, move |mut state| async move {
185            if let Some(next) = state.responses.pop_front() {
186                Ok(Some((next, state)))
187            } else {
188                let cursor = match std::mem::take(&mut state.next_cursor) {
189                    CursorState::Initial => None,
190                    CursorState::Some(x) => Some(x),
191                    CursorState::End => {
192                        return Ok(None);
193                    }
194                };
195                state.req.cursor = cursor;
196                let response = self
197                    .retrieve_rows(db_name, table_name, Some(state.req.clone()))
198                    .await?;
199
200                state.responses.extend(response.items);
201                state.next_cursor = match response.extra_fields.next_cursor {
202                    Some(x) => CursorState::Some(x),
203                    None => CursorState::End,
204                };
205                if let Some(next) = state.responses.pop_front() {
206                    Ok(Some((next, state)))
207                } else {
208                    Ok(None)
209                }
210            }
211        })
212    }
213
214    /// Retrieve all rows from a table, following cursors.
215    ///
216    /// `limit` in the filter only affects how many rows are returned _per request_.
217    ///
218    /// # Arguments
219    ///
220    /// * `db_name` - Database to retrieve rows from.
221    /// * `table_name` - Table to retrieve rows from.
222    /// * `params` - Optional filter parameters. This can set a cursor to start reading from there.
223    pub async fn retrieve_all_rows(
224        &self,
225        db_name: &str,
226        table_name: &str,
227        params: Option<RetrieveRowsQuery>,
228    ) -> Result<Vec<RawRow>> {
229        self.retrieve_all_rows_stream(db_name, table_name, params)
230            .try_collect()
231            .await
232    }
233
234    /// Retrieve all rows from a table, following cursors and reading from multiple streams in parallel.
235    ///
236    /// The order of the returned values is not guaranteed to be in any way consistent.
237    ///
238    /// * `db_name` - Database to retrieve rows from.
239    /// * `table_name` - Table to retrieve rows from.
240    /// * `params` - Optional filter parameters.
241    pub async fn retrieve_all_rows_partitioned(
242        &self,
243        db_name: &str,
244        table_name: &str,
245        params: RetrieveAllPartitionedQuery,
246    ) -> Result<Vec<RawRow>> {
247        self.retrieve_all_rows_partitioned_stream(db_name, table_name, params)
248            .try_collect()
249            .await
250    }
251
252    /// Retrieve all rows from a table, following cursors and reading from multiple streams in parallel.
253    ///
254    /// The order of the returned values is not guaranteed to be in any way consistent.
255    ///
256    /// * `db_name` - Database to retrieve rows from.
257    /// * `table_name` - Table to retrieve rows from.
258    /// * `params` - Optional filter parameters.
259    pub fn retrieve_all_rows_partitioned_stream<'a>(
260        &'a self,
261        db_name: &'a str,
262        table_name: &'a str,
263        params: RetrieveAllPartitionedQuery,
264    ) -> impl TryStream<Ok = RawRow, Error = crate::Error, Item = Result<RawRow>> + CondSend + 'a
265    {
266        self.retrieve_cursors_for_parallel_reads(
267            db_name,
268            table_name,
269            Some(RetrieveCursorsQuery {
270                min_last_updated_time: params.min_last_updated_time,
271                max_last_updated_time: params.max_last_updated_time,
272                number_of_cursors: params.number_of_cursors,
273            }),
274        )
275        .into_stream()
276        .map_ok(move |cursors| {
277            let mut streams = SelectAll::new();
278            for cursor in cursors {
279                let query = RetrieveRowsQuery {
280                    limit: params.limit,
281                    columns: params.columns.clone(),
282                    cursor: Some(cursor),
283                    min_last_updated_time: params.min_last_updated_time,
284                    max_last_updated_time: params.max_last_updated_time,
285                };
286                streams.push(
287                    self.retrieve_all_rows_stream(db_name, table_name, Some(query))
288                        .boxed_cond(),
289                );
290            }
291            streams
292        })
293        .try_flatten()
294    }
295
296    /// Insert rows into a table.
297    ///
298    /// If `ensure_parent` is true, create the database and/or table if they do not exist.
299    ///
300    /// # Arguments
301    ///
302    /// * `db_name` - Database to insert rows into.
303    /// * `table_name` - Table to insert rows into.
304    /// * `ensure_parent` - Create database and/or table if they do not exist.
305    /// * `rows` - Raw rows to create.
306    pub async fn insert_rows(
307        &self,
308        db_name: &str,
309        table_name: &str,
310        ensure_parent: bool,
311        rows: &[RawRowCreate],
312    ) -> Result<()> {
313        let path = format!("raw/dbs/{db_name}/tables/{table_name}/rows");
314        let query = EnsureParentQuery {
315            ensure_parent: Some(ensure_parent),
316        };
317        let items = Items::new(rows);
318        self.api_client
319            .post_with_query::<::serde_json::Value, _, EnsureParentQuery>(
320                &path,
321                &items,
322                Some(query),
323            )
324            .await?;
325        Ok(())
326    }
327
328    /// Retrieve a single row from a raw table.
329    ///
330    /// # Arguments
331    ///
332    /// * `db_name` - Database to retrieve from.
333    /// * `table_name` - Table to retrieve from.
334    /// * `key` - Key of row to retrieve.
335    pub async fn retrieve_row(&self, db_name: &str, table_name: &str, key: &str) -> Result<RawRow> {
336        let path = format!("raw/dbs/{db_name}/tables/{table_name}/rows/{key}");
337        self.api_client.get(&path).await
338    }
339
340    /// Delete rows from a raw table.
341    ///
342    /// # Arguments
343    ///
344    /// * `db_name` - Database to delete from.
345    /// * `table_name` - Table to delete from.
346    /// * `to_delete` - Rows to delete.
347    pub async fn delete_rows(
348        &self,
349        db_name: &str,
350        table_name: &str,
351        to_delete: &[DeleteRow],
352    ) -> Result<()> {
353        let path = format!("raw/dbs/{db_name}/tables/{table_name}/rows/delete");
354        let items = Items::new(to_delete);
355        self.api_client
356            .post::<::serde_json::Value, _>(&path, &items)
357            .await?;
358        Ok(())
359    }
360}