Skip to content
1 change: 1 addition & 0 deletions crates/context_aware_config/src/api/config.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
mod handlers;
pub use handlers::endpoints;
pub use handlers::execute_reduce;
pub mod helpers;
78 changes: 53 additions & 25 deletions crates/context_aware_config/src/api/config/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,17 @@ use actix_web::{
web::{Data, Header, Json, Path, Query},
};
use chrono::{DateTime, Utc};
use diesel::{ExpressionMethods, QueryDsl, RunQueryDsl, SelectableHelper};
use diesel::{
ExpressionMethods, PgConnection, QueryDsl, RunQueryDsl, SelectableHelper,
r2d2::{ConnectionManager, PooledConnection},
};
use itertools::Itertools;
use serde_json::{Map, Value, json};
use service_utils::{
helpers::{fetch_dimensions_info_map, is_not_modified},
kronos_dispatch::submit_job,
redis::{CONFIG_KEY_SUFFIX, LAST_MODIFIED_KEY_SUFFIX, read_through_cache},
service::types::{
AppHeader, AppState, DbConnection, WorkspaceContext, WorkspaceWritePermit,
},
service::types::{AppHeader, AppState, DbConnection, WorkspaceContext},
};
use superposition_core::{
ConfigFormat, JsonFormat, TomlFormat,
Expand All @@ -30,14 +32,15 @@ use superposition_types::{
MergeStrategy, ResolveConfigQuery,
},
context::PutRequest,
jobs::{JobCreateResponse, JobRequest, ReduceRequest},
},
custom_query::{
self as superposition_query, CustomQuery, DimensionQuery, PaginationParams,
QueryMap,
},
database::{
models::{
ChangeReason,
ChangeReason, JobWorkspace,
cac::{ConfigVersion, ConfigVersionListItem},
},
schema::config_versions::dsl as config_versions,
Expand Down Expand Up @@ -444,22 +447,13 @@ async fn reduce_config_key(
})
}

#[authorized]
#[put("/reduce")]
async fn reduce_handler(
workspace_context: WorkspaceContext,
req: HttpRequest,
user: User,
mut write_permit: WorkspaceWritePermit,
state: Data<AppState>,
) -> superposition::Result<HttpResponse> {
let conn = write_permit.connection();
let is_approve = req
.headers()
.get("x-approve")
.and_then(|value| value.to_str().ok().and_then(|s| s.parse::<bool>().ok()))
.unwrap_or(false);

pub async fn execute_reduce(
workspace_context: &WorkspaceContext,
state: &Data<AppState>,
conn: &mut PooledConnection<ConnectionManager<PgConnection>>,
user: &User,
is_approve: bool,
) -> superposition::Result<()> {
let dimensions_info_map =
fetch_dimensions_info_map(conn, &workspace_context.schema_name)?;
let mut config = generate_cac(conn, &workspace_context.schema_name)?;
Expand All @@ -469,24 +463,58 @@ async fn reduce_handler(
let overrides = config.overrides;
let default_config = config.default_configs.into_inner();
config = reduce_config_key(
&user,
user,
conn,
contexts.clone(),
overrides.clone(),
key.as_str(),
&dimensions_info_map,
default_config.clone(),
is_approve,
&workspace_context,
&state,
workspace_context,
state,
)
.await?;
if is_approve {
config = generate_cac(conn, &workspace_context.schema_name)?;
}
}
Ok(())
}

#[authorized]
#[put("/reduce")]
async fn reduce_handler(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we consider locking while doing reduce ?

workspace_context: WorkspaceContext,
db_conn: DbConnection,
state: Data<AppState>,
req: Json<ReduceRequest>,
) -> superposition::Result<Json<JobCreateResponse>> {
let DbConnection(mut conn) = db_conn;
let req = req.into_inner();
let job_request = JobRequest::Reduce(req);
let job_workspace = JobWorkspace::from(&workspace_context.schema_name.0);
let target_workspace = state
.kronos_workspace
.as_deref()
.unwrap_or(&workspace_context.schema_name);

let response = submit_job(
state.kronos_client.as_ref(),
target_workspace,
&job_workspace,
&workspace_context.organisation_id,
&workspace_context.workspace_id,
job_request,
&state.snowflake_generator,
&mut conn,
3,
"Reduce config job",
)
.await
.map_err(|e| unexpected_error!("Failed to submit reduce job: {}", e))?;

Ok(HttpResponse::Ok().json(config))
Ok(Json(response))
}

#[authorized]
Expand Down
1 change: 1 addition & 0 deletions crates/context_aware_config/src/api/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ pub mod operations;
mod types;
pub mod validations;
pub use handlers::endpoints;
pub use handlers::execute_priority_recompute;
pub use operations::delete;
pub use operations::update;
pub use operations::upsert;
Expand Down
98 changes: 56 additions & 42 deletions crates/context_aware_config/src/api/context/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,18 @@ use actix_web::{
use bigdecimal::BigDecimal;
use chrono::Utc;
use diesel::{
Connection, ExpressionMethods, OptionalExtension, QueryDsl, RunQueryDsl,
SelectableHelper,
Connection, ExpressionMethods, OptionalExtension, PgConnection, QueryDsl,
RunQueryDsl, SelectableHelper,
dsl::sql,
r2d2::{ConnectionManager, PooledConnection},
sql_types::{Bool, Text},
};
use serde_json::{Map, Value};
use service_utils::{
helpers::{
WebhookData, execute_webhook_call, fetch_dimensions_info_map, parse_config_tags,
},
kronos_dispatch::submit_job,
middlewares::auth_z::{Action as AuthZAction, AuthZ},
service::types::{
AppHeader, AppState, CustomHeaders, DbConnection, SchemaName, WorkspaceContext,
Expand All @@ -27,23 +29,26 @@ use superposition_core::helpers::{calculate_context_weight, hash};
use superposition_derives::{authorized, declare_resource};
use superposition_macros::{bad_argument, db_error, unexpected_error};
use superposition_types::{
Contextual, DBConnection, DimensionInfo, InternalUserContext, ListResponse,
Overridden, Overrides, PaginatedResponse, PrefixList, Resource, SortBy, User,
Contextual, DBConnection, DimensionInfo, InternalUserContext, Overridden, Overrides,
PaginatedResponse, PrefixList, Resource, SortBy, User,
api::{
DimensionMatchStrategy,
context::{
BulkOperation, BulkOperationResponse, ContextAction, ContextBulkResponse,
ContextListFilters, ContextValidationRequest, Identifier, MoveRequest,
PutRequest, SortOn, UpdateRequest, WeightRecomputeResponse,
},
jobs::{JobCreateResponse, JobRequest, PriorityRecomputeRequest},
webhook::Action,
},
custom_query::{
self as superposition_query, CustomQuery, DimensionQuery, PaginationParams,
QueryMap,
},
database::{
models::{ChangeReason, Description, cac::Context, others::WebhookEvent},
models::{
ChangeReason, Description, JobWorkspace, cac::Context, others::WebhookEvent,
},
schema::contexts::{self, dsl, id},
},
logic::evaluate_local_cohorts_skip_unresolved,
Expand Down Expand Up @@ -1161,21 +1166,16 @@ async fn bulk_operations_handler(
Ok(http_resp)
}

#[authorized]
#[put("/weight/recompute")]
async fn weight_recompute_handler(
workspace_context: WorkspaceContext,
state: Data<AppState>,
custom_headers: CustomHeaders,
mut write_permit: WorkspaceWritePermit,
user: User,
) -> superposition::Result<HttpResponse> {
pub async fn execute_priority_recompute(
workspace_context: &WorkspaceContext,
state: &Data<AppState>,
conn: &mut PooledConnection<ConnectionManager<PgConnection>>,
user: &User,
) -> superposition::Result<()> {
use superposition_types::database::schema::contexts::dsl::{
contexts, last_modified_at, last_modified_by, weight,
};

let conn = write_permit.connection();

let result: Vec<Context> = contexts
.schema_name(&workspace_context.schema_name)
.load(conn)
Expand All @@ -1187,7 +1187,6 @@ async fn weight_recompute_handler(
let dimension_info_map =
fetch_dimensions_info_map(conn, &workspace_context.schema_name)?;
let mut response: Vec<WeightRecomputeResponse> = vec![];
let tags = parse_config_tags(custom_headers.config_tags)?;

let contexts_new_weight = result
.clone()
Expand All @@ -1214,7 +1213,6 @@ async fn weight_recompute_handler(
})
.collect::<superposition::Result<Vec<(BigDecimal, String)>>>()?;

// Update database and add config version
let last_modified_time = Utc::now();
let config_version =
conn.transaction::<_, superposition::AppError, _>(|transaction_conn| {
Expand All @@ -1234,16 +1232,13 @@ async fn weight_recompute_handler(
db_error!(err)
})?;
}
let config_version_desc = Description::try_from("Recomputed weight".to_string()).map_err(|e| unexpected_error!(e))?;
add_config_version(&state, tags, config_version_desc, transaction_conn, &workspace_context.schema_name)
let config_version_desc = Description::try_from("Recomputed weight".to_string())
.map_err(|e| unexpected_error!(e))?;
add_config_version(state, None, config_version_desc, transaction_conn, &workspace_context.schema_name)
})?;
let _ = put_config_in_redis(
&config_version,
&state,
&workspace_context.schema_name,
conn,
)
.await;
let _ =
put_config_in_redis(&config_version, state, &workspace_context.schema_name, conn)
.await;

let data = WebhookData {
payload: &response,
Expand All @@ -1253,22 +1248,41 @@ async fn weight_recompute_handler(
action: Action::Batch(vec![Action::Update; response.len()]),
};

let webhook_status =
execute_webhook_call(data, &workspace_context, &state, conn).await;
let _ = execute_webhook_call(data, workspace_context, state, conn).await;
Ok(())
}

let mut http_resp = if webhook_status {
HttpResponse::Ok()
} else {
HttpResponse::build(
actix_web::http::StatusCode::from_u16(512)
.unwrap_or(actix_web::http::StatusCode::INTERNAL_SERVER_ERROR),
)
};
http_resp.insert_header((
AppHeader::XConfigVersion.to_string(),
config_version.id.to_string(),
));
Ok(http_resp.json(ListResponse::new(response)))
#[authorized]
#[put("/weight/recompute")]
async fn weight_recompute_handler(
workspace_context: WorkspaceContext,
state: Data<AppState>,
db_conn: DbConnection,
) -> superposition::Result<Json<JobCreateResponse>> {
let DbConnection(mut conn) = db_conn;
let job_request = JobRequest::PriorityRecompute(PriorityRecomputeRequest);
let job_workspace = JobWorkspace::from(&workspace_context.schema_name.0);
let target_workspace = state
.kronos_workspace
.as_deref()
.unwrap_or(&workspace_context.schema_name);

let response = submit_job(
state.kronos_client.as_ref(),
target_workspace,
&job_workspace,
&workspace_context.organisation_id,
&workspace_context.workspace_id,
job_request,
&state.snowflake_generator,
&mut conn,
3,
"Priority recompute job",
)
.await
.map_err(|e| unexpected_error!("Failed to submit priority recompute job: {}", e))?;

Ok(Json(response))
}

#[authorized]
Expand Down
Loading
Loading