1use std::sync::Arc;
2
3use alopex_core::kv::KVTransaction;
4use alopex_core::kv::{KeySearchPage, KeySearchRequest};
5use alopex_core::types::TxnMode;
6use alopex_core::KVStore;
7use axum::extract::Extension;
8use axum::response::Response;
9use axum::Json;
10use serde::{Deserialize, Serialize};
11use std::time::{SystemTime, UNIX_EPOCH};
12
13use crate::error::{Result, ServerError};
14use crate::http::{error_response, json_response, RequestContext};
15use crate::server::ServerState;
16
17#[derive(Debug, Deserialize)]
18pub struct KvGetRequest {
19 pub key: String,
20}
21
22#[derive(Debug, Deserialize)]
23pub struct KvPutRequest {
24 pub key: String,
25 pub value: Vec<u8>,
26}
27
28#[derive(Debug, Deserialize)]
29pub struct KvDeleteRequest {
30 pub key: String,
31}
32
33#[derive(Debug, Deserialize)]
34pub struct KvListRequest {
35 pub prefix: Option<String>,
36}
37
38#[derive(Debug, Deserialize)]
39pub struct KvTxnBeginRequest {
40 pub timeout_secs: Option<u64>,
41}
42
43#[derive(Debug, Deserialize)]
44pub struct KvTxnRequest {
45 pub txn_id: String,
46 pub key: Option<String>,
47 pub value: Option<Vec<u8>>,
48}
49
50#[derive(Debug, Serialize)]
51pub struct KvGetResponse {
52 pub key: Vec<u8>,
53 pub value: Option<Vec<u8>>,
54}
55
56#[derive(Debug, Serialize)]
57pub struct KvListEntry {
58 pub key: Vec<u8>,
59 pub value: Vec<u8>,
60}
61
62#[derive(Debug, Serialize)]
63pub struct KvListResponse {
64 pub entries: Vec<KvListEntry>,
65}
66
67#[derive(Debug, Serialize)]
68pub struct KvStatusResponse {
69 pub success: bool,
70}
71
72#[derive(Debug, Serialize)]
73pub struct KvTxnBeginResponse {
74 pub txn_id: String,
75}
76
77pub async fn get(
78 Extension(state): Extension<Arc<ServerState>>,
79 Extension(ctx): Extension<RequestContext>,
80 Json(request): Json<KvGetRequest>,
81) -> Response {
82 match get_impl(state.clone(), request) {
83 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
84 Err(err) => error_response(err, &ctx),
85 }
86}
87
88pub async fn put(
89 Extension(state): Extension<Arc<ServerState>>,
90 Extension(ctx): Extension<RequestContext>,
91 Json(request): Json<KvPutRequest>,
92) -> Response {
93 match put_impl(state.clone(), request) {
94 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
95 Err(err) => error_response(err, &ctx),
96 }
97}
98
99pub async fn delete(
100 Extension(state): Extension<Arc<ServerState>>,
101 Extension(ctx): Extension<RequestContext>,
102 Json(request): Json<KvDeleteRequest>,
103) -> Response {
104 match delete_impl(state.clone(), request) {
105 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
106 Err(err) => error_response(err, &ctx),
107 }
108}
109
110pub async fn list(
111 Extension(state): Extension<Arc<ServerState>>,
112 Extension(ctx): Extension<RequestContext>,
113 Json(request): Json<KvListRequest>,
114) -> Response {
115 match list_impl(state.clone(), request) {
116 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
117 Err(err) => error_response(err, &ctx),
118 }
119}
120
121pub async fn search(
123 Extension(state): Extension<Arc<ServerState>>,
124 Extension(ctx): Extension<RequestContext>,
125 Json(request): Json<KeySearchRequest>,
126) -> Response {
127 match search_impl(state.clone(), request) {
128 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
129 Err(err) => error_response(err, &ctx),
130 }
131}
132
133pub async fn txn_begin(
134 Extension(state): Extension<Arc<ServerState>>,
135 Extension(ctx): Extension<RequestContext>,
136 Json(request): Json<KvTxnBeginRequest>,
137) -> Response {
138 match txn_begin_impl(state.clone(), request) {
139 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
140 Err(err) => error_response(err, &ctx),
141 }
142}
143
144pub async fn txn_get(
145 Extension(state): Extension<Arc<ServerState>>,
146 Extension(ctx): Extension<RequestContext>,
147 Json(request): Json<KvTxnRequest>,
148) -> Response {
149 match txn_get_impl(state.clone(), request) {
150 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
151 Err(err) => error_response(err, &ctx),
152 }
153}
154
155pub async fn txn_put(
156 Extension(state): Extension<Arc<ServerState>>,
157 Extension(ctx): Extension<RequestContext>,
158 Json(request): Json<KvTxnRequest>,
159) -> Response {
160 match txn_put_impl(state.clone(), request) {
161 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
162 Err(err) => error_response(err, &ctx),
163 }
164}
165
166pub async fn txn_delete(
167 Extension(state): Extension<Arc<ServerState>>,
168 Extension(ctx): Extension<RequestContext>,
169 Json(request): Json<KvTxnRequest>,
170) -> Response {
171 match txn_delete_impl(state.clone(), request) {
172 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
173 Err(err) => error_response(err, &ctx),
174 }
175}
176
177pub async fn txn_commit(
178 Extension(state): Extension<Arc<ServerState>>,
179 Extension(ctx): Extension<RequestContext>,
180 Json(request): Json<KvTxnRequest>,
181) -> Response {
182 match txn_commit_impl(state.clone(), request) {
183 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
184 Err(err) => error_response(err, &ctx),
185 }
186}
187
188pub async fn txn_rollback(
189 Extension(state): Extension<Arc<ServerState>>,
190 Extension(ctx): Extension<RequestContext>,
191 Json(request): Json<KvTxnRequest>,
192) -> Response {
193 match txn_rollback_impl(state.clone(), request) {
194 Ok(resp) => json_response(resp, state.config.max_response_size, &ctx),
195 Err(err) => error_response(err, &ctx),
196 }
197}
198
199fn get_impl(state: Arc<ServerState>, request: KvGetRequest) -> Result<KvGetResponse> {
200 let mut txn = state.store.begin(TxnMode::ReadOnly)?;
201 let key_bytes = request.key.into_bytes();
202 let value = txn.get(&key_bytes)?;
203 txn.commit_self()?;
204 Ok(KvGetResponse {
205 key: key_bytes,
206 value,
207 })
208}
209
210fn put_impl(state: Arc<ServerState>, request: KvPutRequest) -> Result<KvStatusResponse> {
211 state.lifecycle_state.check_write_allowed()?;
212 let mut txn = state.store.begin(TxnMode::ReadWrite)?;
213 txn.put(request.key.into_bytes(), request.value)?;
214 txn.commit_self()?;
215 Ok(KvStatusResponse { success: true })
216}
217
218fn delete_impl(state: Arc<ServerState>, request: KvDeleteRequest) -> Result<KvStatusResponse> {
219 state.lifecycle_state.check_write_allowed()?;
220 let mut txn = state.store.begin(TxnMode::ReadWrite)?;
221 txn.delete(request.key.into_bytes())?;
222 txn.commit_self()?;
223 Ok(KvStatusResponse { success: true })
224}
225
226fn list_impl(state: Arc<ServerState>, request: KvListRequest) -> Result<KvListResponse> {
227 let mut txn = state.store.begin(TxnMode::ReadOnly)?;
228 let prefix = request.prefix.unwrap_or_default();
229 let mut entries = Vec::new();
230 for (key, value) in txn.scan_prefix(prefix.as_bytes())? {
231 entries.push(KvListEntry { key, value });
232 }
233 txn.commit_self()?;
234 Ok(KvListResponse { entries })
235}
236
237fn search_impl(state: Arc<ServerState>, mut request: KeySearchRequest) -> Result<KeySearchPage> {
238 let transport_raw_limit = (state.config.max_response_size / 4).max(1);
241 request.max_bytes = request.max_bytes.min(transport_raw_limit);
242 let mut txn = state.store.begin(TxnMode::ReadOnly)?;
243 let page = txn.search_keys(&request).map_err(map_search_error)?;
244 txn.commit_self()?;
245 Ok(page)
246}
247
248fn map_search_error(error: alopex_core::Error) -> ServerError {
249 match error {
250 alopex_core::Error::InvalidParameter { param, reason } => {
251 ServerError::BadRequest(format!("invalid {param}: {reason}"))
252 }
253 alopex_core::Error::SearchBudgetExceeded { limit } => {
254 ServerError::PayloadTooLarge(format!("search scan budget exceeded: limit={limit}"))
255 }
256 alopex_core::Error::SearchResponseTooLarge { limit, requested } => {
257 ServerError::PayloadTooLarge(format!(
258 "search response size exceeded: limit={limit}, requested={requested}"
259 ))
260 }
261 other => ServerError::Core(other),
262 }
263}
264
265fn txn_begin_impl(
266 state: Arc<ServerState>,
267 request: KvTxnBeginRequest,
268) -> Result<KvTxnBeginResponse> {
269 state.lifecycle_state.check_write_allowed()?;
270 let timeout_secs = request.timeout_secs.unwrap_or(DEFAULT_TXN_TIMEOUT_SECS);
271 let meta = TxnMeta {
272 started_at_secs: current_timestamp_secs(),
273 timeout_secs,
274 };
275 let txn_id = generate_txn_id();
276 let mut txn = state.store.begin(TxnMode::ReadWrite)?;
277 txn.put(txn_meta_key(&txn_id), encode_meta(meta))?;
278 txn.commit_self()?;
279 Ok(KvTxnBeginResponse { txn_id })
280}
281
282fn txn_get_impl(state: Arc<ServerState>, request: KvTxnRequest) -> Result<KvGetResponse> {
283 let key = request
284 .key
285 .ok_or_else(|| ServerError::BadRequest("key is required".into()))?;
286 let mut txn = state.store.begin(TxnMode::ReadOnly)?;
287 let meta = load_meta(&mut txn, &request.txn_id)?;
288 if is_expired_from_meta(meta, current_timestamp_secs()) {
289 txn.commit_self()?;
290 rollback_transaction(state.clone(), &request.txn_id)?;
291 return Err(ServerError::SessionExpired("transaction expired".into()));
292 }
293 let value = if let Some(raw) = txn.get(&txn_write_key(&request.txn_id, key.as_bytes()))? {
294 match decode_write(&request.txn_id, &raw)? {
295 TxnWrite::Put(value) => Some(value),
296 TxnWrite::Delete => None,
297 }
298 } else {
299 txn.get(&key.as_bytes().to_vec())?
300 };
301 txn.commit_self()?;
302 Ok(KvGetResponse {
303 key: key.into_bytes(),
304 value,
305 })
306}
307
308fn txn_put_impl(state: Arc<ServerState>, request: KvTxnRequest) -> Result<KvStatusResponse> {
309 state.lifecycle_state.check_write_allowed()?;
310 let key = request
311 .key
312 .ok_or_else(|| ServerError::BadRequest("key is required".into()))?;
313 let value = request
314 .value
315 .ok_or_else(|| ServerError::BadRequest("value is required".into()))?;
316 let mut txn = state.store.begin(TxnMode::ReadWrite)?;
317 let meta = load_meta(&mut txn, &request.txn_id)?;
318 if is_expired_from_meta(meta, current_timestamp_secs()) {
319 txn.rollback_self()?;
320 rollback_transaction(state.clone(), &request.txn_id)?;
321 return Err(ServerError::SessionExpired("transaction expired".into()));
322 }
323 txn.put(
324 txn_write_key(&request.txn_id, key.as_bytes()),
325 encode_write(TxnWrite::Put(value)),
326 )?;
327 txn.commit_self()?;
328 Ok(KvStatusResponse { success: true })
329}
330
331fn txn_delete_impl(state: Arc<ServerState>, request: KvTxnRequest) -> Result<KvStatusResponse> {
332 state.lifecycle_state.check_write_allowed()?;
333 let key = request
334 .key
335 .ok_or_else(|| ServerError::BadRequest("key is required".into()))?;
336 let mut txn = state.store.begin(TxnMode::ReadWrite)?;
337 let meta = load_meta(&mut txn, &request.txn_id)?;
338 if is_expired_from_meta(meta, current_timestamp_secs()) {
339 txn.rollback_self()?;
340 rollback_transaction(state.clone(), &request.txn_id)?;
341 return Err(ServerError::SessionExpired("transaction expired".into()));
342 }
343 txn.put(
344 txn_write_key(&request.txn_id, key.as_bytes()),
345 encode_write(TxnWrite::Delete),
346 )?;
347 txn.commit_self()?;
348 Ok(KvStatusResponse { success: true })
349}
350
351fn txn_commit_impl(state: Arc<ServerState>, request: KvTxnRequest) -> Result<KvStatusResponse> {
352 state.lifecycle_state.check_write_allowed()?;
353 commit_transaction(state, &request.txn_id)?;
354 Ok(KvStatusResponse { success: true })
355}
356
357fn txn_rollback_impl(state: Arc<ServerState>, request: KvTxnRequest) -> Result<KvStatusResponse> {
358 state.lifecycle_state.check_write_allowed()?;
359 rollback_transaction(state, &request.txn_id)?;
360 Ok(KvStatusResponse { success: true })
361}
362
363fn commit_transaction(state: Arc<ServerState>, txn_id: &str) -> Result<()> {
364 let mut txn = state.store.begin(TxnMode::ReadWrite)?;
365 let meta = load_meta(&mut txn, txn_id)?;
366 if is_expired_from_meta(meta, current_timestamp_secs()) {
367 txn.rollback_self()?;
368 rollback_transaction(state, txn_id)?;
369 return Err(ServerError::SessionExpired("transaction expired".into()));
370 }
371 let prefix = txn_write_prefix(txn_id);
372 let staged: Vec<(Vec<u8>, Vec<u8>)> = txn.scan_prefix(&prefix)?.collect();
373 for (staged_key, raw) in &staged {
374 let user_key = extract_user_key(txn_id, staged_key)?;
375 match decode_write(txn_id, raw)? {
376 TxnWrite::Put(value) => {
377 txn.put(user_key, value)?;
378 }
379 TxnWrite::Delete => {
380 txn.delete(user_key)?;
381 }
382 }
383 }
384 for (staged_key, _) in staged {
385 txn.delete(staged_key)?;
386 }
387 txn.delete(txn_meta_key(txn_id))?;
388 txn.commit_self()?;
389 Ok(())
390}
391
392fn rollback_transaction(state: Arc<ServerState>, txn_id: &str) -> Result<()> {
393 let mut txn = state.store.begin(TxnMode::ReadWrite)?;
394 let _ = load_meta(&mut txn, txn_id)?;
395 let prefix = txn_write_prefix(txn_id);
396 let staged: Vec<(Vec<u8>, Vec<u8>)> = txn.scan_prefix(&prefix)?.collect();
397 for (staged_key, _) in staged {
398 txn.delete(staged_key)?;
399 }
400 txn.delete(txn_meta_key(txn_id))?;
401 txn.commit_self()?;
402 Ok(())
403}
404
405const DEFAULT_TXN_TIMEOUT_SECS: u64 = 60;
406const TXN_META_PREFIX: &[u8] = b"__alopex_txn_meta__:";
407const TXN_WRITE_PREFIX: &[u8] = b"__alopex_txn_write__:";
408const TXN_WRITE_DELETE: u8 = 0;
409const TXN_WRITE_PUT: u8 = 1;
410
411#[derive(Debug, Clone, Copy)]
412struct TxnMeta {
413 started_at_secs: u64,
414 timeout_secs: u64,
415}
416
417enum TxnWrite {
418 Put(Vec<u8>),
419 Delete,
420}
421
422fn current_timestamp_secs() -> u64 {
423 SystemTime::now()
424 .duration_since(UNIX_EPOCH)
425 .unwrap_or_default()
426 .as_secs()
427}
428
429fn generate_txn_id() -> String {
430 let nanos = SystemTime::now()
431 .duration_since(UNIX_EPOCH)
432 .unwrap_or_default()
433 .as_nanos();
434 format!("txn-{}-{}", nanos, std::process::id())
435}
436
437fn txn_meta_key(txn_id: &str) -> Vec<u8> {
438 let mut key = Vec::with_capacity(TXN_META_PREFIX.len() + txn_id.len());
439 key.extend_from_slice(TXN_META_PREFIX);
440 key.extend_from_slice(txn_id.as_bytes());
441 key
442}
443
444fn txn_write_prefix(txn_id: &str) -> Vec<u8> {
445 let mut key = Vec::with_capacity(TXN_WRITE_PREFIX.len() + txn_id.len() + 1);
446 key.extend_from_slice(TXN_WRITE_PREFIX);
447 key.extend_from_slice(txn_id.as_bytes());
448 key.push(b':');
449 key
450}
451
452fn txn_write_key(txn_id: &str, key: &[u8]) -> Vec<u8> {
453 let mut full = txn_write_prefix(txn_id);
454 full.extend_from_slice(key);
455 full
456}
457
458fn encode_meta(meta: TxnMeta) -> Vec<u8> {
459 let mut payload = Vec::with_capacity(16);
460 payload.extend_from_slice(&meta.started_at_secs.to_le_bytes());
461 payload.extend_from_slice(&meta.timeout_secs.to_le_bytes());
462 payload
463}
464
465fn decode_meta(txn_id: &str, raw: &[u8]) -> Result<TxnMeta> {
466 if raw.len() < 16 {
467 return Err(ServerError::BadRequest(format!(
468 "transaction metadata invalid: {}",
469 txn_id
470 )));
471 }
472 let started_at_secs = u64::from_le_bytes(raw[0..8].try_into().unwrap());
473 let timeout_secs = u64::from_le_bytes(raw[8..16].try_into().unwrap());
474 Ok(TxnMeta {
475 started_at_secs,
476 timeout_secs,
477 })
478}
479
480fn load_meta(
481 txn: &mut alopex_core::kv::any::AnyKVTransaction<'_>,
482 txn_id: &str,
483) -> Result<TxnMeta> {
484 let Some(raw) = txn.get(&txn_meta_key(txn_id))? else {
485 return Err(ServerError::NotFound("transaction not found".into()));
486 };
487 decode_meta(txn_id, &raw)
488}
489
490fn is_expired_from_meta(meta: TxnMeta, now_secs: u64) -> bool {
491 now_secs.saturating_sub(meta.started_at_secs) >= meta.timeout_secs
492}
493
494fn encode_write(entry: TxnWrite) -> Vec<u8> {
495 match entry {
496 TxnWrite::Put(value) => {
497 let mut payload = Vec::with_capacity(1 + value.len());
498 payload.push(TXN_WRITE_PUT);
499 payload.extend_from_slice(&value);
500 payload
501 }
502 TxnWrite::Delete => vec![TXN_WRITE_DELETE],
503 }
504}
505
506fn decode_write(txn_id: &str, raw: &[u8]) -> Result<TxnWrite> {
507 let Some((&tag, rest)) = raw.split_first() else {
508 return Err(ServerError::BadRequest(format!(
509 "transaction write entry invalid: {}",
510 txn_id
511 )));
512 };
513 match tag {
514 TXN_WRITE_PUT => Ok(TxnWrite::Put(rest.to_vec())),
515 TXN_WRITE_DELETE => Ok(TxnWrite::Delete),
516 _ => Err(ServerError::BadRequest(format!(
517 "transaction write entry invalid: {}",
518 txn_id
519 ))),
520 }
521}
522
523fn extract_user_key(txn_id: &str, staged_key: &[u8]) -> Result<Vec<u8>> {
524 let prefix = txn_write_prefix(txn_id);
525 if !staged_key.starts_with(&prefix) {
526 return Err(ServerError::BadRequest(format!(
527 "transaction write key invalid: {}",
528 txn_id
529 )));
530 }
531 Ok(staged_key[prefix.len()..].to_vec())
532}