Skip to main content

alopex_server/http/
kv.rs

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
121/// Searches raw KV key bytes with an explicit bounded glob or regex request.
122pub 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    // JSON byte arrays can expand every raw byte to three decimal digits plus a comma.
239    // Bound the raw page before serialization so the transport cap is not an allocation cap only.
240    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}