Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
177 changes: 30 additions & 147 deletions Cargo.lock

Large diffs are not rendered by default.

1 change: 0 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@ documentation = "https://docs.rs/nvisy-server"

# Runtime crates
nvisy-engine = { git = "https://github.com/nvisycom/runtime", branch = "main" }
nvisy-schema = { git = "https://github.com/nvisycom/runtime", branch = "main", default-features = false }

# Internal crates
nvisy-core = { path = "./crates/nvisy-core", version = "0.1.0" }
Expand Down
4 changes: 1 addition & 3 deletions crates/nvisy-postgres/src/model/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@ mod workspace;
mod workspace_activity;
mod workspace_connection;
mod workspace_connection_run;
mod workspace_context;
mod workspace_file;
mod workspace_invite;
mod workspace_member;
Expand All @@ -27,7 +26,7 @@ pub use account_api_token::{AccountApiToken, NewAccountApiToken, UpdateAccountAp
pub use account_notification::{
AccountNotification, NewAccountNotification, UpdateAccountNotification,
};
pub use pipeline_reference::{PipelineContext, PipelinePolicy};
pub use pipeline_reference::PipelinePolicy;
// Workspace models
pub use workspace::{NewWorkspace, UpdateWorkspace, Workspace};
pub use workspace_activity::{NewWorkspaceActivity, WorkspaceActivity};
Expand All @@ -37,7 +36,6 @@ pub use workspace_connection::{
pub use workspace_connection_run::{
NewWorkspaceConnectionRun, UpdateWorkspaceConnectionRun, WorkspaceConnectionRun,
};
pub use workspace_context::{NewWorkspaceContext, UpdateWorkspaceContext, WorkspaceContext};
pub use workspace_file::{NewWorkspaceFile, UpdateWorkspaceFile, WorkspaceFile};
pub use workspace_invite::{NewWorkspaceInvite, UpdateWorkspaceInvite, WorkspaceInvite};
pub use workspace_member::{NewWorkspaceMember, UpdateWorkspaceMember, WorkspaceMember};
Expand Down
19 changes: 3 additions & 16 deletions crates/nvisy-postgres/src/model/pipeline_reference.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
//! Join models linking a pipeline to the policies and contexts it references.
//! Join model linking a pipeline to the policies it references.
//!
//! References are relational (real foreign keys) rather than embedded in the
//! pipeline's JSON definition, so the database enforces integrity and cleans up
//! on cascade. The composite keys pin every reference to a single workspace.
//! on cascade. The composite key pins every reference to a single workspace.

use diesel::prelude::*;
use uuid::Uuid;

use crate::schema::{workspace_pipeline_contexts, workspace_pipeline_policies};
use crate::schema::workspace_pipeline_policies;

/// A pipeline → policy reference row.
#[derive(Debug, Clone, PartialEq, Queryable, Selectable, Insertable)]
Expand All @@ -21,16 +21,3 @@ pub struct PipelinePolicy {
/// Referenced policy.
pub policy_id: Uuid,
}

/// A pipeline → context reference row.
#[derive(Debug, Clone, PartialEq, Queryable, Selectable, Insertable)]
#[diesel(table_name = workspace_pipeline_contexts)]
#[diesel(check_for_backend(diesel::pg::Pg))]
pub struct PipelineContext {
/// Workspace both the pipeline and context belong to.
pub workspace_id: Uuid,
/// Referencing pipeline.
pub pipeline_id: Uuid,
/// Referenced context.
pub context_id: Uuid,
}
110 changes: 0 additions & 110 deletions crates/nvisy-postgres/src/model/workspace_context.rs

This file was deleted.

2 changes: 0 additions & 2 deletions crates/nvisy-postgres/src/query/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@ mod workspace;
mod workspace_activity;
mod workspace_connection;
mod workspace_connection_run;
mod workspace_context;
mod workspace_file;
mod workspace_invite;
mod workspace_member;
Expand All @@ -39,7 +38,6 @@ pub use workspace::WorkspaceRepository;
pub use workspace_activity::WorkspaceActivityRepository;
pub use workspace_connection::WorkspaceConnectionRepository;
pub use workspace_connection_run::WorkspaceConnectionRunRepository;
pub use workspace_context::WorkspaceContextRepository;
pub use workspace_file::WorkspaceFileRepository;
pub use workspace_invite::WorkspaceInviteRepository;
pub use workspace_member::WorkspaceMemberRepository;
Expand Down
133 changes: 6 additions & 127 deletions crates/nvisy-postgres/src/query/pipeline_reference.rs
Original file line number Diff line number Diff line change
@@ -1,18 +1,17 @@
//! Repository for a pipeline's policy and context references.
//! Repository for a pipeline's policy references.
//!
//! References live in join tables (`workspace_pipeline_policies`, `workspace_pipeline_contexts`)
//! rather than the pipeline's JSON definition, so foreign keys enforce that
//! every referenced policy/context exists in the pipeline's workspace. The
//! `replace_*` operations are delete-then-insert and expect to run inside a
//! caller-owned transaction.
//! References live in the `workspace_pipeline_policies` join table rather than
//! the pipeline's JSON definition, so foreign keys enforce that every referenced
//! policy exists in the pipeline's workspace. The `replace_*` operation is
//! delete-then-insert and expects to run inside a caller-owned transaction.

use std::future::Future;

use diesel::prelude::*;
use diesel_async::RunQueryDsl;
use uuid::Uuid;

use crate::model::{PipelineContext, PipelinePolicy};
use crate::model::PipelinePolicy;
use crate::types::Slug;
use crate::{PgConnection, PgError, PgResult, schema};

Expand All @@ -29,14 +28,6 @@ pub trait PipelineReferenceRepository {
policy_ids: &[Uuid],
) -> impl Future<Output = PgResult<()>> + Send;

/// Replaces a pipeline's context references with the given set.
fn replace_workspace_pipeline_contexts(
&mut self,
workspace_id: Uuid,
pipeline_id: Uuid,
context_ids: &[Uuid],
) -> impl Future<Output = PgResult<()>> + Send;

/// Lists the ids of the policies a pipeline references.
///
/// Used by the run path to resolve each referenced policy to its record for
Expand All @@ -46,24 +37,12 @@ pub trait PipelineReferenceRepository {
pipeline_id: Uuid,
) -> impl Future<Output = PgResult<Vec<Uuid>>> + Send;

/// Lists the ids of the contexts a pipeline references.
fn list_pipeline_context_ids(
&mut self,
pipeline_id: Uuid,
) -> impl Future<Output = PgResult<Vec<Uuid>>> + Send;

/// Lists the slugs of the policies a pipeline references.
fn list_pipeline_policy_slugs(
&mut self,
pipeline_id: Uuid,
) -> impl Future<Output = PgResult<Vec<Slug>>> + Send;

/// Lists the slugs of the contexts a pipeline references.
fn list_pipeline_context_slugs(
&mut self,
pipeline_id: Uuid,
) -> impl Future<Output = PgResult<Vec<Slug>>> + Send;

/// Resolves policy slugs to their ids within a workspace, preserving order.
///
/// Returns `None` if any slug does not match a live policy in the workspace,
Expand All @@ -74,13 +53,6 @@ pub trait PipelineReferenceRepository {
workspace_id: Uuid,
slugs: &[Slug],
) -> impl Future<Output = PgResult<Option<Vec<Uuid>>>> + Send;

/// Resolves context slugs to their ids within a workspace, preserving order.
fn resolve_context_slugs(
&mut self,
workspace_id: Uuid,
slugs: &[Slug],
) -> impl Future<Output = PgResult<Option<Vec<Uuid>>>> + Send;
}

impl PipelineReferenceRepository for PgConnection {
Expand Down Expand Up @@ -117,39 +89,6 @@ impl PipelineReferenceRepository for PgConnection {
Ok(())
}

async fn replace_workspace_pipeline_contexts(
&mut self,
workspace_id: Uuid,
pipeline_id: Uuid,
context_ids: &[Uuid],
) -> PgResult<()> {
use schema::workspace_pipeline_contexts::{self, dsl};

diesel::delete(workspace_pipeline_contexts::table.filter(dsl::pipeline_id.eq(pipeline_id)))
.execute(self)
.await
.map_err(PgError::from)?;

if !context_ids.is_empty() {
let rows: Vec<PipelineContext> = dedup(context_ids)
.into_iter()
.map(|context_id| PipelineContext {
workspace_id,
pipeline_id,
context_id,
})
.collect();

diesel::insert_into(workspace_pipeline_contexts::table)
.values(&rows)
.execute(self)
.await
.map_err(PgError::from)?;
}

Ok(())
}

async fn list_pipeline_policy_ids(&mut self, pipeline_id: Uuid) -> PgResult<Vec<Uuid>> {
use schema::{workspace_pipeline_policies, workspace_policies};

Expand All @@ -168,24 +107,6 @@ impl PipelineReferenceRepository for PgConnection {
Ok(ids)
}

async fn list_pipeline_context_ids(&mut self, pipeline_id: Uuid) -> PgResult<Vec<Uuid>> {
use schema::{workspace_contexts, workspace_pipeline_contexts};

let ids = workspace_pipeline_contexts::table
.inner_join(
workspace_contexts::table
.on(workspace_contexts::id.eq(workspace_pipeline_contexts::context_id)),
)
.filter(workspace_pipeline_contexts::pipeline_id.eq(pipeline_id))
.filter(workspace_contexts::deleted_at.is_null())
.select(workspace_pipeline_contexts::context_id)
.load(self)
.await
.map_err(PgError::from)?;

Ok(ids)
}

async fn list_pipeline_policy_slugs(&mut self, pipeline_id: Uuid) -> PgResult<Vec<Slug>> {
use schema::{workspace_pipeline_policies, workspace_policies};

Expand All @@ -206,24 +127,6 @@ impl PipelineReferenceRepository for PgConnection {
Ok(slugs)
}

async fn list_pipeline_context_slugs(&mut self, pipeline_id: Uuid) -> PgResult<Vec<Slug>> {
use schema::{workspace_contexts, workspace_pipeline_contexts};

let slugs = workspace_pipeline_contexts::table
.inner_join(
workspace_contexts::table
.on(workspace_contexts::id.eq(workspace_pipeline_contexts::context_id)),
)
.filter(workspace_pipeline_contexts::pipeline_id.eq(pipeline_id))
.filter(workspace_contexts::deleted_at.is_null())
.select(workspace_contexts::slug)
.load(self)
.await
.map_err(PgError::from)?;

Ok(slugs)
}

async fn resolve_policy_slugs(
&mut self,
workspace_id: Uuid,
Expand All @@ -247,30 +150,6 @@ impl PipelineReferenceRepository for PgConnection {

Ok(map_slugs_to_ids(slugs, found))
}

async fn resolve_context_slugs(
&mut self,
workspace_id: Uuid,
slugs: &[Slug],
) -> PgResult<Option<Vec<Uuid>>> {
use schema::workspace_contexts::{self, dsl};

if slugs.is_empty() {
return Ok(Some(Vec::new()));
}

let wanted: Vec<String> = slugs.iter().map(|slug| slug.as_str().to_owned()).collect();
let found: Vec<(Slug, Uuid)> = workspace_contexts::table
.filter(dsl::workspace_id.eq(workspace_id))
.filter(dsl::deleted_at.is_null())
.filter(dsl::slug.eq_any(&wanted))
.select((dsl::slug, dsl::id))
.load(self)
.await
.map_err(PgError::from)?;

Ok(map_slugs_to_ids(slugs, found))
}
}

/// Maps the requested slugs to ids in request order, returning `None` if any
Expand Down
Loading