Files
PiCloud/crates/picloud/src/lib.rs
MechaCat02 a432091191 feat(groups): tree repos, hierarchy-aware RBAC, admin API
Server-side foundation for Phase-2 groups (no group-owned resources yet):

Shared types:
- GroupId, Group; App gains group_id; AppRole::{precedence,max} for
  folding the highest effective role across the membership chain.

Repos:
- group_repo: tree CRUD with reparent (ancestor-walk cycle guard under a
  coarse instance-wide structural advisory lock; slug frozen; bumps
  structure_version) and delete=RESTRICT (refuses non-empty groups).
- group_members_repo: per-(user, group) role grants, mirroring app_members.

Hierarchy-aware authz (§5.3):
- AuthzRepo gains effective_app_role / effective_group_role (default to
  direct membership / none, so the ~18 existing test stubs are untouched);
  the Postgres impl resolves each via one depth-bounded recursive CTE that
  MAXes the app's own row with every ancestor group_members row.
- can(): the Member path now folds inherited group roles, so a group_admin
  on any ancestor is implicitly app_admin beneath it. New Capability
  variants InstanceCreateGroup / Group{Read,Write,Admin}; group caps carry
  no app_id (bound API keys can't manage groups). 8 new unit tests.

Admin API:
- groups_api: group CRUD + reparent (admin at both source and destination
  parent, §5.6) + per-group members, all capability-gated.
- apps: POST /apps takes an optional parent group (default root); app
  responses carry group_id; my_role now reflects the effective role.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-24 19:59:00 +02:00

727 lines
33 KiB
Rust

//! Library half of the picloud all-in-one. `main.rs` is a thin wrapper
//! that opens the pool, runs migrations, calls `build_app`, and binds
//! the listener. Tests use the same `build_app` against an
//! ephemeral test database.
use std::sync::Arc;
use std::time::Duration;
use axum::middleware::from_fn_with_state;
use axum::{routing::get, Json, Router};
use picloud_executor_core::{Engine, Limits};
use picloud_manager_core::{
admin_router, admins_router, api_keys_router, app_members_router, apply_router, apps_api,
apps_router, attach_principal_if_present, auth_router, compile_routes, dead_letters_router,
dev_emails_router, email_inbound_router, files_admin_router, groups_router, kv_admin_router,
migrations, require_authenticated, route_admin_router, secrets_router, topics_router,
triggers_router, AbandonedRepo, AdminPrincipalResolver, AdminSessionRepository, AdminState,
AdminUserRepository, AdminsState, ApiKeyRepository, ApiKeysState, AppDomainRepository,
AppMembersRepository, AppMembersState, AppRepository, ApplyService, AppsState, AuthState,
AuthzRepo, DeadLetterRepo, DeadLettersState, DevEmailState, Dispatcher, DocsServiceImpl,
EmailInboundState, EmailServiceImpl, FilesAdminState, FilesConfig, FilesServiceImpl,
FsFilesRepo, GroupMembersRepository, GroupRepository, GroupsState, HttpConfig, HttpServiceImpl,
InboundNonceDedup, KvAdminState, KvServiceImpl, OutboxEventEmitter, OutboxRepo,
PostgresAbandonedRepo, PostgresAdminSessionRepository, PostgresAdminUserRepository,
PostgresApiKeyRepository, PostgresAppDomainRepository, PostgresAppMembersRepository,
PostgresAppRepository, PostgresAppSecretsRepo, PostgresAppUserInvitationRepo,
PostgresAppUserPasswordResetRepo, PostgresAppUserRepository, PostgresAppUserRoleRepo,
PostgresAppUserSessionRepository, PostgresAppUserVerificationRepo, PostgresDeadLetterRepo,
PostgresDeadLetterService, PostgresDocsRepo, PostgresExecutionLogRepository,
PostgresExecutionLogSink, PostgresGroupMembersRepository, PostgresGroupRepository,
PostgresKvRepo, PostgresOutboxRepo, PostgresPubsubRepo, PostgresRouteRepository,
PostgresScriptRepository, PostgresSecretsRepo, PostgresTopicRepo, PostgresTriggerRepo,
PrincipalResolver, PubsubServiceImpl, RealtimeAuthorityImpl, RepoResolver, RouteAdminState,
RouteRepository, SandboxCeiling, ScriptRepository, SecretsConfig, SecretsServiceImpl,
SecretsState, SubscriberTokenConfig, TopicRepo, TopicsState, TriggerConfig, TriggerRepo,
TriggersState, UsersServiceConfig, UsersServiceImpl,
};
use picloud_orchestrator_core::realtime::DEFAULT_GC_INTERVAL_SECS;
use picloud_orchestrator_core::routing::{AppDomainTable, RouteTable};
use picloud_orchestrator_core::{
data_plane_router, realtime_router, spawn_realtime_gc, user_routes_router, DataPlaneState,
ExecutionGate, InProcessBroadcaster, InboxRegistry, LocalExecutorClient, RealtimeState,
};
use picloud_shared::{
DeadLetterService, DocsService, EmailService, ExecutionLogSink, FilesService, HttpService,
InboxResolver, KvService, MasterKey, OutboxWriter, PubsubService, RealtimeAuthority,
RealtimeBroadcaster, SecretsService, ServiceEventEmitter, Services, UsersService, API_VERSION,
PRODUCT_VERSION, SDK_VERSION, WIRE_VERSION,
};
use sqlx::postgres::PgPoolOptions;
use sqlx::PgPool;
use tower_http::trace::TraceLayer;
/// Default session TTL when `PICLOUD_SESSION_TTL_HOURS` isn't set.
const DEFAULT_SESSION_TTL_HOURS: u64 = 24;
/// Bundles the auth-related dependencies that both `build_app` and the
/// startup bootstrap need. Built once in `main.rs` from the shared pool.
pub struct AuthDeps {
pub users: Arc<dyn AdminUserRepository>,
pub sessions: Arc<dyn AdminSessionRepository>,
pub keys: Arc<dyn ApiKeyRepository>,
pub ttl: Duration,
}
impl AuthDeps {
/// Construct from a pool with the binary's standard defaults.
#[must_use]
pub fn from_pool(pool: PgPool) -> Self {
Self {
users: Arc::new(PostgresAdminUserRepository::new(pool.clone())),
sessions: Arc::new(PostgresAdminSessionRepository::new(pool.clone())),
keys: Arc::new(PostgresApiKeyRepository::new(pool)),
ttl: read_session_ttl(),
}
}
}
fn read_session_ttl() -> Duration {
let hours = std::env::var("PICLOUD_SESSION_TTL_HOURS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|h| *h > 0)
.unwrap_or(DEFAULT_SESSION_TTL_HOURS);
Duration::from_secs(hours * 3600)
}
/// Compose the manager + orchestrator routes on top of a shared
/// Postgres pool, returning an Axum router ready to be served.
///
/// All API routes live under `/api/v{API_VERSION}/...`. The dashboard
/// is mounted by Caddy at `/admin/*` (its base path). Anything else
/// falls through to the user-route table — user scripts can bind to
/// arbitrary paths (subject to the reserved-prefix list).
///
/// `auth` carries the admin user/session repositories and the
/// configured session TTL. The manager-side admin endpoints
/// (`/api/v1/admin/scripts/*`, `/api/v1/admin/routes/*`,
/// `/api/v1/admin/admins/*`, `/api/v1/admin/auth/me`) are guarded by
/// the `require_admin` middleware. The data plane
/// (`/api/v1/execute/{id}`, the user-route fallthrough, `/healthz`,
/// `/version`) stays open — it's the public ingress for user scripts.
#[allow(clippy::too_many_lines)]
pub async fn build_app(
pool: PgPool,
auth: AuthDeps,
master_key: MasterKey,
) -> anyhow::Result<Router> {
let script_repo = Arc::new(PostgresScriptRepository::new(pool.clone()));
let log_repo = Arc::new(PostgresExecutionLogRepository::new(pool.clone()));
let log_sink: Arc<dyn ExecutionLogSink> = Arc::new(PostgresExecutionLogSink::new(pool.clone()));
let route_repo = Arc::new(PostgresRouteRepository::new(pool.clone()));
let apps_repo: Arc<dyn AppRepository> = Arc::new(PostgresAppRepository::new(pool.clone()));
let groups_repo: Arc<dyn GroupRepository> =
Arc::new(PostgresGroupRepository::new(pool.clone()));
let group_members_repo: Arc<dyn GroupMembersRepository> =
Arc::new(PostgresGroupMembersRepository::new(pool.clone()));
let domains_repo: Arc<dyn AppDomainRepository> =
Arc::new(PostgresAppDomainRepository::new(pool.clone()));
// The Postgres app_members repo implements both `AppMembersRepository`
// (CRUD over the table) and `AuthzRepo` (single-row membership lookup
// for capability checks). Construct it once and clone the Arc into
// both trait views — same allocation, two vtables.
let members_concrete = Arc::new(PostgresAppMembersRepository::new(pool.clone()));
let members: Arc<dyn AppMembersRepository> = members_concrete.clone();
let authz: Arc<dyn AuthzRepo> = members_concrete;
// Triggers framework storage. The outbox event emitter routes
// KV mutations into the outbox; the dispatcher fans them out.
let trigger_repo: Arc<dyn TriggerRepo> = Arc::new(PostgresTriggerRepo::new(pool.clone()));
// PostgresOutboxRepo implements both `OutboxRepo` (the dispatcher
// surface) and `OutboxWriter` (the orchestrator surface). Construct
// the concrete Arc once, clone it into each trait view — same
// allocation, two vtables (mirrors how `members_concrete` above is
// used as both `AppMembersRepository` and `AuthzRepo`).
let outbox_concrete = Arc::new(PostgresOutboxRepo::new(pool.clone()));
let outbox_repo: Arc<dyn OutboxRepo> = outbox_concrete.clone();
let outbox_writer: Arc<dyn OutboxWriter> = outbox_concrete;
let dl_repo: Arc<dyn DeadLetterRepo> = Arc::new(PostgresDeadLetterRepo::new(pool.clone()));
let abandoned_repo: Arc<dyn AbandonedRepo> = Arc::new(PostgresAbandonedRepo::new(pool.clone()));
let trigger_config = TriggerConfig::from_env();
// SDK services bundle. v1.1.1 added KV + dead-letter; v1.1.2 added
// the docs store; v1.1.3 adds the module source backing the Rhai
// resolver. All bound services share the outbox-backed event
// emitter so KV and docs mutations both fan out through the same
// dispatcher.
let kv_repo = Arc::new(PostgresKvRepo::new(pool.clone()));
let docs_repo = Arc::new(PostgresDocsRepo::new(pool.clone()));
let events: Arc<dyn ServiceEventEmitter> = Arc::new(OutboxEventEmitter::new(
trigger_repo.clone(),
outbox_repo.clone(),
));
let kv: Arc<dyn KvService> = Arc::new(KvServiceImpl::with_max_value_bytes(
kv_repo.clone(),
authz.clone(),
events.clone(),
picloud_manager_core::kv_service::kv_max_value_bytes_from_env(),
));
let docs: Arc<dyn DocsService> = Arc::new(DocsServiceImpl::with_max_value_bytes(
docs_repo,
authz.clone(),
events.clone(),
picloud_manager_core::docs_service::docs_max_value_bytes_from_env(),
));
let dl_service: Arc<dyn DeadLetterService> = Arc::new(PostgresDeadLetterService::new(
dl_repo.clone(),
outbox_repo.clone(),
authz.clone(),
));
let modules: Arc<dyn picloud_shared::ModuleSource> = Arc::new(
picloud_manager_core::PostgresModuleSource::new(pool.clone()),
);
// v1.1.4 outbound HTTP. The reqwest client is built once here with
// the SSRF deny-list resolver. `PICLOUD_HTTP_ALLOW_PRIVATE=true`
// disables the deny-list entirely — dev/test only, so warn loudly.
let http_config = HttpConfig::from_env();
if http_config.allow_private {
tracing::warn!(
"PICLOUD_HTTP_ALLOW_PRIVATE is set — the outbound-HTTP SSRF deny-list is DISABLED. \
Scripts can reach loopback/private/link-local addresses. Do NOT use in production."
);
}
let http: Arc<dyn HttpService> = Arc::new(HttpServiceImpl::new(http_config, authz.clone()));
// v1.1.5 filesystem-backed blob storage. Metadata lives in Postgres;
// the bytes live on disk under `PICLOUD_FILES_ROOT` (default ./data).
let files_config = FilesConfig::from_env();
let files_max_size = files_config.max_file_size_bytes;
// Kept for the v1.1.6 orphan sweeper (cleans stale `*.tmp.*` files).
let files_root = files_config.root.clone();
let files_repo = Arc::new(FsFilesRepo::new(pool.clone(), files_config));
let files: Arc<dyn FilesService> = Arc::new(FilesServiceImpl::new(
files_repo.clone(),
authz.clone(),
events.clone(),
files_max_size,
));
// v1.1.6 realtime: the in-process broadcaster is shared between the
// publish path (PubsubServiceImpl fans out to SSE subscribers after
// the durable outbox fan-out) and the SSE endpoint (subscribe side).
// The topic registry + app-secrets repo back the subscriber-token
// mint + SSE subscribe-authorization.
let broadcaster_concrete = Arc::new(InProcessBroadcaster::from_env());
let broadcaster: Arc<dyn RealtimeBroadcaster> = broadcaster_concrete.clone();
let topic_repo: Arc<dyn TopicRepo> = Arc::new(PostgresTopicRepo::new(pool.clone()));
let app_secrets_repo = Arc::new(PostgresAppSecretsRepo::new(
pool.clone(),
master_key.clone(),
));
// v1.1.7's migrate_plaintext_keys startup sweep was deleted in
// v1.1.8 F1 alongside migration 0032, which drops the plaintext
// column. Operators upgrading from v1.1.6 or earlier MUST apply
// v1.1.7 first; the 0032 migration's guard refuses to run
// otherwise.
// v1.1.5 durable pub/sub, extended in v1.1.6 with the realtime
// broadcast + subscriber-token mint. Publishes fan out to matching
// pubsub triggers at publish time (one outbox row each, delivered by
// the same dispatcher as every other async trigger) AND, best-effort,
// to in-process SSE subscribers.
let pubsub_repo = Arc::new(PostgresPubsubRepo::new(pool.clone()));
let pubsub: Arc<dyn PubsubService> = Arc::new(
PubsubServiceImpl::new(pubsub_repo, authz.clone())
.with_max_message_bytes(
picloud_manager_core::pubsub_service::pubsub_max_message_bytes_from_env(),
)
.with_realtime(
broadcaster.clone(),
topic_repo.clone(),
app_secrets_repo.clone(),
SubscriberTokenConfig::from_env(),
),
);
// v1.1.7 encrypted per-app secrets. Values are AES-256-GCM-sealed
// with the process master key before they touch Postgres; the repo
// only ever sees ciphertext + nonce. The admin surface reuses the
// same repo + master key (see `secrets_state` below).
let secrets_config = SecretsConfig::from_env();
let secrets_max_value_bytes = secrets_config.max_value_bytes;
let secrets_repo: Arc<dyn picloud_manager_core::SecretsRepo> =
Arc::new(PostgresSecretsRepo::new(pool.clone()));
let secrets: Arc<dyn SecretsService> = Arc::new(SecretsServiceImpl::new(
secrets_repo.clone(),
authz.clone(),
master_key.clone(),
secrets_config,
));
// v1.1.7 outbound email. Builds a lettre SMTP transport from
// PICLOUD_SMTP_* env (disabled mode + warning if unconfigured). G5: in
// dev mode with no relay, captures mail in memory instead of erroring;
// `dev_email_sink` is `Some` then, and we mount the dev inspection
// endpoint below.
let (email_impl, dev_email_sink) = EmailServiceImpl::from_env_with_dev_capture(authz.clone());
let email: Arc<dyn EmailService> = Arc::new(email_impl);
// v1.1.8 data-plane user management. Wires Argon2id-hashed user
// rows + SHA-256-hashed sliding-window sessions to the Rhai
// `users::*` namespace and the admin /apps/{id}/users HTTP surface.
let app_users_repo = Arc::new(PostgresAppUserRepository::new(pool.clone()));
let app_user_sessions_repo = Arc::new(PostgresAppUserSessionRepository::new(pool.clone()));
let app_user_verifications_repo = Arc::new(PostgresAppUserVerificationRepo::new(pool.clone()));
let app_user_password_resets_repo =
Arc::new(PostgresAppUserPasswordResetRepo::new(pool.clone()));
let app_user_invitations_repo = Arc::new(PostgresAppUserInvitationRepo::new(pool.clone()));
let app_user_roles_repo = Arc::new(PostgresAppUserRoleRepo::new(pool.clone()));
let users: Arc<dyn UsersService> = Arc::new(UsersServiceImpl::new(
app_users_repo.clone(),
app_user_sessions_repo.clone(),
app_user_verifications_repo.clone(),
app_user_password_resets_repo.clone(),
app_user_invitations_repo.clone(),
app_user_roles_repo.clone(),
email.clone(),
authz.clone(),
events.clone(),
UsersServiceConfig::from_env(),
));
// v1.1.8 F3: build the realtime authority now that UsersService
// is constructed; the Session arm of authorize_subscribe calls
// users.verify_session_for_realtime.
let realtime_authority: Arc<dyn RealtimeAuthority> = Arc::new(RealtimeAuthorityImpl::new(
topic_repo.clone(),
app_secrets_repo.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::with_max_payload_bytes(
queue_repo.clone(),
authz.clone(),
picloud_manager_core::queue_service::queue_max_payload_bytes_from_env(),
),
);
// Route table created early (before Services) so InvokeServiceImpl
// can use it for path resolution. It's populated below from the
// route_repo, then re-populated whenever the admin layer writes
// routes.
let route_table = Arc::new(RouteTable::new());
let initial = route_repo.list_all().await?;
// Lenient: a single un-compilable stored route (e.g. one whose path
// became reserved under stricter validation) is skipped-with-warning
// inside compile_routes, never aborting boot. (H1)
let compiled = compile_routes(&initial);
route_table.replace_all(compiled);
// v1.1.9 function composition. InvokeServiceImpl resolves targets
// by Id, Name (via get_by_name), or Path (via the orchestrator's
// RouteTable). invoke_async writes an OutboxSourceKind::Invoke row
// that the dispatcher fires through the standard executor path.
let invoke: Arc<dyn picloud_shared::InvokeService> = Arc::new(
picloud_manager_core::invoke_service::InvokeServiceImpl::new(
script_repo.clone(),
route_table.clone(),
outbox_repo.clone(),
)
// F-S-012: gate authenticated invoke() callers on AppInvoke so
// an editor-role principal can't trigger admin-only worker
// scripts in the same app. Anonymous callers skip the check
// under script-as-gate semantics.
.with_authz(authz.clone()),
);
let services = Services::new(
kv,
docs,
dl_service.clone(),
events,
modules,
http,
files,
pubsub,
secrets,
email,
users.clone(),
queue,
invoke,
);
// v1.1.9: keep the invoke depth bound aligned with the dispatcher's
// trigger-depth bound (same counter under the hood).
let engine_limits = Limits {
trigger_depth_max: trigger_config.max_trigger_depth,
..Limits::default()
};
let engine = Arc::new(Engine::new(engine_limits, services));
// v1.1.9: install the back-reference the `invoke` SDK bridge needs
// for synchronous re-entry. Weak so the Arc-cycle stays loose;
// OnceLock-backed so it's idempotent.
engine.set_self_weak(Arc::downgrade(&engine));
// Same shape for app domains (Host → app_id cache).
let app_domain_table = Arc::new(AppDomainTable::new());
let initial_domains = domains_repo.list_all().await?;
let compiled_domains: Vec<_> = initial_domains
.iter()
.filter_map(|d| {
picloud_orchestrator_core::routing::parse_app_domain(&d.pattern)
.ok()
.map(|p| picloud_orchestrator_core::routing::CompiledAppDomain {
app_id: d.app_id,
pattern: p.pattern,
shape_key: p.shape_key,
})
})
.collect();
app_domain_table.replace(compiled_domains);
let resolver = Arc::new(RepoResolver::new(script_repo.clone()));
// Single global gate — overflow is rejected with 503 + Retry-After.
// See `ExecutionGate` docs and `PICLOUD_MAX_CONCURRENT_EXECUTIONS`.
let gate = Arc::new(ExecutionGate::from_env());
let executor = Arc::new(LocalExecutorClient::new(engine.clone(), gate.clone()));
// Dispatcher — single tokio task that polls the outbox and routes
// due rows to the executor. Shares the `ExecutionGate` with sync
// HTTP per design notes §2 (one cap for everything).
let dispatcher_script_repo: Arc<dyn ScriptRepository> = script_repo.clone();
let principals: Arc<dyn PrincipalResolver> =
Arc::new(AdminPrincipalResolver::new(auth.users.clone()));
// The InboxRegistry is constructed once and shared between the
// orchestrator (registers receivers, awaits) and the dispatcher
// (delivers results). Two Arc views on the same allocation.
let inbox_registry = Arc::new(InboxRegistry::new());
let inbox_resolver: Arc<dyn InboxResolver> = inbox_registry.clone();
Dispatcher {
outbox: outbox_repo.clone(),
triggers: trigger_repo.clone(),
scripts: dispatcher_script_repo,
dead_letters: dl_repo.clone(),
abandoned: abandoned_repo.clone(),
principals,
executor: executor.clone(),
gate,
log_sink: log_sink.clone(),
inbox: inbox_resolver,
queue: queue_repo.clone(),
config: trigger_config,
instance_id: format!("picloud-{}", std::process::id()),
}
.spawn();
let admin = AdminState {
repo: script_repo.clone(),
logs: log_repo,
apps: apps_repo.clone(),
authz: authz.clone(),
validator: engine.clone(),
sandbox_ceiling: SandboxCeiling::from_env(),
};
let route_admin = RouteAdminState {
routes: route_repo.clone(),
scripts: script_repo.clone(),
domains: domains_repo.clone(),
table: route_table.clone(),
authz: authz.clone(),
};
let data_plane = DataPlaneState {
executor,
resolver,
log_sink,
app_domains: app_domain_table.clone(),
routes: route_table.clone(),
inbox: inbox_registry,
outbox: outbox_writer,
};
// Weekly retention sweepers for dead_letters + abandoned_executions.
// Defaults: 30 days / 7 days (design notes §3 #9 + §4 retention).
picloud_manager_core::spawn_dead_letter_gc(
dl_repo.clone(),
trigger_config.dead_letter_retention_days,
);
picloud_manager_core::spawn_abandoned_gc(
abandoned_repo.clone(),
trigger_config.abandoned_retention_days,
);
// v1.1.8 weekly sweep for app-user sessions + tokens. No
// configurable retention — rows are pruned as soon as they're
// either expired (by TTL) or consumed/revoked.
picloud_manager_core::spawn_app_user_token_gc(
app_user_sessions_repo.clone(),
app_user_verifications_repo.clone(),
app_user_password_resets_repo.clone(),
app_user_invitations_repo.clone(),
);
// v1.1.4: cron scheduler. Polls cron_trigger_details on a tick and
// enqueues due triggers into the outbox; the dispatcher above
// delivers them like any other async trigger.
picloud_manager_core::spawn_cron_scheduler(pool.clone(), trigger_config.cron_tick_interval_ms);
// v1.1.6: GC empty realtime broadcast channels (one-shot subscribers)
// and sweep orphaned `*.tmp.*` blobs left by crashed file writes.
spawn_realtime_gc(broadcaster_concrete, DEFAULT_GC_INTERVAL_SECS);
picloud_manager_core::spawn_files_orphan_sweep(files_root);
let triggers_state = TriggersState {
triggers: trigger_repo.clone(),
apps: apps_repo.clone(),
authz: authz.clone(),
scripts: script_repo.clone(),
config: trigger_config,
master_key: master_key.clone(),
};
// Declarative reconcile engine (pic plan / apply). Trait-object repos
// for the read/diff path; shares the same handles as the CRUD routers.
let apply_service = ApplyService {
pool: pool.clone(),
scripts: script_repo.clone(),
routes: route_repo.clone(),
triggers: trigger_repo.clone(),
secrets: secrets_repo.clone(),
apps: apps_repo.clone(),
domains: domains_repo.clone(),
authz: authz.clone(),
validator: engine.clone(),
sandbox_ceiling: SandboxCeiling::from_env(),
trigger_config,
route_table: route_table.clone(),
master_key: master_key.clone(),
};
// v1.1.9: keep a clone for the queues-api state (built later).
let trigger_repo_for_queues = trigger_repo.clone();
// v1.1.7 public inbound-email receiver. Outside the admin auth layer
// (the URL + per-trigger HMAC secret are the security boundary).
let email_inbound_state = EmailInboundState {
triggers: trigger_repo,
outbox: outbox_repo.clone(),
master_key: master_key.clone(),
nonce_dedup: Arc::new(InboundNonceDedup::new()),
bad_sig_limiter: Arc::new(
picloud_manager_core::email_inbound_api::BadSignatureLimiter::new(),
),
};
let dead_letters_state = DeadLettersState {
repo: dl_repo,
service: dl_service,
apps: apps_repo.clone(),
authz: authz.clone(),
};
let files_admin_state = FilesAdminState {
files: files_repo,
apps: apps_repo.clone(),
authz: authz.clone(),
};
let topics_state = TopicsState {
topics: topic_repo,
apps: apps_repo.clone(),
authz: authz.clone(),
broadcaster: broadcaster.clone(),
};
let secrets_state = SecretsState {
repo: secrets_repo,
apps: apps_repo.clone(),
authz: authz.clone(),
master_key,
max_value_bytes: secrets_max_value_bytes,
};
let apps_state = AppsState {
apps: apps_repo.clone(),
domains: domains_repo,
routes: route_repo,
domain_table: app_domain_table.clone(),
authz: authz.clone(),
groups: groups_repo.clone(),
};
// Audit 2026-06-11 H-B1 — login DoS defenses. PICLOUD_LOGIN_ARGON2_PARALLELISM
// overrides the per-process Argon2-during-login concurrency cap.
let argon2_login_parallelism = std::env::var("PICLOUD_LOGIN_ARGON2_PARALLELISM")
.ok()
.and_then(|s| s.parse::<usize>().ok())
.filter(|n| *n > 0)
.unwrap_or(2);
// Audit 2026-06-11 (PrincipalCache revocation-lag) — one shared
// cache instance threaded into every state that performs auth or
// revocation, so an evict on the revocation side is visible to the
// middleware's resolve side. A per-state cache would let a revoked
// principal keep authenticating from the middleware's own copy.
let principal_cache = Arc::new(picloud_manager_core::auth_middleware::PrincipalCache::new());
let auth_state = AuthState {
users: auth.users.clone(),
sessions: auth.sessions.clone(),
keys: auth.keys.clone(),
ttl: auth.ttl,
principal_cache: principal_cache.clone(),
login_rate_limiter: Arc::new(
picloud_manager_core::login_rate_limit::LoginRateLimiter::new(),
),
argon2_login_semaphore: Arc::new(tokio::sync::Semaphore::new(argon2_login_parallelism)),
};
let admins_state = AdminsState {
users: auth.users.clone(),
sessions: auth.sessions,
keys: auth.keys.clone(),
authz: authz.clone(),
principal_cache: principal_cache.clone(),
};
let app_members_state = AppMembersState {
apps: apps_state.apps.clone(),
users: auth.users.clone(),
members,
authz: authz.clone(),
};
let groups_state = GroupsState {
groups: groups_repo.clone(),
group_members: group_members_repo.clone(),
apps: apps_state.apps.clone(),
users: auth.users.clone(),
authz: authz.clone(),
};
let app_users_admin_state = picloud_manager_core::AppUsersState {
apps: apps_state.apps.clone(),
authz: authz.clone(),
users: users.clone(),
};
let api_keys_state = ApiKeysState {
keys: auth.keys,
principal_cache: principal_cache.clone(),
};
// /admin/auth/login + /logout are unguarded by design (login is how
// you get in). /admin/auth/me applies the middleware internally so
// the same Router::with_state machinery composes cleanly. Everything
// else under /admin gets the require_authenticated layer; capability
// checks live in each handler (after the resource is loaded so the
// capability binds to the resource's actual app_id).
let mut guarded_admin = Router::new()
.merge(admin_router(admin))
.merge(route_admin_router(route_admin))
.merge(admins_router(admins_state))
.merge(apps_router(apps_state))
.merge(app_members_router(app_members_state))
.merge(groups_router(groups_state))
.merge(picloud_manager_core::app_users_router(
app_users_admin_state,
))
.merge(api_keys_router(api_keys_state))
.merge(triggers_router(triggers_state))
.merge(apply_router(apply_service))
.merge(picloud_manager_core::queues_api::queues_router(
picloud_manager_core::queues_api::QueuesState {
queues: queue_repo.clone(),
triggers: trigger_repo_for_queues.clone(),
apps: apps_repo.clone(),
authz: authz.clone(),
scripts: script_repo.clone(),
},
))
.merge(files_admin_router(files_admin_state))
.merge(kv_admin_router(KvAdminState {
kv: kv_repo.clone(),
apps: apps_repo.clone(),
authz: authz.clone(),
}))
.merge(topics_router(topics_state))
.merge(secrets_router(secrets_state))
.merge(dead_letters_router(dead_letters_state));
// G5: dev-only mail inspection — mounted exactly when the email
// service is capturing in memory (dev mode + no relay). Same
// `require_authenticated` layer as the rest of /admin applies below.
if let Some(sink) = dev_email_sink {
guarded_admin = guarded_admin.merge(dev_emails_router(DevEmailState { sink }));
}
let guarded_admin = guarded_admin.layer(from_fn_with_state(
auth_state.clone(),
require_authenticated,
));
// Silence "unused import" lint on `apps_api` — we re-export via the
// facade above; the bare module path is retained so it's discoverable.
let _ = apps_api::AppsState::clone;
// Opportunistic principal extraction on every data-plane request.
// Always inserts `Extension<Option<Principal>>`: Some for authed
// ingress (bearer / cookie), None otherwise. Handlers depend on
// this layer being applied — scoped to the data-plane routers so
// the admin path (which uses `require_authenticated`) doesn't
// double-resolve the same token.
let data_plane_routed = data_plane_router(data_plane.clone()).layer(from_fn_with_state(
auth_state.clone(),
attach_principal_if_present,
));
let user_routes = user_routes_router(data_plane).layer(from_fn_with_state(
auth_state.clone(),
attach_principal_if_present,
));
let api_v1 = Router::new()
.nest("/admin", auth_router(auth_state))
.nest("/admin", guarded_admin)
.merge(email_inbound_router(email_inbound_state))
.merge(data_plane_routed);
// v1.1.6 SSE realtime surface, merged at the root (deliberately NOT
// under /api/ — realtime is its own versioning surface). Public auth
// is per-topic; no principal middleware (token verification is the
// gate, handled inside the authority).
let realtime = realtime_router(RealtimeState::new(
app_domain_table,
broadcaster,
realtime_authority,
));
Ok(Router::new()
.route("/healthz", get(healthz))
.route("/version", get(version))
.nest(&format!("/api/v{API_VERSION}"), api_v1)
.merge(realtime)
.merge(user_routes)
.layer(TraceLayer::new_for_http()))
}
/// Default Postgres pool size. Scaled to match the execution gate
/// default (32 concurrent script executions) so a hot pool doesn't
/// starve the data plane. Override with `PICLOUD_DB_MAX_CONNECTIONS`.
/// Sizing relationship: each script run does multiple sequential DB
/// calls (script resolve, SDK calls, log sink, outbox emit) plus the
/// dispatcher tick / cron tick / GC sweeps / auth middleware all draw
/// from the same pool, so this is a floor not a ceiling.
pub const DEFAULT_DB_MAX_CONNECTIONS: u32 = 32;
/// Open a Postgres pool with the binary's standard timeout settings.
/// Exposed so tests reach for the same configuration when needed.
pub async fn init_db(url: &str) -> anyhow::Result<PgPool> {
let max_connections = std::env::var("PICLOUD_DB_MAX_CONNECTIONS")
.ok()
.and_then(|v| v.trim().parse::<u32>().ok())
.filter(|n| *n > 0)
.unwrap_or(DEFAULT_DB_MAX_CONNECTIONS);
let pool = PgPoolOptions::new()
.max_connections(max_connections)
.acquire_timeout(Duration::from_secs(5))
.connect(url)
.await?;
Ok(pool)
}
async fn healthz() -> &'static str {
"ok"
}
/// Snapshot of every compatibility-surface version this process speaks
/// plus the operator-configured public base URL (so the dashboard can
/// render full URLs for user routes).
///
/// Source of truth: `shared::version`, the embedded migrations, and
/// the `PICLOUD_PUBLIC_BASE_URL` env var (default
/// `http://localhost:8000`).
async fn version() -> Json<serde_json::Value> {
let public_base_url = std::env::var("PICLOUD_PUBLIC_BASE_URL")
.unwrap_or_else(|_| "http://localhost:8000".to_string());
Json(serde_json::json!({
"product": PRODUCT_VERSION,
"sdk": SDK_VERSION,
"api": API_VERSION,
"schema": migrations::latest_version(),
"wire": WIRE_VERSION,
"public_base_url": public_base_url,
}))
}
// F-Q-011: the PostgresScriptRepoHandle newtype lived here so the
// resolver could take an owned `impl ScriptRepository` from the shared
// Arc<PostgresScriptRepository>. The blanket
// `impl<T: ScriptRepository + ?Sized> ScriptRepository for Arc<T>` in
// manager-core::repo lets us pass `script_repo.clone()` directly,
// eliminating 70 lines of hand-delegated boilerplate that had to be
// rewritten every time the trait gained a method.