fix(manager-core): F-S-001 cap kv/docs/pubsub/queue payload sizes (default 256 KB)
Files (per-file cap), secrets (64 KB default), and email (25 MB default)
already enforce limits; kv::set, docs::create/update, pubsub::publish_durable
and queue::enqueue accepted any JSON value straight to a JSONB column with
no size validation. An anonymous public-HTTP script could fill disk via
queue::enqueue or amplify a single publish into N outbox rows × payload bytes.
Adds four new error variants:
- KvError::ValueTooLarge { limit, actual }
- DocsError::ValueTooLarge { limit, actual }
- PubsubError::MessageTooLarge { limit, actual }
- QueueError::PayloadTooLarge { limit, actual }
Each stateful service grows a `max_value_bytes` field with:
- Conservative 256 KB default (DEFAULT_KV_MAX_VALUE_BYTES etc.).
- New `with_max_*` constructor preserving the old `new()` signature.
- Env-knob reader (`*_max_*_from_env()`) — mirrors SecretsConfig::from_env.
Validation runs at the entry point BEFORE authz so an anonymous DoS doesn't
pay a membership lookup per attempt.
Wired via env knobs:
- PICLOUD_KV_MAX_VALUE_BYTES
- PICLOUD_DOCS_MAX_VALUE_BYTES
- PICLOUD_PUBSUB_MAX_MESSAGE_BYTES
- PICLOUD_QUEUE_MAX_PAYLOAD_BYTES
Documented in CLAUDE.md runtime config table.
New unit test in queue_service verifying the cap fires before authz.
AUDIT.md anchor: F-S-001. Depends on F-Q-004 (Backend variant).
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -34,10 +34,31 @@ use crate::authz::{self, AuthzRepo, Capability};
|
||||
use crate::docs_filter::{parse_filter, FilterParseError};
|
||||
use crate::docs_repo::{DocsRepo, DocsRepoError};
|
||||
|
||||
/// Default per-document JSON-encoded data cap (256 KB). Override with
|
||||
/// `PICLOUD_DOCS_MAX_VALUE_BYTES`.
|
||||
pub const DEFAULT_DOCS_MAX_VALUE_BYTES: usize = 256 * 1024;
|
||||
|
||||
/// Read `PICLOUD_DOCS_MAX_VALUE_BYTES`; invalid values fall back to the
|
||||
/// conservative default with a warning.
|
||||
#[must_use]
|
||||
pub fn docs_max_value_bytes_from_env() -> usize {
|
||||
if let Ok(v) = std::env::var("PICLOUD_DOCS_MAX_VALUE_BYTES") {
|
||||
match v.trim().parse::<usize>() {
|
||||
Ok(n) if n > 0 => return n,
|
||||
_ => tracing::warn!(
|
||||
value = %v,
|
||||
"ignoring invalid PICLOUD_DOCS_MAX_VALUE_BYTES (want a positive integer)"
|
||||
),
|
||||
}
|
||||
}
|
||||
DEFAULT_DOCS_MAX_VALUE_BYTES
|
||||
}
|
||||
|
||||
pub struct DocsServiceImpl {
|
||||
repo: Arc<dyn DocsRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
events: Arc<dyn ServiceEventEmitter>,
|
||||
max_value_bytes: usize,
|
||||
}
|
||||
|
||||
impl DocsServiceImpl {
|
||||
@@ -46,14 +67,38 @@ impl DocsServiceImpl {
|
||||
repo: Arc<dyn DocsRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
events: Arc<dyn ServiceEventEmitter>,
|
||||
) -> Self {
|
||||
Self::with_max_value_bytes(repo, authz, events, DEFAULT_DOCS_MAX_VALUE_BYTES)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn with_max_value_bytes(
|
||||
repo: Arc<dyn DocsRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
events: Arc<dyn ServiceEventEmitter>,
|
||||
max_value_bytes: usize,
|
||||
) -> Self {
|
||||
Self {
|
||||
repo,
|
||||
authz,
|
||||
events,
|
||||
max_value_bytes,
|
||||
}
|
||||
}
|
||||
|
||||
fn check_data_size(&self, data: &serde_json::Value) -> Result<(), DocsError> {
|
||||
let encoded_len = serde_json::to_vec(data)
|
||||
.map(|v| v.len())
|
||||
.map_err(|e| DocsError::Backend(format!("encode doc data: {e}")))?;
|
||||
if encoded_len > self.max_value_bytes {
|
||||
return Err(DocsError::ValueTooLarge {
|
||||
limit: self.max_value_bytes,
|
||||
actual: encoded_len,
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn check_read(&self, cx: &SdkCallCx) -> Result<(), DocsError> {
|
||||
authz::script_gate(
|
||||
&*self.authz,
|
||||
@@ -116,6 +161,7 @@ impl DocsService for DocsServiceImpl {
|
||||
) -> Result<DocId, DocsError> {
|
||||
validate_collection(collection)?;
|
||||
validate_data(&data)?;
|
||||
self.check_data_size(&data)?;
|
||||
self.check_write(cx).await?;
|
||||
let row = self
|
||||
.repo
|
||||
@@ -193,6 +239,7 @@ impl DocsService for DocsServiceImpl {
|
||||
) -> Result<(), DocsError> {
|
||||
validate_collection(collection)?;
|
||||
validate_data(&data)?;
|
||||
self.check_data_size(&data)?;
|
||||
self.check_write(cx).await?;
|
||||
let previous = self
|
||||
.repo
|
||||
|
||||
@@ -25,10 +25,32 @@ use picloud_shared::{
|
||||
use crate::authz::{self, AuthzRepo, Capability};
|
||||
use crate::kv_repo::{KvRepo, KvRepoError};
|
||||
|
||||
/// Default per-key JSON-encoded value cap (256 KB). Override with
|
||||
/// `PICLOUD_KV_MAX_VALUE_BYTES`.
|
||||
pub const DEFAULT_KV_MAX_VALUE_BYTES: usize = 256 * 1024;
|
||||
|
||||
/// Read `PICLOUD_KV_MAX_VALUE_BYTES`; invalid values fall back to the
|
||||
/// conservative default with a warning. Public so the picloud binary can
|
||||
/// thread it through `KvServiceImpl::new`.
|
||||
#[must_use]
|
||||
pub fn kv_max_value_bytes_from_env() -> usize {
|
||||
if let Ok(v) = std::env::var("PICLOUD_KV_MAX_VALUE_BYTES") {
|
||||
match v.trim().parse::<usize>() {
|
||||
Ok(n) if n > 0 => return n,
|
||||
_ => tracing::warn!(
|
||||
value = %v,
|
||||
"ignoring invalid PICLOUD_KV_MAX_VALUE_BYTES (want a positive integer)"
|
||||
),
|
||||
}
|
||||
}
|
||||
DEFAULT_KV_MAX_VALUE_BYTES
|
||||
}
|
||||
|
||||
pub struct KvServiceImpl {
|
||||
repo: Arc<dyn KvRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
events: Arc<dyn ServiceEventEmitter>,
|
||||
max_value_bytes: usize,
|
||||
}
|
||||
|
||||
impl KvServiceImpl {
|
||||
@@ -37,11 +59,22 @@ impl KvServiceImpl {
|
||||
repo: Arc<dyn KvRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
events: Arc<dyn ServiceEventEmitter>,
|
||||
) -> Self {
|
||||
Self::with_max_value_bytes(repo, authz, events, DEFAULT_KV_MAX_VALUE_BYTES)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn with_max_value_bytes(
|
||||
repo: Arc<dyn KvRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
events: Arc<dyn ServiceEventEmitter>,
|
||||
max_value_bytes: usize,
|
||||
) -> Self {
|
||||
Self {
|
||||
repo,
|
||||
authz,
|
||||
events,
|
||||
max_value_bytes,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -102,6 +135,15 @@ impl KvService for KvServiceImpl {
|
||||
value: serde_json::Value,
|
||||
) -> Result<(), KvError> {
|
||||
validate_collection(collection)?;
|
||||
let encoded_len = serde_json::to_vec(&value)
|
||||
.map(|v| v.len())
|
||||
.map_err(|e| KvError::Backend(format!("encode value: {e}")))?;
|
||||
if encoded_len > self.max_value_bytes {
|
||||
return Err(KvError::ValueTooLarge {
|
||||
limit: self.max_value_bytes,
|
||||
actual: encoded_len,
|
||||
});
|
||||
}
|
||||
self.check_write(cx).await?;
|
||||
let previous = self
|
||||
.repo
|
||||
|
||||
@@ -70,6 +70,26 @@ fn load_i64(dst: &mut i64, key: &str) {
|
||||
}
|
||||
}
|
||||
|
||||
/// Default per-message JSON-encoded payload cap (256 KB). Override with
|
||||
/// `PICLOUD_PUBSUB_MAX_MESSAGE_BYTES`.
|
||||
pub const DEFAULT_PUBSUB_MAX_MESSAGE_BYTES: usize = 256 * 1024;
|
||||
|
||||
/// Read `PICLOUD_PUBSUB_MAX_MESSAGE_BYTES`; invalid values fall back to
|
||||
/// the conservative default with a warning.
|
||||
#[must_use]
|
||||
pub fn pubsub_max_message_bytes_from_env() -> usize {
|
||||
if let Ok(v) = std::env::var("PICLOUD_PUBSUB_MAX_MESSAGE_BYTES") {
|
||||
match v.trim().parse::<usize>() {
|
||||
Ok(n) if n > 0 => return n,
|
||||
_ => tracing::warn!(
|
||||
value = %v,
|
||||
"ignoring invalid PICLOUD_PUBSUB_MAX_MESSAGE_BYTES (want a positive integer)"
|
||||
),
|
||||
}
|
||||
}
|
||||
DEFAULT_PUBSUB_MAX_MESSAGE_BYTES
|
||||
}
|
||||
|
||||
pub struct PubsubServiceImpl {
|
||||
repo: Arc<dyn PubsubRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
@@ -80,6 +100,7 @@ pub struct PubsubServiceImpl {
|
||||
topics: Option<Arc<dyn TopicRepo>>,
|
||||
secrets: Option<Arc<dyn AppSecretsRepo>>,
|
||||
token_config: SubscriberTokenConfig,
|
||||
max_message_bytes: usize,
|
||||
}
|
||||
|
||||
impl PubsubServiceImpl {
|
||||
@@ -92,9 +113,16 @@ impl PubsubServiceImpl {
|
||||
topics: None,
|
||||
secrets: None,
|
||||
token_config: SubscriberTokenConfig::conservative(),
|
||||
max_message_bytes: DEFAULT_PUBSUB_MAX_MESSAGE_BYTES,
|
||||
}
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn with_max_message_bytes(mut self, n: usize) -> Self {
|
||||
self.max_message_bytes = n;
|
||||
self
|
||||
}
|
||||
|
||||
/// Attach the v1.1.6 realtime surface: the in-process broadcaster
|
||||
/// (publish fan-out to SSE subscribers), the topic registry +
|
||||
/// app-secrets repo (subscriber-token minting), and the TTL config.
|
||||
@@ -142,6 +170,17 @@ impl PubsubService for PubsubServiceImpl {
|
||||
if topic.trim().is_empty() {
|
||||
return Err(PubsubError::EmptyTopic);
|
||||
}
|
||||
// Reject oversized messages BEFORE the authz check so an
|
||||
// anonymous public-script DoS doesn't pay the membership lookup.
|
||||
let encoded_len = serde_json::to_vec(&message)
|
||||
.map(|v| v.len())
|
||||
.map_err(|e| PubsubError::Rejected(format!("encode message: {e}")))?;
|
||||
if encoded_len > self.max_message_bytes {
|
||||
return Err(PubsubError::MessageTooLarge {
|
||||
limit: self.max_message_bytes,
|
||||
actual: encoded_len,
|
||||
});
|
||||
}
|
||||
self.check_publish(cx).await?;
|
||||
|
||||
// `published_at` is stamped once on the manager side so every
|
||||
|
||||
@@ -15,16 +15,50 @@ use picloud_shared::{EnqueueOpts, QueueError, QueueMessageId, QueueService, SdkC
|
||||
use crate::authz::{self, AuthzRepo, Capability};
|
||||
use crate::queue_repo::{NewQueueMessage, QueueRepo};
|
||||
|
||||
/// Default per-message JSON-encoded payload cap (256 KB). Override with
|
||||
/// `PICLOUD_QUEUE_MAX_PAYLOAD_BYTES`.
|
||||
pub const DEFAULT_QUEUE_MAX_PAYLOAD_BYTES: usize = 256 * 1024;
|
||||
|
||||
/// Read `PICLOUD_QUEUE_MAX_PAYLOAD_BYTES`; invalid values fall back to
|
||||
/// the conservative default with a warning.
|
||||
#[must_use]
|
||||
pub fn queue_max_payload_bytes_from_env() -> usize {
|
||||
if let Ok(v) = std::env::var("PICLOUD_QUEUE_MAX_PAYLOAD_BYTES") {
|
||||
match v.trim().parse::<usize>() {
|
||||
Ok(n) if n > 0 => return n,
|
||||
_ => tracing::warn!(
|
||||
value = %v,
|
||||
"ignoring invalid PICLOUD_QUEUE_MAX_PAYLOAD_BYTES (want a positive integer)"
|
||||
),
|
||||
}
|
||||
}
|
||||
DEFAULT_QUEUE_MAX_PAYLOAD_BYTES
|
||||
}
|
||||
|
||||
/// Production impl: authz gate → repo. Trivial wrapper.
|
||||
pub struct QueueServiceImpl {
|
||||
repo: Arc<dyn QueueRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
max_payload_bytes: usize,
|
||||
}
|
||||
|
||||
impl QueueServiceImpl {
|
||||
#[must_use]
|
||||
pub fn new(repo: Arc<dyn QueueRepo>, authz: Arc<dyn AuthzRepo>) -> Self {
|
||||
Self { repo, authz }
|
||||
Self::with_max_payload_bytes(repo, authz, DEFAULT_QUEUE_MAX_PAYLOAD_BYTES)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn with_max_payload_bytes(
|
||||
repo: Arc<dyn QueueRepo>,
|
||||
authz: Arc<dyn AuthzRepo>,
|
||||
max_payload_bytes: usize,
|
||||
) -> Self {
|
||||
Self {
|
||||
repo,
|
||||
authz,
|
||||
max_payload_bytes,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -40,6 +74,17 @@ impl QueueService for QueueServiceImpl {
|
||||
if queue_name.is_empty() {
|
||||
return Err(QueueError::EmptyName);
|
||||
}
|
||||
// Reject oversized payloads BEFORE the authz check so an
|
||||
// anonymous public-script DoS doesn't pay the membership lookup.
|
||||
let encoded_len = serde_json::to_vec(&payload)
|
||||
.map(|v| v.len())
|
||||
.map_err(|e| QueueError::Rejected(format!("encode payload: {e}")))?;
|
||||
if encoded_len > self.max_payload_bytes {
|
||||
return Err(QueueError::PayloadTooLarge {
|
||||
limit: self.max_payload_bytes,
|
||||
actual: encoded_len,
|
||||
});
|
||||
}
|
||||
let max_attempts = opts.max_attempts.unwrap_or(3);
|
||||
if !(1..=20).contains(&max_attempts) {
|
||||
return Err(QueueError::InvalidOpts(
|
||||
@@ -311,6 +356,24 @@ mod tests {
|
||||
assert!(captured.deliver_after.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn oversized_payload_is_rejected_before_authz() {
|
||||
let repo = Arc::new(CapturingRepo {
|
||||
last: tokio::sync::Mutex::new(None),
|
||||
});
|
||||
let svc = QueueServiceImpl::with_max_payload_bytes(repo, Arc::new(AlwaysAllowAuthz), 16);
|
||||
let cx = anon_cx();
|
||||
let payload = serde_json::json!({ "blob": "x".repeat(200) });
|
||||
let err = svc
|
||||
.enqueue(&cx, "jobs", payload, EnqueueOpts::default())
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(
|
||||
matches!(err, QueueError::PayloadTooLarge { .. }),
|
||||
"expected PayloadTooLarge, got {err:?}"
|
||||
);
|
||||
}
|
||||
|
||||
struct DenyAllAuthz;
|
||||
#[async_trait]
|
||||
impl AuthzRepo for DenyAllAuthz {
|
||||
|
||||
Reference in New Issue
Block a user