feat(v1.1.9): QueueService trait + Postgres impl + Services bundle wiring
- shared/services.rs: Services::new gains queue: Arc<dyn QueueService> positionally after users (mirrors v1.1.8's append of users). with_noop_services adds NoopQueueService. - manager-core/queue_service.rs: QueueServiceImpl wraps QueueRepo with script-as-gate authz on AppQueueEnqueue. enqueue clamps max_attempts to [1,20] and delay_ms to [0, 86_400_000ms]. depth/depth_pending are read-only — no authz check (scripts in the app can see their own queue depths). cx.principal threads through as enqueued_by_principal (forensic only). - manager-core/authz.rs: AppQueueEnqueue(AppId) capability — script:write scope, granted to editor+ (same trust shape as AppPubsubPublish). - picloud/lib.rs: wires PostgresQueueRepo + QueueServiceImpl into Services::new alongside the existing v1.1.7+ services. - 11 sdk test binaries + manager-core/realtime_authority.rs updated to pass NoopQueueService to Services::new. Unit tests cover empty queue_name reject, max_attempts clamping (0/21 → invalid), delay_ms negative-reject, anonymous principal skips authz, depth/depth_pending pass-through. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -104,6 +104,7 @@ async fn original_backend_error_is_logged_at_error_level() {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
let engine = Engine::new(Limits::default(), services);
|
let engine = Engine::new(Limits::default(), services);
|
||||||
|
|
||||||
|
|||||||
@@ -102,6 +102,7 @@ fn services_with(modules: Arc<dyn ModuleSource>) -> Services {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -233,6 +233,7 @@ fn make_engine() -> Arc<Engine> {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -42,6 +42,7 @@ fn engine_with(rec: Arc<RecordingEmail>) -> Arc<Engine> {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
rec,
|
rec,
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -170,6 +170,7 @@ fn make_engine() -> Arc<Engine> {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -93,6 +93,7 @@ fn engine_with(http: Arc<dyn HttpService>) -> Arc<Engine> {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -112,6 +112,7 @@ fn make_engine() -> Arc<Engine> {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -50,6 +50,7 @@ fn make_engine(svc: Arc<RecordingPubsub>) -> Arc<Engine> {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -103,6 +103,7 @@ fn make_engine() -> Arc<Engine> {
|
|||||||
Arc::new(InMemorySecrets::default()),
|
Arc::new(InMemorySecrets::default()),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -97,6 +97,7 @@ fn make_engine() -> Arc<Engine> {
|
|||||||
Arc::new(picloud_shared::NoopSecretsService),
|
Arc::new(picloud_shared::NoopSecretsService),
|
||||||
Arc::new(picloud_shared::NoopEmailService),
|
Arc::new(picloud_shared::NoopEmailService),
|
||||||
Arc::new(picloud_shared::NoopUsersService),
|
Arc::new(picloud_shared::NoopUsersService),
|
||||||
|
Arc::new(picloud_shared::NoopQueueService),
|
||||||
);
|
);
|
||||||
Arc::new(Engine::new(Limits::default(), services))
|
Arc::new(Engine::new(Limits::default(), services))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -89,6 +89,12 @@ pub enum Capability {
|
|||||||
/// (v1.1.5). Maps to `script:write` on API keys (a publish is a
|
/// (v1.1.5). Maps to `script:write` on API keys (a publish is a
|
||||||
/// write that fans out to subscribers). Granted to `editor`+.
|
/// write that fans out to subscribers). Granted to `editor`+.
|
||||||
AppPubsubPublish(AppId),
|
AppPubsubPublish(AppId),
|
||||||
|
/// Enqueue a message onto this app's queue from a script (v1.1.9).
|
||||||
|
/// Maps to `script:write` on API keys (an enqueue is a write that
|
||||||
|
/// fans out to the registered consumer). Granted to `editor`+.
|
||||||
|
/// `depth` / `depth_pending` are read-only inspection and don't gate
|
||||||
|
/// — scripts in the app can always see their own queue depths.
|
||||||
|
AppQueueEnqueue(AppId),
|
||||||
/// Read a decrypted secret from this app's secrets store (v1.1.7).
|
/// Read a decrypted secret from this app's secrets store (v1.1.7).
|
||||||
/// Same trust shape as KV/docs/files read — granted to `viewer`+,
|
/// Same trust shape as KV/docs/files read — granted to `viewer`+,
|
||||||
/// maps to `script:read` on API keys. Honors the seven-scope
|
/// maps to `script:read` on API keys. Honors the seven-scope
|
||||||
@@ -156,6 +162,7 @@ impl Capability {
|
|||||||
| Self::AppFilesRead(id)
|
| Self::AppFilesRead(id)
|
||||||
| Self::AppFilesWrite(id)
|
| Self::AppFilesWrite(id)
|
||||||
| Self::AppPubsubPublish(id)
|
| Self::AppPubsubPublish(id)
|
||||||
|
| Self::AppQueueEnqueue(id)
|
||||||
| Self::AppSecretsRead(id)
|
| Self::AppSecretsRead(id)
|
||||||
| Self::AppSecretsWrite(id)
|
| Self::AppSecretsWrite(id)
|
||||||
| Self::AppEmailSend(id)
|
| Self::AppEmailSend(id)
|
||||||
@@ -191,6 +198,7 @@ impl Capability {
|
|||||||
| Self::AppHttpRequest(_)
|
| Self::AppHttpRequest(_)
|
||||||
| Self::AppFilesWrite(_)
|
| Self::AppFilesWrite(_)
|
||||||
| Self::AppPubsubPublish(_)
|
| Self::AppPubsubPublish(_)
|
||||||
|
| Self::AppQueueEnqueue(_)
|
||||||
| Self::AppSecretsWrite(_)
|
| Self::AppSecretsWrite(_)
|
||||||
| Self::AppEmailSend(_)
|
| Self::AppEmailSend(_)
|
||||||
| Self::AppUsersWrite(_)
|
| Self::AppUsersWrite(_)
|
||||||
@@ -358,6 +366,7 @@ const fn role_satisfies(role: AppRole, cap: Capability) -> bool {
|
|||||||
| Capability::AppHttpRequest(_)
|
| Capability::AppHttpRequest(_)
|
||||||
| Capability::AppFilesWrite(_)
|
| Capability::AppFilesWrite(_)
|
||||||
| Capability::AppPubsubPublish(_)
|
| Capability::AppPubsubPublish(_)
|
||||||
|
| Capability::AppQueueEnqueue(_)
|
||||||
| Capability::AppSecretsWrite(_)
|
| Capability::AppSecretsWrite(_)
|
||||||
| Capability::AppEmailSend(_)
|
| Capability::AppEmailSend(_)
|
||||||
| Capability::AppUsersWrite(_)
|
| Capability::AppUsersWrite(_)
|
||||||
|
|||||||
@@ -56,6 +56,7 @@ pub mod principal_resolver;
|
|||||||
pub mod pubsub_repo;
|
pub mod pubsub_repo;
|
||||||
pub mod pubsub_service;
|
pub mod pubsub_service;
|
||||||
pub mod queue_repo;
|
pub mod queue_repo;
|
||||||
|
pub mod queue_service;
|
||||||
pub mod realtime_authority;
|
pub mod realtime_authority;
|
||||||
pub mod repo;
|
pub mod repo;
|
||||||
pub mod route_admin;
|
pub mod route_admin;
|
||||||
|
|||||||
325
crates/manager-core/src/queue_service.rs
Normal file
325
crates/manager-core/src/queue_service.rs
Normal file
@@ -0,0 +1,325 @@
|
|||||||
|
//! `QueueServiceImpl` — wires `QueueRepo` underneath the
|
||||||
|
//! `picloud_shared::QueueService` trait scripts see via the Rhai bridge.
|
||||||
|
//!
|
||||||
|
//! Mirrors the other stateful services: script-as-gate authz
|
||||||
|
//! (`AppQueueEnqueue`, skipped when `cx.principal` is `None`), with the
|
||||||
|
//! backend doing the actual write. No `ServiceEventEmitter` here —
|
||||||
|
//! enqueue is observable via `queue:receive` triggers (commit 6).
|
||||||
|
|
||||||
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
use async_trait::async_trait;
|
||||||
|
use chrono::Utc;
|
||||||
|
use picloud_shared::{EnqueueOpts, QueueError, QueueMessageId, QueueService, SdkCallCx};
|
||||||
|
|
||||||
|
use crate::authz::{self, AuthzRepo, Capability};
|
||||||
|
use crate::queue_repo::{NewQueueMessage, QueueRepo};
|
||||||
|
|
||||||
|
/// Production impl: authz gate → repo. Trivial wrapper.
|
||||||
|
pub struct QueueServiceImpl {
|
||||||
|
repo: Arc<dyn QueueRepo>,
|
||||||
|
authz: Arc<dyn AuthzRepo>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl QueueServiceImpl {
|
||||||
|
#[must_use]
|
||||||
|
pub fn new(repo: Arc<dyn QueueRepo>, authz: Arc<dyn AuthzRepo>) -> Self {
|
||||||
|
Self { repo, authz }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait]
|
||||||
|
impl QueueService for QueueServiceImpl {
|
||||||
|
async fn enqueue(
|
||||||
|
&self,
|
||||||
|
cx: &SdkCallCx,
|
||||||
|
queue_name: &str,
|
||||||
|
payload: serde_json::Value,
|
||||||
|
opts: EnqueueOpts,
|
||||||
|
) -> Result<QueueMessageId, QueueError> {
|
||||||
|
if queue_name.is_empty() {
|
||||||
|
return Err(QueueError::EmptyName);
|
||||||
|
}
|
||||||
|
let max_attempts = opts.max_attempts.unwrap_or(3);
|
||||||
|
if !(1..=20).contains(&max_attempts) {
|
||||||
|
return Err(QueueError::InvalidOpts(
|
||||||
|
"queue::enqueue: max_attempts must be in [1, 20]".into(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
if let Some(delay) = opts.delay_ms {
|
||||||
|
// 24h cap on delay; matches what a cron tick would do anyway.
|
||||||
|
if !(0..=86_400_000).contains(&delay) {
|
||||||
|
return Err(QueueError::InvalidOpts(
|
||||||
|
"queue::enqueue: delay_ms must be in [0, 86_400_000]".into(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Script-as-gate authz: anonymous public-HTTP scripts skip the
|
||||||
|
// check (cx.principal is None); authenticated callers must hold
|
||||||
|
// AppQueueEnqueue.
|
||||||
|
if let Some(principal) = cx.principal.as_ref() {
|
||||||
|
authz::require(
|
||||||
|
&*self.authz,
|
||||||
|
principal,
|
||||||
|
Capability::AppQueueEnqueue(cx.app_id),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.map_err(|_| QueueError::Rejected("forbidden".into()))?;
|
||||||
|
}
|
||||||
|
|
||||||
|
let deliver_after = opts
|
||||||
|
.delay_ms
|
||||||
|
.and_then(chrono::Duration::try_milliseconds)
|
||||||
|
.map(|d| Utc::now() + d);
|
||||||
|
|
||||||
|
// Forensic only — never used for authz at consume-time (the
|
||||||
|
// queue:receive trigger runs as the registering principal per
|
||||||
|
// design notes §4).
|
||||||
|
let enqueued_by_principal = cx.principal.as_ref().map(|p| p.user_id);
|
||||||
|
|
||||||
|
self.repo
|
||||||
|
.enqueue(NewQueueMessage {
|
||||||
|
app_id: cx.app_id,
|
||||||
|
queue_name: queue_name.to_string(),
|
||||||
|
payload,
|
||||||
|
deliver_after,
|
||||||
|
max_attempts,
|
||||||
|
enqueued_by_principal,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(|e| QueueError::Unavailable(e.to_string()))
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn depth(&self, cx: &SdkCallCx, queue_name: &str) -> Result<u64, QueueError> {
|
||||||
|
if queue_name.is_empty() {
|
||||||
|
return Err(QueueError::EmptyName);
|
||||||
|
}
|
||||||
|
self.repo
|
||||||
|
.depth(cx.app_id, queue_name)
|
||||||
|
.await
|
||||||
|
.map_err(|e| QueueError::Unavailable(e.to_string()))
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn depth_pending(&self, cx: &SdkCallCx, queue_name: &str) -> Result<u64, QueueError> {
|
||||||
|
if queue_name.is_empty() {
|
||||||
|
return Err(QueueError::EmptyName);
|
||||||
|
}
|
||||||
|
self.repo
|
||||||
|
.depth_pending(cx.app_id, queue_name)
|
||||||
|
.await
|
||||||
|
.map_err(|e| QueueError::Unavailable(e.to_string()))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use crate::authz::AuthzError;
|
||||||
|
use picloud_shared::{AppId, AppRole, ExecutionId, RequestId, ScriptId};
|
||||||
|
|
||||||
|
struct AlwaysAllowAuthz;
|
||||||
|
#[async_trait]
|
||||||
|
impl AuthzRepo for AlwaysAllowAuthz {
|
||||||
|
async fn membership(
|
||||||
|
&self,
|
||||||
|
_user_id: picloud_shared::UserId,
|
||||||
|
_app_id: AppId,
|
||||||
|
) -> Result<Option<AppRole>, AuthzError> {
|
||||||
|
Ok(Some(AppRole::Editor))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
struct CapturingRepo {
|
||||||
|
last: tokio::sync::Mutex<Option<NewQueueMessage>>,
|
||||||
|
}
|
||||||
|
#[async_trait]
|
||||||
|
impl QueueRepo for CapturingRepo {
|
||||||
|
async fn enqueue(
|
||||||
|
&self,
|
||||||
|
msg: NewQueueMessage,
|
||||||
|
) -> Result<QueueMessageId, crate::queue_repo::QueueRepoError> {
|
||||||
|
let id = QueueMessageId::new();
|
||||||
|
*self.last.lock().await = Some(msg);
|
||||||
|
Ok(id)
|
||||||
|
}
|
||||||
|
async fn claim(
|
||||||
|
&self,
|
||||||
|
_app_id: AppId,
|
||||||
|
_queue_name: &str,
|
||||||
|
) -> Result<Option<crate::queue_repo::ClaimedMessage>, crate::queue_repo::QueueRepoError>
|
||||||
|
{
|
||||||
|
Ok(None)
|
||||||
|
}
|
||||||
|
async fn ack(
|
||||||
|
&self,
|
||||||
|
_message_id: QueueMessageId,
|
||||||
|
_claim_token: uuid::Uuid,
|
||||||
|
) -> Result<bool, crate::queue_repo::QueueRepoError> {
|
||||||
|
Ok(true)
|
||||||
|
}
|
||||||
|
async fn nack(
|
||||||
|
&self,
|
||||||
|
_message_id: QueueMessageId,
|
||||||
|
_claim_token: uuid::Uuid,
|
||||||
|
_retry_delay: chrono::Duration,
|
||||||
|
) -> Result<bool, crate::queue_repo::QueueRepoError> {
|
||||||
|
Ok(true)
|
||||||
|
}
|
||||||
|
async fn reclaim_visibility_timeouts(
|
||||||
|
&self,
|
||||||
|
) -> Result<u64, crate::queue_repo::QueueRepoError> {
|
||||||
|
Ok(0)
|
||||||
|
}
|
||||||
|
async fn depth(
|
||||||
|
&self,
|
||||||
|
_app_id: AppId,
|
||||||
|
_queue_name: &str,
|
||||||
|
) -> Result<u64, crate::queue_repo::QueueRepoError> {
|
||||||
|
Ok(42)
|
||||||
|
}
|
||||||
|
async fn depth_pending(
|
||||||
|
&self,
|
||||||
|
_app_id: AppId,
|
||||||
|
_queue_name: &str,
|
||||||
|
) -> Result<u64, crate::queue_repo::QueueRepoError> {
|
||||||
|
Ok(7)
|
||||||
|
}
|
||||||
|
async fn list_for_app(
|
||||||
|
&self,
|
||||||
|
_app_id: AppId,
|
||||||
|
) -> Result<Vec<(String, crate::queue_repo::QueueStats)>, crate::queue_repo::QueueRepoError>
|
||||||
|
{
|
||||||
|
Ok(vec![])
|
||||||
|
}
|
||||||
|
async fn dead_letter(
|
||||||
|
&self,
|
||||||
|
_message_id: QueueMessageId,
|
||||||
|
_claim_token: uuid::Uuid,
|
||||||
|
_app_id: AppId,
|
||||||
|
_queue_name: &str,
|
||||||
|
_trigger_id: Option<picloud_shared::TriggerId>,
|
||||||
|
_script_id: Option<ScriptId>,
|
||||||
|
_attempt: u32,
|
||||||
|
_first_attempt_at: chrono::DateTime<chrono::Utc>,
|
||||||
|
_last_error: &str,
|
||||||
|
) -> Result<picloud_shared::DeadLetterId, crate::queue_repo::QueueRepoError> {
|
||||||
|
Ok(picloud_shared::DeadLetterId::new())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn anon_cx() -> SdkCallCx {
|
||||||
|
let exec = ExecutionId::new();
|
||||||
|
SdkCallCx {
|
||||||
|
app_id: AppId::new(),
|
||||||
|
script_id: ScriptId::new(),
|
||||||
|
principal: None,
|
||||||
|
execution_id: exec,
|
||||||
|
request_id: RequestId::new(),
|
||||||
|
trigger_depth: 0,
|
||||||
|
root_execution_id: exec,
|
||||||
|
is_dead_letter_handler: false,
|
||||||
|
event: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn empty_queue_name_rejects() {
|
||||||
|
let repo = Arc::new(CapturingRepo {
|
||||||
|
last: tokio::sync::Mutex::new(None),
|
||||||
|
});
|
||||||
|
let svc = QueueServiceImpl::new(repo.clone(), Arc::new(AlwaysAllowAuthz));
|
||||||
|
let cx = anon_cx();
|
||||||
|
let err = svc
|
||||||
|
.enqueue(&cx, "", serde_json::Value::Null, EnqueueOpts::default())
|
||||||
|
.await
|
||||||
|
.unwrap_err();
|
||||||
|
assert!(matches!(err, QueueError::EmptyName));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn invalid_max_attempts_rejects() {
|
||||||
|
let repo = Arc::new(CapturingRepo {
|
||||||
|
last: tokio::sync::Mutex::new(None),
|
||||||
|
});
|
||||||
|
let svc = QueueServiceImpl::new(repo, Arc::new(AlwaysAllowAuthz));
|
||||||
|
let cx = anon_cx();
|
||||||
|
let err = svc
|
||||||
|
.enqueue(
|
||||||
|
&cx,
|
||||||
|
"x",
|
||||||
|
serde_json::Value::Null,
|
||||||
|
EnqueueOpts {
|
||||||
|
delay_ms: None,
|
||||||
|
max_attempts: Some(0),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap_err();
|
||||||
|
assert!(matches!(err, QueueError::InvalidOpts(_)));
|
||||||
|
let err = svc
|
||||||
|
.enqueue(
|
||||||
|
&cx,
|
||||||
|
"x",
|
||||||
|
serde_json::Value::Null,
|
||||||
|
EnqueueOpts {
|
||||||
|
delay_ms: None,
|
||||||
|
max_attempts: Some(21),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap_err();
|
||||||
|
assert!(matches!(err, QueueError::InvalidOpts(_)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn invalid_delay_rejects() {
|
||||||
|
let repo = Arc::new(CapturingRepo {
|
||||||
|
last: tokio::sync::Mutex::new(None),
|
||||||
|
});
|
||||||
|
let svc = QueueServiceImpl::new(repo, Arc::new(AlwaysAllowAuthz));
|
||||||
|
let cx = anon_cx();
|
||||||
|
let err = svc
|
||||||
|
.enqueue(
|
||||||
|
&cx,
|
||||||
|
"x",
|
||||||
|
serde_json::Value::Null,
|
||||||
|
EnqueueOpts {
|
||||||
|
delay_ms: Some(-1),
|
||||||
|
max_attempts: None,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap_err();
|
||||||
|
assert!(matches!(err, QueueError::InvalidOpts(_)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn anon_principal_skips_authz_and_writes() {
|
||||||
|
let repo = Arc::new(CapturingRepo {
|
||||||
|
last: tokio::sync::Mutex::new(None),
|
||||||
|
});
|
||||||
|
let svc = QueueServiceImpl::new(repo.clone(), Arc::new(AlwaysAllowAuthz));
|
||||||
|
let cx = anon_cx();
|
||||||
|
let payload = serde_json::json!({ "x": 1 });
|
||||||
|
svc.enqueue(&cx, "jobs", payload.clone(), EnqueueOpts::default())
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let captured = repo.last.lock().await.clone().unwrap();
|
||||||
|
assert_eq!(captured.queue_name, "jobs");
|
||||||
|
assert_eq!(captured.payload, payload);
|
||||||
|
assert_eq!(captured.max_attempts, 3);
|
||||||
|
assert!(captured.deliver_after.is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn depth_passes_through() {
|
||||||
|
let repo = Arc::new(CapturingRepo {
|
||||||
|
last: tokio::sync::Mutex::new(None),
|
||||||
|
});
|
||||||
|
let svc = QueueServiceImpl::new(repo, Arc::new(AlwaysAllowAuthz));
|
||||||
|
let cx = anon_cx();
|
||||||
|
assert_eq!(svc.depth(&cx, "any").await.unwrap(), 42);
|
||||||
|
assert_eq!(svc.depth_pending(&cx, "any").await.unwrap(), 7);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -261,6 +261,19 @@ pub async fn build_app(
|
|||||||
app_secrets_repo.clone(),
|
app_secrets_repo.clone(),
|
||||||
users.clone(),
|
users.clone(),
|
||||||
));
|
));
|
||||||
|
// v1.1.9 durable per-app queues. Producers write to queue_messages
|
||||||
|
// via QueueService; the dispatcher's queue arm (commit 6) consumes
|
||||||
|
// claimed messages and fires the registered queue:receive trigger.
|
||||||
|
let queue_repo: Arc<dyn picloud_manager_core::queue_repo::QueueRepo> =
|
||||||
|
Arc::new(picloud_manager_core::queue_repo::PostgresQueueRepo::new(
|
||||||
|
pool.clone(),
|
||||||
|
));
|
||||||
|
let queue: Arc<dyn picloud_shared::QueueService> = Arc::new(
|
||||||
|
picloud_manager_core::queue_service::QueueServiceImpl::new(
|
||||||
|
queue_repo.clone(),
|
||||||
|
authz.clone(),
|
||||||
|
),
|
||||||
|
);
|
||||||
let services = Services::new(
|
let services = Services::new(
|
||||||
kv,
|
kv,
|
||||||
docs,
|
docs,
|
||||||
@@ -273,6 +286,7 @@ pub async fn build_app(
|
|||||||
secrets,
|
secrets,
|
||||||
email,
|
email,
|
||||||
users.clone(),
|
users.clone(),
|
||||||
|
queue,
|
||||||
);
|
);
|
||||||
let engine = Arc::new(Engine::new(Limits::default(), services));
|
let engine = Arc::new(Engine::new(Limits::default(), services));
|
||||||
|
|
||||||
|
|||||||
@@ -23,8 +23,8 @@ use crate::{
|
|||||||
DeadLetterService, DocsService, EmailService, FilesService, HttpService, KvService,
|
DeadLetterService, DocsService, EmailService, FilesService, HttpService, KvService,
|
||||||
ModuleSource, NoopDeadLetterService, NoopDocsService, NoopEmailService, NoopEventEmitter,
|
ModuleSource, NoopDeadLetterService, NoopDocsService, NoopEmailService, NoopEventEmitter,
|
||||||
NoopFilesService, NoopHttpService, NoopKvService, NoopModuleSource, NoopPubsubService,
|
NoopFilesService, NoopHttpService, NoopKvService, NoopModuleSource, NoopPubsubService,
|
||||||
NoopSecretsService, NoopUsersService, PubsubService, SecretsService, ServiceEventEmitter,
|
NoopQueueService, NoopSecretsService, NoopUsersService, PubsubService, QueueService,
|
||||||
UsersService,
|
SecretsService, ServiceEventEmitter, UsersService,
|
||||||
};
|
};
|
||||||
|
|
||||||
/// SDK service bundle. See module docs for the lifecycle and the v1.1.x
|
/// SDK service bundle. See module docs for the lifecycle and the v1.1.x
|
||||||
@@ -95,6 +95,13 @@ pub struct Services {
|
|||||||
/// session tokens in the picloud binary; `NoopUsersService` in
|
/// session tokens in the picloud binary; `NoopUsersService` in
|
||||||
/// tests that don't exercise users.
|
/// tests that don't exercise users.
|
||||||
pub users: Arc<dyn UsersService>,
|
pub users: Arc<dyn UsersService>,
|
||||||
|
|
||||||
|
/// Durable per-app named queues (v1.1.9). Scripts get
|
||||||
|
/// `queue::{enqueue,depth,depth_pending}`. Backed by Postgres in
|
||||||
|
/// the picloud binary; `NoopQueueService` in tests that don't
|
||||||
|
/// touch queues. Consumers register via `queue:receive` triggers
|
||||||
|
/// (one per `(app_id, queue_name)`).
|
||||||
|
pub queue: Arc<dyn QueueService>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Services {
|
impl Services {
|
||||||
@@ -115,6 +122,7 @@ impl Services {
|
|||||||
secrets: Arc<dyn SecretsService>,
|
secrets: Arc<dyn SecretsService>,
|
||||||
email: Arc<dyn EmailService>,
|
email: Arc<dyn EmailService>,
|
||||||
users: Arc<dyn UsersService>,
|
users: Arc<dyn UsersService>,
|
||||||
|
queue: Arc<dyn QueueService>,
|
||||||
) -> Self {
|
) -> Self {
|
||||||
Self {
|
Self {
|
||||||
kv,
|
kv,
|
||||||
@@ -128,6 +136,7 @@ impl Services {
|
|||||||
secrets,
|
secrets,
|
||||||
email,
|
email,
|
||||||
users,
|
users,
|
||||||
|
queue,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -150,6 +159,7 @@ impl Services {
|
|||||||
Arc::new(NoopSecretsService),
|
Arc::new(NoopSecretsService),
|
||||||
Arc::new(NoopEmailService),
|
Arc::new(NoopEmailService),
|
||||||
Arc::new(NoopUsersService),
|
Arc::new(NoopUsersService),
|
||||||
|
Arc::new(NoopQueueService),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user