alopex_server/http/
session.rs1use std::sync::Arc;
2
3use axum::extract::{Extension, Path};
4use axum::response::Response;
5use serde::Serialize;
6
7use crate::error::{Result, ServerError};
8use crate::http::sql::sync_catalog_to_store;
9use crate::http::{error_response, json_response, RequestContext};
10use crate::server::ServerState;
11use crate::session::SessionId;
12
13#[derive(Serialize)]
14struct SessionBeginResponse {
15 session_id: String,
16 expires_at: String,
17}
18
19#[derive(Serialize)]
20struct SessionActionResponse {
21 success: bool,
22}
23
24pub async fn begin(
25 Extension(state): Extension<Arc<ServerState>>,
26 Extension(ctx): Extension<RequestContext>,
27) -> Response {
28 match begin_session(state.clone()).await {
29 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
30 Err(err) => error_response(err, &ctx),
31 }
32}
33
34pub async fn commit(
35 Extension(state): Extension<Arc<ServerState>>,
36 Extension(ctx): Extension<RequestContext>,
37 Path(id): Path<String>,
38) -> Response {
39 match session_action(state.clone(), &id, Action::Commit).await {
40 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
41 Err(err) => error_response(err, &ctx),
42 }
43}
44
45pub async fn rollback(
46 Extension(state): Extension<Arc<ServerState>>,
47 Extension(ctx): Extension<RequestContext>,
48 Path(id): Path<String>,
49) -> Response {
50 match session_action(state.clone(), &id, Action::Rollback).await {
51 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
52 Err(err) => error_response(err, &ctx),
53 }
54}
55
56async fn begin_session(state: Arc<ServerState>) -> Result<SessionBeginResponse> {
57 let session_id = state.session_manager.create_session().await?;
58 state.session_manager.begin_transaction(&session_id).await?;
59 let snapshot = state.session_manager.get_session(&session_id).await?;
60 let expires_at = chrono::DateTime::<chrono::Utc>::from(snapshot.expires_at);
61 Ok(SessionBeginResponse {
62 session_id: session_id.to_string(),
63 expires_at: expires_at.to_rfc3339(),
64 })
65}
66
67enum Action {
68 Commit,
69 Rollback,
70}
71
72async fn session_action(
73 state: Arc<ServerState>,
74 id: &str,
75 action: Action,
76) -> Result<SessionActionResponse> {
77 let session_id = id
78 .parse::<SessionId>()
79 .map_err(|_| ServerError::BadRequest("invalid session id".into()))?;
80 match action {
81 Action::Commit => {
82 let effects = state.session_manager.commit(&session_id).await?;
83 if !effects.is_empty() {
84 state.apply_table_lifecycle_effects(effects)?;
85 sync_catalog_to_store(&state)?;
86 }
87 }
88 Action::Rollback => {
89 let effects = state.session_manager.rollback(&session_id).await?;
90 state.apply_catalog_rollback_effects(effects)?;
91 }
92 }
93 Ok(SessionActionResponse { success: true })
94}