oxen-server 0.53.0

Oxen server is a fast data version control backend, supporting local disk and S3. Self host your repositories on your own storage, or use the hosted platform on Oxen.ai. Stores, syncs, and serves versioned datasets, model checkpoints, game assets, studio media, and any large data. Use the oxen CLI to push and pull from the oxen server.
use std::path::PathBuf;

use crate::errors::OxenHttpError;
use crate::helpers::get_repo;
use crate::params::{app_data, path_param};

use actix_web::{HttpRequest, HttpResponse, web};
use futures_util::stream::StreamExt as _;
use liboxen::core;
use liboxen::core::repo_locks;
use liboxen::repositories;
use liboxen::view::StatusMessage;
use liboxen::view::versions::CompleteVersionUploadRequest;
use serde::Deserialize;

#[derive(Deserialize, Debug)]
pub struct ChunkQuery {
    pub offset: Option<u64>,
    pub size: Option<u64>,
}

pub async fn upload(
    req: HttpRequest,
    query: web::Query<ChunkQuery>,
    mut body: web::Payload,
) -> Result<HttpResponse, OxenHttpError> {
    let app_data = app_data(&req)?;
    let namespace = path_param(&req, "namespace")?.to_string();
    let repo_name = path_param(&req, "repo_name")?.to_string();
    let version_id = path_param(&req, "version_id")?.to_string();

    let offset = query.offset.unwrap_or(0);

    let repo = get_repo(app_data, namespace, repo_name)?;
    let _write = repo_locks::acquire_write(&repo)?;

    log::debug!(
        "/upload version {} chunk offset{} to repo: {:?}",
        version_id,
        offset,
        repo.path
    );

    let version_store = repo.version_store();

    // Read the chunk from the payload. Chunks are <= the streamed-transfer segment size.
    let mut chunk = web::BytesMut::new();
    while let Some(part_result) = body.next().await {
        let part = part_result.map_err(|e| OxenHttpError::BadRequest(e.to_string().into()))?;
        chunk.extend_from_slice(&part);
    }

    version_store
        .store_version_chunk(&version_id, offset, chunk.freeze())
        .await?;

    Ok(HttpResponse::Ok().json(StatusMessage::resource_created()))
}

pub async fn complete(req: HttpRequest, body: String) -> Result<HttpResponse, OxenHttpError> {
    let app_data = app_data(&req)?;
    let namespace = path_param(&req, "namespace")?.to_string();
    let repo_name = path_param(&req, "repo_name")?.to_string();
    let version_id = path_param(&req, "version_id")?.to_string();
    let repo = get_repo(app_data, namespace, repo_name)?;
    let _write = repo_locks::acquire_write(&repo)?;

    log::debug!("/complete version chunk upload to repo: {:?}", repo.path);

    // Try to deserialize the body
    let request: Result<CompleteVersionUploadRequest, serde_json::Error> =
        serde_json::from_str(&body);
    if let Ok(request) = request {
        // There should only be a single file in the request
        if request.files.len() != 1 {
            return Ok(HttpResponse::BadRequest().json(StatusMessage::error(
                "Expected a single file in the request",
            )));
        }

        let file = &request.files[0];
        // Support both new clients (num_chunks) and old clients (upload_results)
        let num_chunks = file
            .num_chunks
            .or_else(|| file.upload_results.as_ref().map(|r| r.len()))
            .ok_or_else(|| {
                OxenHttpError::BadRequest(
                    "Missing both num_chunks and upload_results in request".into(),
                )
            })?;
        log::debug!("Client uploaded {num_chunks} chunks");
        let version_store = repo.version_store();

        let chunks = version_store.list_version_chunks(&version_id).await?;
        log::debug!("Found {} chunks on server", chunks.len());

        if chunks.len() != num_chunks {
            return Ok(
                HttpResponse::BadRequest().json(StatusMessage::error(format!(
                    "Number of chunks does not match expected number of chunks: {} != {}",
                    chunks.len(),
                    num_chunks
                ))),
            );
        }

        // Combine all the chunks for a version file into a single file
        version_store.combine_version_chunks(&version_id).await?;

        // If the workspace id is provided, stage the file
        if let Some(workspace_id) = request.workspace_id {
            let Some(workspace) = repositories::workspaces::get(&repo, &workspace_id)? else {
                return Ok(HttpResponse::NotFound().json(StatusMessage::error(format!(
                    "Workspace not found: {workspace_id}"
                ))));
            };
            let dst_path = if let Some(dst_dir) = &file.dst_dir {
                dst_dir.join(file.file_name.clone())
            } else {
                PathBuf::from(file.file_name.clone())
            };

            core::v_latest::workspaces::files::add_version_file(&workspace, &dst_path, &version_id)
                .await?;
        }

        return Ok(HttpResponse::Ok().json(StatusMessage::resource_found()));
    }

    Ok(HttpResponse::BadRequest().json(StatusMessage::error("Invalid request body")))
}

pub async fn create(_req: HttpRequest, _body: String) -> Result<HttpResponse, OxenHttpError> {
    Ok(HttpResponse::Ok().json(StatusMessage::resource_found()))
}