cognite-sdk 0.6.2

SDK for the Cognite Data Fusion API
Documentation
#![cfg(feature = "integration_tests")]

use bytes::Bytes;
use cognite::models::instances::{
    CogniteExtractorFile, ExtractorFileObject, InstanceId, NodeOrEdgeSpecification, SlimNodeOrEdge,
};
use cognite::models::ItemId;
use cognite::prelude::*;
use cognite::utils::lease::ResourceLease;
use cognite::{files::*, Identity};
use futures::TryStreamExt;
#[cfg(test)]
use std::time::Duration;
use tokio::fs::File;
use tokio_util::codec::{BytesCodec, FramedRead};

mod common;
pub use common::*;
use uuid::Uuid;

async fn ensure_test_file(client: &CogniteClient) {
    let id = "rust-sdk-test-file".to_string();
    let new_file = AddFile {
        name: "Rust SDK test file".to_string(),
        external_id: Some(id),
        mime_type: Some("text/plain".to_string()),
        ..Default::default()
    };

    let file = match client.files.upload(false, &new_file).await {
        Err(cognite::Error::Conflict(_)) => return,
        Err(e) => panic!("{}", e),
        Ok(f) => f,
    };

    let chunks: Vec<Result<_, ::std::io::Error>> = vec![Ok("test "), Ok("file "), Ok("contents")];
    let stream = futures::stream::iter(chunks);

    client
        .files
        .upload_stream("text/plain", &file.extra.upload_url, stream, false)
        .await
        .unwrap();
}

#[tokio::test]
async fn create_upload_delete_file() {
    let id = format!("{}-file1", PREFIX.as_str());
    let new_file = AddFile {
        name: "File 1".to_string(),
        external_id: Some(id),
        mime_type: Some("text/plain".to_string()),
        ..Default::default()
    };

    let client = get_client();

    let res = client.files.upload(true, &new_file).await.unwrap();
    let mut lease = ResourceLease::new_println(client.files.clone());
    lease.add_resources([res.metadata.clone()]);

    let chunks: Vec<Result<_, ::std::io::Error>> = vec![Ok("test "), Ok("file "), Ok("contents")];

    let stream = futures::stream::iter(chunks);

    client
        .files
        .upload_stream("text/plain", &res.extra.upload_url, stream, false)
        .await
        .unwrap();

    lease.clean().await.unwrap();
}

#[tokio::test]
async fn create_upload_delete_actual_file() {
    let id = format!("{}-file3", PREFIX.as_str());
    let new_file = AddFile {
        name: "File 1".to_string(),
        external_id: Some(id),
        mime_type: Some("text/plain".to_string()),
        ..Default::default()
    };

    let client = get_client();

    let res = client.files.upload(true, &new_file).await.unwrap();
    let mut lease = ResourceLease::new_println(client.files.clone());
    lease.add_resources([res.metadata.clone()]);

    let size = tokio::fs::metadata("tests/dummyfile.txt")
        .await
        .unwrap()
        .len();
    let file = File::open("tests/dummyfile.txt").await.unwrap();
    let stream = FramedRead::new(file, BytesCodec::new());
    client
        .files
        .upload_stream_known_size("text/plain", &res.extra.upload_url, stream, size)
        .await
        .unwrap();

    lease.clean().await.unwrap();
}

#[tokio::test]
async fn create_update_delete_file() {
    let id = format!("{}-file2", PREFIX.as_str());
    let new_file = AddFile {
        name: "File 2".to_string(),
        external_id: Some(id.clone()),
        mime_type: Some("text/plain".to_string()),
        ..Default::default()
    };

    let client = get_client();

    let mut res = client.files.upload(true, &new_file).await.unwrap();
    let mut lease = ResourceLease::new_println(client.files.clone());
    lease.add_resources([res.metadata.clone()]);

    res.metadata.source = Some("New source".to_string());

    let upd_res = client.files.update_from(&vec![res.metadata]).await.unwrap();

    let upd_res = upd_res.first().unwrap();

    assert_eq!(Some("New source".to_string()), upd_res.source);

    lease.clean().await.unwrap();
}

#[tokio::test]
async fn download_test_file() {
    let client = get_client();

    ensure_test_file(&client).await;

    let data: Vec<Bytes> = client
        .files
        .download_file(IdentityOrInstance::Identity(Identity::ExternalId {
            external_id: "rust-sdk-test-file".to_string(),
        }))
        .await
        .unwrap()
        .try_collect()
        .await
        .unwrap();

    let data: Vec<u8> = data.into_iter().flatten().collect();
    let contents = String::from_utf8(data).unwrap();

    assert_eq!("test file contents", contents.as_str());
}

#[tokio::test]
async fn create_multipart_file() {
    let id = format!("{}-multipart", PREFIX.as_str());
    let new_file = AddFile {
        name: "Multipart File".to_owned(),
        external_id: Some(id.clone()),
        mime_type: Some("text/plain".to_owned()),
        ..Default::default()
    };

    let client = get_client();

    let content_1 = "abcde".repeat(1_200_000);
    let content_2 = "fghij";
    let (session, filemeta) = client
        .files
        .multipart_upload(true, 2, &new_file)
        .await
        .unwrap();
    let mut lease = ResourceLease::new_println(client.files.clone());
    lease.add_resources([filemeta]);
    session.upload_part_blob(0, content_1).await.unwrap();
    session.upload_part_blob(1, content_2).await.unwrap();

    session.complete().await.unwrap();

    let data: Vec<Bytes> = client
        .files
        .download_file(IdentityOrInstance::Identity(Identity::ExternalId {
            external_id: id.to_owned(),
        }))
        .await
        .unwrap()
        .try_collect()
        .await
        .unwrap();
    let data: Vec<u8> = data.into_iter().flatten().collect();
    assert_eq!(1_200_000 * 5 + 5, data.len());
    assert!(data.ends_with("fghij".as_bytes()));

    lease.clean().await.unwrap();
}

#[tokio::test]
async fn create_delete_dm_files() {
    let _permit = CDM_CONCURRENCY_PERMITS.acquire().await.unwrap();
    let client = CogniteClient::new_oidc("testing_instances", None).unwrap();
    let external_id = Uuid::new_v4().to_string();
    let space = std::env::var("CORE_DM_TEST_SPACE").unwrap();
    let col = CogniteExtractorFile::new(
        space.to_string(),
        external_id.to_string(),
        ExtractorFileObject::new(),
    );
    let res = client
        .models
        .instances
        .apply(&[col], None, None, None, None, false)
        .await
        .unwrap();
    let res = res.first().unwrap();
    assert!(matches!(res, SlimNodeOrEdge::Node(_)));

    let res_node = match res {
        cognite::models::instances::SlimNodeOrEdge::Node(slim_node_definition) => {
            slim_node_definition
        }
        cognite::models::instances::SlimNodeOrEdge::Edge(_) => {
            panic!("Invalid type received.")
        }
    };
    assert_eq!(external_id.to_string(), res_node.external_id);
    let id = IdentityOrInstance::InstanceId {
        instance_id: InstanceId {
            space: space.to_string(),
            external_id: external_id.to_string(),
        },
    };

    let res = client.files.get_upload_link(&id).await.unwrap();
    let size = tokio::fs::metadata("tests/dummyfile.txt")
        .await
        .unwrap()
        .len();
    let file = File::open("tests/dummyfile.txt").await.unwrap();
    let stream = FramedRead::new(file, BytesCodec::new());
    client
        .files
        .upload_stream_known_size("text/plain", &res.extra.upload_url, stream, size)
        .await
        .unwrap();

    let node_specs = NodeOrEdgeSpecification::Node(ItemId {
        space: space.to_string(),
        external_id: external_id.to_string(),
    });

    let mut backoff = Backoff::default();
    let mut deleted: Option<ItemsVec<NodeOrEdgeSpecification>> = None;
    for _ in 0..10 {
        match client.models.instances.delete(&[node_specs.clone()]).await {
            Ok(res) => {
                deleted = Some(res);
                break;
            }
            Err(_) => {
                tokio::time::sleep(backoff.next().unwrap()).await;
                continue;
            }
        }
    }
    let deleted = deleted.unwrap();
    let deleted = deleted.items.first().unwrap();
    assert!(matches!(deleted, NodeOrEdgeSpecification::Node(_)));
}

#[tokio::test]
async fn create_core_dm_multipart_file() {
    let _permit = CDM_CONCURRENCY_PERMITS.acquire().await.unwrap();
    let client = CogniteClient::new_oidc("testing_instances", None).unwrap();
    let external_id = Uuid::new_v4().to_string();
    let space = std::env::var("CORE_DM_TEST_SPACE").unwrap();
    let col = CogniteExtractorFile::new(
        space.to_string(),
        external_id.to_string(),
        ExtractorFileObject::new(),
    );
    let res = client
        .models
        .instances
        .apply(&[col], None, None, None, None, false)
        .await
        .unwrap();
    let res = res.first().unwrap();
    assert!(matches!(res, SlimNodeOrEdge::Node(_)));

    let res_node = match res {
        cognite::models::instances::SlimNodeOrEdge::Node(slim_node_definition) => {
            slim_node_definition
        }
        cognite::models::instances::SlimNodeOrEdge::Edge(_) => {
            panic!("Invalid type received.")
        }
    };
    assert_eq!(external_id.to_string(), res_node.external_id);
    let id = IdentityOrInstance::InstanceId {
        instance_id: InstanceId {
            space: space.to_string(),
            external_id: external_id.to_string(),
        },
    };

    let content_1 = "abcde".repeat(1_200_000);
    let content_2 = "fghij";
    let (session, _) = client
        .files
        .multipart_upload_existing(&id, 2)
        .await
        .unwrap();
    session.upload_part_blob(0, content_1).await.unwrap();
    session.upload_part_blob(1, content_2).await.unwrap();

    session.complete().await.unwrap();
    tokio::time::sleep(Duration::from_secs(3)).await;
    let id_json = serde_json::to_string(&id).unwrap();
    println!("{id_json}");
    let mut data: Option<_> = None;
    let mut backoff = Backoff::default();
    for _ in 0..=3 {
        match client.files.download_file(id.clone()).await {
            Ok(d) => data = Some(d),
            Err(_) => {
                tokio::time::sleep(backoff.next().unwrap()).await;
                continue;
            }
        }
    }

    let data: Vec<Bytes> = data.unwrap().try_collect().await.unwrap();
    let data: Vec<u8> = data.into_iter().flatten().collect();
    assert_eq!(1_200_000 * 5 + 5, data.len());
    assert!(data.ends_with("fghij".as_bytes()));

    let node_specs = NodeOrEdgeSpecification::Node(ItemId {
        space: space.to_string(),
        external_id: external_id.to_string(),
    });

    let mut backoff = Backoff::default();
    let mut deleted: Option<ItemsVec<NodeOrEdgeSpecification>> = None;
    for _ in 0..10 {
        match client.models.instances.delete(&[node_specs.clone()]).await {
            Ok(res) => {
                deleted = Some(res);
                break;
            }
            Err(_) => {
                tokio::time::sleep(backoff.next().unwrap()).await;
                continue;
            }
        }
    }
    let deleted = deleted.unwrap();
    let deleted = deleted.items.first().unwrap();
    assert!(matches!(deleted, NodeOrEdgeSpecification::Node(_)));
}