Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
648d82f
fix(mst2): make retention acquisition atomic
Ivanbeethoven Oct 4, 2026
0ef0bd0
feat(mst2): add durable retention GC schema
Ivanbeethoven Oct 5, 2026
c8482bc
test(mst2): include retention migration in ordering check
Ivanbeethoven Oct 5, 2026
a39be11
fix(mst2): bound JSON body reads and preserve typed rejections
Ivanbeethoven Oct 5, 2026
c0391c4
Merge branch 'fix/mst2-config-schema' into fix/mst2-bounded-json-input
Ivanbeethoven Oct 5, 2026
ba5438d
fix(mst2): reject duplicate JSON keys including null and escapes
Ivanbeethoven Oct 5, 2026
c5e812f
Merge branch 'fix/mst2-config-schema' into fix/mst2-bounded-json-input
Ivanbeethoven Oct 5, 2026
8483226
feat(mst2): add atomic PostgreSQL retention repository
Ivanbeethoven Oct 5, 2026
481a867
Merge branch 'fix/mst2-config-schema' into fix/mst2-runtime-retain-at…
Ivanbeethoven Oct 5, 2026
8a87aa1
Merge branch 'fix/mst2-runtime-retain-atomic' into test/mst2-real-com…
Ivanbeethoven Oct 5, 2026
1301c00
Merge branch 'fix/mst2-bounded-json-input' into test/mst2-real-commit…
Ivanbeethoven Oct 5, 2026
01ac68f
fix(mst2): remove obsolete body extractor import
Ivanbeethoven Oct 5, 2026
2307730
Merge branch 'fix/mst2-bounded-json-input' into test/mst2-real-commit…
Ivanbeethoven Oct 5, 2026
c6c1d5c
test(mst2): account for integrated native and retention migrations
Ivanbeethoven Oct 5, 2026
19965e4
fix(mst2): repair native publication compile gates
Ivanbeethoven Oct 5, 2026
74f4373
test(mst2): require successful bounded-input lease fixture registration
Ivanbeethoven Oct 5, 2026
3400d62
Merge native publication compile fixes
Ivanbeethoven Oct 5, 2026
04cf79e
test(mst2): seed retention migration timestamps
Ivanbeethoven Oct 5, 2026
2375645
Merge retention migration fixture correction
Ivanbeethoven Oct 5, 2026
0873479
test(mst2): seed legacy retention edge timestamps
Ivanbeethoven Oct 5, 2026
5c6c24f
Merge legacy retention edge fixture correction
Ivanbeethoven Oct 5, 2026
7087a44
test(mst2): assert the GC receipt completion constraint
Ivanbeethoven Oct 5, 2026
77ad897
Merge precise GC completion constraint fixture
Ivanbeethoven Oct 5, 2026
37b6108
Merge native HTTP and receipt migration gate repairs
Ivanbeethoven Oct 5, 2026
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
7 changes: 3 additions & 4 deletions src/api/router/snapshot_content.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,11 @@ use axum::{
http::HeaderMap,
response::{IntoResponse, Response},
};
use bytes::Bytes;
use futures::stream::StreamExt;
use serde::Deserialize;
use serde_json::json;

use super::{abs_view_path, internal, mst2_error_response, treeframe_response};
use super::{abs_view_path, internal, mst2_error_response, request::Mst2Bytes, treeframe_response};
use crate::ceres::snapshot::{
chunks::{ChunkProjection, get_or_project},
error::{SnapshotError, SnapshotErrorCode},
Expand Down Expand Up @@ -174,7 +173,7 @@ const OBJECT_TOTAL_MAX: usize = 8 * 1024 * 1024;
pub(super) async fn objects(
state: State<crate::api::MonoApiServiceState>,
AxumPath(snapshot_id): AxumPath<String>,
body: Bytes,
Mst2Bytes(body): Mst2Bytes,
) -> Result<Response, Response> {
ensure(&state)?;
let ctx = runtime()
Expand Down Expand Up @@ -469,7 +468,7 @@ struct Planned {
pub(super) async fn chunks(
state: State<crate::api::MonoApiServiceState>,
AxumPath(snapshot_id): AxumPath<String>,
body: Bytes,
Mst2Bytes(body): Mst2Bytes,
) -> Result<Response, Response> {
ensure(&state)?;
let ctx = runtime()
Expand Down
114 changes: 114 additions & 0 deletions src/api/router/snapshot_request.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
//! Bounded raw JSON input for the MST/2 POST surface (spec 14).

use std::{collections::HashSet, fmt, time::Duration};

use axum::{
extract::{FromRequest, Request},
http::StatusCode,
};
use bytes::Bytes;
use serde::{
Deserialize, Deserializer,
de::{self, MapAccess, SeqAccess, Visitor},
};

use crate::ceres::snapshot::error::{SnapshotError, SnapshotErrorCode};

/// One overall read deadline, including a body which keeps trickling bytes.
/// This bounds input collection only, not handler work or response streams.
pub(super) const JSON_REQUEST_TIMEOUT: Duration = Duration::from_secs(10);

/// Preserve the original bytes for TreeFrame request-body digests. The router
/// supplies DefaultBodyLimit; its rejection is converted to the MST envelope.
pub(super) struct Mst2Bytes(pub(super) Bytes);

/// Decode keys before comparing them, including escaped spellings. This
/// separate pass also catches duplicate optional fields whose first value is
/// null, which a derived DTO can otherwise treat as an absent field.
pub(super) fn validate_json_keys(body: &[u8]) -> Result<(), serde_json::Error> {
serde_json::from_slice::<UniqueKeys>(body).map(|_| ())
}

struct UniqueKeys;

impl<'de> Deserialize<'de> for UniqueKeys {
fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
deserializer.deserialize_any(UniqueKeyVisitor)
}
}

struct UniqueKeyVisitor;

impl<'de> Visitor<'de> for UniqueKeyVisitor {
type Value = UniqueKeys;

fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("JSON without duplicate object keys")
}

fn visit_bool<E: de::Error>(self, _: bool) -> Result<UniqueKeys, E> {
Ok(UniqueKeys)
}
fn visit_i64<E: de::Error>(self, _: i64) -> Result<UniqueKeys, E> {
Ok(UniqueKeys)
}
fn visit_u64<E: de::Error>(self, _: u64) -> Result<UniqueKeys, E> {
Ok(UniqueKeys)
}
fn visit_f64<E: de::Error>(self, _: f64) -> Result<UniqueKeys, E> {
Ok(UniqueKeys)
}
fn visit_str<E: de::Error>(self, _: &str) -> Result<UniqueKeys, E> {
Ok(UniqueKeys)
}
fn visit_unit<E: de::Error>(self) -> Result<UniqueKeys, E> {
Ok(UniqueKeys)
}

fn visit_seq<A: SeqAccess<'de>>(self, mut sequence: A) -> Result<UniqueKeys, A::Error> {
while sequence.next_element::<UniqueKeys>()?.is_some() {}
Ok(UniqueKeys)
}

fn visit_map<A: MapAccess<'de>>(self, mut object: A) -> Result<UniqueKeys, A::Error> {
let mut keys = HashSet::new();
while let Some(key) = object.next_key::<String>()? {
if !keys.insert(key) {
return Err(de::Error::custom("duplicate JSON object key"));
}
object.next_value::<UniqueKeys>()?;
}
Ok(UniqueKeys)
}
}

impl<S> FromRequest<S> for Mst2Bytes
where
S: Send + Sync,
{
type Rejection = SnapshotError;

async fn from_request(req: Request, state: &S) -> Result<Self, Self::Rejection> {
match tokio::time::timeout(JSON_REQUEST_TIMEOUT, Bytes::from_request(req, state)).await {
Ok(Ok(body)) => Ok(Self(body)),
Ok(Err(error)) => {
let (code, message) = if error.status() == StatusCode::PAYLOAD_TOO_LARGE {
(
SnapshotErrorCode::LimitExceeded,
"request body over the spec 14 limit",
)
} else {
(
SnapshotErrorCode::InvalidRequest,
"could not read request body",
)
};
Err(SnapshotError::new(code, message))
}
Err(_) => Err(SnapshotError::new(
SnapshotErrorCode::TemporaryUnavailable,
"request body read deadline exceeded",
)),
}
}
}
Loading
Loading