From 2fc9476f9efbdc2292df4c8dbd69959a4e5c281f Mon Sep 17 00:00:00 2001 From: MechaCat02 Date: Mon, 13 Jul 2026 20:44:50 +0200 Subject: [PATCH] =?UTF-8?q?feat(interceptors):=20=C2=A79.4=20service=20int?= =?UTF-8?q?erceptors=20=E2=80=94=20thin=20KV=20allow/deny=20slice?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Smallest honest vertical slice of §9.4: a `[[interceptors]]` block (app OR group) binds a script to run BEFORE `kv::set`/`delete`; it reads the operation context (`ctx.request.body`: service, action, collection, key, value, caller ids) and returns `#{ allowed, reason }` — `allowed == false` denies the op (the write never runs, the caller gets a runtime error). Reuses two existing mechanisms rather than inventing new ones: - Registration mirrors extension points (§5.5): a marker table `0073_interceptors.sql` (owner-polymorphic app_id/group_id XOR, keyed (service, op) → script), `interceptor_repo` (insert/delete/list + the nearest-owner-wins `resolve_before` chain walk), reconciled through the declarative apply exactly like `vars` (create/update/delete, prunable). - Execution reuses the `invoke()` re-entry path: the new `InterceptorService` (shared trait + Postgres-backed impl) only RESOLVES the script name (keeping executor-core Postgres-free); the executor's `sdk::interceptor::run_before` resolves that name and runs it via `run_resolved_blocking` (extracted from `invoke_blocking` — shared depth bound + AST cache). An un-hooked write pays one indexed `Ok(None)` resolve; no interceptor ⇒ zero overhead. Nearest-owner-wins so an app overrides a group's interceptor, and a group interceptor is inherited by every descendant app — the chain walk is the isolation boundary (a sibling subtree never matches). `validate_bundle_for` restricts the MVP to `service = "kv"`, `op ∈ {set, delete}`, one marker per (service, op). Deferred (documented in §9.4): the `data` transform return, services other than kv, `after_*` hooks, chaining + circular-dependency guard, the timeout policy, and a `pic interceptors ls` read surface (needs a server route). Pinned by `tests/interceptors.rs` (deny blocks the write; allow passes; group→app inheritance), schema snapshot re-blessed. 154/154 journeys pass. Co-Authored-By: Claude Opus 4.8 (1M context) --- CLAUDE.md | 4 +- crates/executor-core/src/sdk/interceptor.rs | 111 ++++++++++ crates/executor-core/src/sdk/invoke.rs | 51 ++--- crates/executor-core/src/sdk/kv.rs | 60 +++++- crates/executor-core/src/sdk/mod.rs | 3 +- .../migrations/0073_interceptors.sql | 42 ++++ crates/manager-core/src/apply_service.rs | 170 +++++++++++++++ crates/manager-core/src/interceptor_repo.rs | 184 +++++++++++++++++ .../manager-core/src/interceptor_service.rs | 37 ++++ crates/manager-core/src/lib.rs | 2 + crates/manager-core/tests/expected_schema.txt | 24 +++ crates/picloud-cli/src/client.rs | 10 + crates/picloud-cli/src/cmds/apply.rs | 18 ++ crates/picloud-cli/src/cmds/init.rs | 1 + crates/picloud-cli/src/cmds/plan.rs | 17 +- crates/picloud-cli/src/cmds/pull.rs | 3 + crates/picloud-cli/src/manifest.rs | 25 +++ crates/picloud-cli/tests/cli.rs | 1 + crates/picloud-cli/tests/interceptors.rs | 193 ++++++++++++++++++ crates/picloud/src/lib.rs | 8 +- crates/shared/src/interceptor.rs | 51 +++++ crates/shared/src/lib.rs | 2 + crates/shared/src/services.rs | 27 ++- serverless_cloud_blueprint.md | 16 +- 24 files changed, 1015 insertions(+), 45 deletions(-) create mode 100644 crates/executor-core/src/sdk/interceptor.rs create mode 100644 crates/manager-core/migrations/0073_interceptors.sql create mode 100644 crates/manager-core/src/interceptor_repo.rs create mode 100644 crates/manager-core/src/interceptor_service.rs create mode 100644 crates/picloud-cli/tests/interceptors.rs create mode 100644 crates/shared/src/interceptor.rs diff --git a/CLAUDE.md b/CLAUDE.md index 19da0ae..b047ab7 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -12,7 +12,7 @@ Authoritative design: [serverless_cloud_blueprint.md](serverless_cloud_blueprint **Current focus: v1.2 _Hierarchies_ — groups + the declarative project tool** ([docs/design/groups-and-project-tool.md](docs/design/groups-and-project-tool.md)). That doc's §11 uses its own **Phase 1–6 numbering, distinct from the blueprint product-phase numbering above — do not conflate them** (its "Phase 3" = group-inherited config, not admin auth). Implemented on `feat/groups-*` branches: §11 Phase 1 (declarative `pic plan`/`apply`/`prune` + env overlays), Phase 2 (single-parent groups tree + hierarchy-aware RBAC), Phase 3 (group-inherited, env-scoped `vars` + secrets resolved **live** via a recursive CTE — no materialized cache), Phase 4-lite (group-owned **endpoint** scripts: `scripts` polymorphic owner in `0050_group_scripts.sql`, `get_by_name_inherited`/`is_invocable_by_app` chain resolution, inherited `invoke()` + declarative route/trigger binding — all **live**, no body materialization), Phase 5 (the **declarative project tool maps onto the group tree**: the reconcile engine generalized to `ApplyOwner{App|Group}`, a `[group]` manifest kind, and a single atomic **tree apply** — `pic plan/apply --dir` reconciles a whole directory tree of `picloud.toml` nodes in one Postgres transaction, groups-before-apps so an app route can bind a group script created in the same tx; the bound token folds in each group's `structure_version`. The single-owner ownership **claim** shipped as §7 M1, the attach-point ceiling + blast-radius as §7 M2/M3, and per-env approval gating as §3 M3 — all server-authoritative (see the tail); declarative group **create/reparent** + structural-divergence detection shipped as §6 (`reconcile_group_structure_tx`/`reparent_group_tx`/`StructureMode`), so groups no longer need to pre-exist), Phase 4b (group **modules** + the **lexical (sealed-by-default) import resolver**, §5.5: owner-polymorphic `ModuleScript`, origin-rooted `ModuleSource::resolve` walking the importing node's chain, `ExecRequest.script_owner` threaded from every dispatch + `invoke()` site, `_source`-driven lexical chaining in `PicloudModuleResolver` with the compiled-module cache re-keyed by `ScriptId`, group modules/imports allowed, single-node dangling-import `plan` check — an inherited group script's imports **seal to the group**, a leaf can't shadow them), §5.5 **extension points** (opt-in polymorphism — **§5.5 now complete**: marker table `0051_extension_points.sql` (owner-polymorphic, CASCADE — structurally a `secrets` name; default body = a co-located `kind=module` script), `ModuleSource::resolve_policy` with **nearest-declaration-kind-wins** — a concrete module resolves lexically, an EP marker resolves **dynamically against the inheriting app** (its override else the default body up-chain), `NoProvider` is a hard error; declarative-only authoring via the `[app]`/`[group]` manifest key `extension_points = [...]`, reconcile mirrors `secrets`, single-node no-provider `plan` check, read-only `pic extension-points ls` + `pull` round-trip — the app can **override** a group default, the deliberate inverse of the Phase 4b sealed import), §11.6 **group-level collections — KV + DOCS + FILES slices** (full cross-app shared read/write: a group declares a collection shared via the `[group]` manifest `collections = [...]` → owner-polymorphic marker `0052_group_collections.sql` with a `kind` discriminator + a per-kind group-keyed store: `0053_group_kv_entries.sql` (`kind='kv'`), `0054_group_docs.sql` (`kind='docs'`, the queryable-JSON store), and `0055_group_files.sql` (`kind='files'`, blob metadata in Postgres + bytes on disk under `/files/groups//...`, a `groups/` infix disjoint from the per-app `files//` subtree so the existing recursive orphan sweeper covers both with zero change) — no `app_id`, a shared row belongs to the group; CASCADE on group delete, an app delete leaves the data. Scripts use the **explicit** `kv::shared_collection("name")` / `docs::shared_collection("name")` / `files::shared_collection("name")` handles (`shared` alone is a Rhai reserved word); `GroupKv`/`GroupDocs`/`GroupFilesServiceImpl` resolve the owning group from `cx.app_id`'s ancestor chain **filtered by kind** (nearest-wins) — **that walk is the isolation boundary**, a foreign app gets `CollectionNotShared`; a `kv`, a `docs`, and a `files` collection of the same name are distinct stores. The docs slice reuses the `docs_filter` DSL — `build_find_query` generalized on its owner column (`docs`/`app_id` vs `group_docs`/`group_id`, both literals); the files slice likewise generalized the atomic-write + checksum-on-read path helpers on an owner-relative dir (one source for the security-sensitive disk mechanics). **Reads open** to any subtree script (anonymous incl. — the declaration is the grant), **writes require an authenticated editor+** on the owning group (`GroupKvRead/Write`, `GroupDocsRead/Write`, `GroupFilesRead/Write`, `script_gate_require_principal` fails closed on anon). Declarative authoring is the **string-or-table** form `collections = ["catalog", { name = "articles", kind = "docs" }, { name = "assets", kind = "files" }]` (bare string = kv); reconcile keys markers by `(name, kind)`; read-only `pic collections ls --group` shows a kind column. Topic shared collections shipped as D2 (storeless), queue as D3 — see below), and **§4.5 group TRIGGER templates** (live, event kinds — a `[group]` declares a `[[triggers.kv|docs|files|pubsub]]` template binding a group-owned handler; `triggers` gained a polymorphic owner `0056_group_triggers.sql` mirroring `0050`; the dispatcher's `list_matching_kv/docs/files` + the pubsub publish fan-out prepend `CHAIN_LEVELS_CTE` + `JOIN chain c ON (t.app_id = c.app_owner OR t.group_id = c.group_owner)` so a descendant app's event matches its own triggers **plus** ancestor-group templates in one query, the handler running under the firing `app_id` — **the chain walk is the isolation boundary**, a sibling-subtree app never matches; stateful kinds cron/queue/email need materialization — see M5 below; per-app opt-out deferred; read-only `pic triggers ls --group`), and **§4.5 group ROUTE templates** (live, inherited — a `[group]` declares a `[[routes]]` template binding a group-owned endpoint; `routes` gained a polymorphic owner `0057_group_routes.sql` mirroring `0056`. Unlike triggers (per-event SQL), routes serve from the in-memory `RouteTable`, so the HTTP hot path can't resolve inheritance per request — instead the table **rebuild** expands templates into each descendant app's slice via `RouteRepository::list_effective` (all-apps generalization of `CHAIN_LEVELS_CTE`: every app × its ancestor chain ⋈ routes), and `compile_effective_routes` applies **nearest-owner-wins shadowing** (an app's own identical binding shadows the inherited template; non-identical bindings coexist under the matcher's existing precedence — a route picks one winner, unlike a fanning trigger). Because the table is a cache, inheritance is rebuilt **full-live** through the single `rebuild_route_table` chokepoint on every edge that changes it: route CRUD, apply, **and tree mutations** — app create/delete (`apps_api`) + group reparent (`groups_api`) — so a new app under a group serves its templates instantly. Host-claim validation is skipped for a group template (descendants serve it on their own host claim; templates use `host_kind = any`). **The chain expansion is the isolation boundary** — a sibling-subtree app never inherits (pinned by `tests/group_route_templates.rs` + the `group_routes` journey); read-only `pic routes ls --group`. Deferred: multi-node snapshot propagation), and **§11 tail per-app opt-out (template suppression)** (a descendant declines an inherited group template: an `[app]` declares `[suppress]` with `triggers = [...]` (handler script names) + `routes = [...]` (paths) — **coarse by reference**, not a full definition, since template row-ids churn on re-apply but a reference is stable (re-apply NoOp, may decline several templates bound to the same script/path). App-only marker `0058_template_suppressions.sql` (`app_id NOT NULL` CASCADE, a `target_kind` discriminator), reconciled with the extension-point marker pattern (prunable → re-inherits). Consumed at the two resolution points: the trigger dispatch queries gain a correlated `NOT EXISTS` anti-join (gated to `t.group_id IS NOT NULL`), and `compile_effective_routes` drops an inherited (`depth > 0`) route at a suppressed path (loaded via `RouteRepository::list_route_suppressions`). **Inheritance-only** — the `group_id IS NOT NULL` / `depth > 0` gates mean an app can only decline what it inherits, never its own or a sibling's (pinned by `tests/template_suppression.rs` + the `suppress` journey); a dangling suppress is an apply-time warning; read-only `pic suppress ls --app`. **Trust-model consequence:** group templates are advisory-by-default (run *unless* a descendant declines) — a footgun for compliance hooks (audit/security triggers a tenant can opt out of)), and **§11 tail `sealed` (mandatory) group templates** (closes that footgun: a `[group]` marks a route/event-trigger template `sealed = true` and the two suppression filters skip it, so a descendant's `[suppress]` is ignored — it fires/serves on every descendant. Column `sealed BOOLEAN` on `triggers` + `routes` (`0059_sealed_templates.sql`); the trigger anti-join gains `AND t.sealed = FALSE` (a sealed row is never excluded → fires through), and `compile_effective_routes` gates its suppression `continue` on `!er.route.sealed`. `sealed` lives on the shared `Route` DTO + manager-core `Trigger` (both pure data) so the apply diff sees the current value — part of the route Update comparison + the trigger identity, so toggling it re-applies. Authored per-template (`sealed = true` on a `[[routes]]`/`[[triggers.kv|docs|files|pubsub]]`), **group-only** — `validate_bundle_for` rejects it on an app owner (an app resource is never inherited). Sealing only *strengthens* the guarantee (a sealed template can't be declined; it never grants new reach — the chain walk is still the isolation boundary). The dangling-suppress warning also flags a suppression matching only sealed templates ("… is sealed — the suppression has no effect"), and `pic triggers/routes ls --group` show a `sealed` column; pinned by `tests/sealed_templates.rs` + the `sealed` journey. Deferred: multi-node snapshot propagation), and **§11 tail M1 group-level suppression** (`template_suppressions` gained a polymorphic owner `0060_group_suppressions.sql` — a `[group]` declares `[suppress]` to decline a template it inherits from a higher ancestor for its **whole subtree**. Both filters generalized to the chain: the trigger anti-join joins the `chain` CTE (`ts.app_id = sc.app_owner OR ts.group_id = sc.group_owner`), and `list_route_suppressions` expands group suppressions across descendants via the all-apps `app_chain` CTE (`compile_effective_routes` unchanged). Still inheritance-only + `sealed` overrides; owner-polymorphic `suppression_repo::list_for_owner`/`insert`/`delete`; read-only `pic suppress ls --group`; the ineffective-suppress warning walks `GROUP_CHAIN_LEVELS_CTE` for a group node. Pinned by `tests/group_suppression.rs` + the `suppress` journey), and **§11.6 M2 shared-collection triggers** (a write to a group SHARED collection now fires a trigger — closes the "group trigger has no app to watch" gap. A group-owned trigger marked `shared = true` (`shared BOOLEAN` on `triggers`, `0061_shared_triggers.sql`) watches the group's shared collection; the per-app `list_matching_kv/docs/files` add `AND t.shared = FALSE`, new `list_matching_shared_{kv,docs,files}(owning_group,…)` select `shared = TRUE` triggers on the owning group — the `shared` flag is the namespace boundary. `ServiceEventEmitter` gained `emit_shared(cx, owning_group, event)`; `GroupKv/Docs/FilesServiceImpl` (via a `with_events` builder) emit on write, and the outbox emitter stamps the WRITER app_id so the handler runs under the writer. `shared` is authored on a group `[[triggers.kv|docs|files]]`, is part of the diff identity, and `validate_bundle_for` rejects it on an app owner / a non-collection kind / an undeclared collection. Read-only `shared` column in `pic triggers ls --group`; pinned by `tests/shared_triggers.rs` + the `shared_triggers` journey. Shared pubsub triggers shipped as D2 — see below), and **§11.6 M3 per-group quotas** (global env-var ceilings enforced in the group write path: `PICLOUD_GROUP_KV_MAX_ROWS`/`_DOCS_MAX_ROWS` (per-group row count, new-key only), `PICLOUD_GROUP_{KV,DOCS}_MAX_TOTAL_BYTES` (per-group total stored bytes, projected-total check — Track A M4) + `_FILES_MAX_TOTAL_BYTES` (per-group total bytes); `group_quota` helpers + `count_rows`/`total_bytes` repo methods + `QuotaExceeded`/`TotalBytesQuotaExceeded` errors), and **§11.6 M4 read-only operator admin API** (`group_blobs_api` mirrors the per-app kv/files admin surface for a group's shared collections: `GET /groups/{id}/{kv,docs,files}[/{collection}/{key|id}]`, authz `GroupKvRead`/`GroupDocsRead`/`GroupFilesRead`; the host hoists the group repos to share one instance; `pic kv ls/get --group`; reads-only, writes stay script-only), and **§4.5 M5 stateful group trigger templates via materialization** (cron + queue + email — the kinds that can't resolve live because each needs a per-app row: cron `last_fired_at`, the queue one-consumer advisory lock, the email sealed inbound secret. A group `[[triggers.cron|queue|email]]` template is **materialized** into an app-owned copy per descendant (`materialized_from` col, `0062_materialized_triggers.sql` + the `0063_materialized_unique.sql` partial-unique index that makes rematerialization idempotent under concurrency) by `materialize::rematerialize_stateful_templates` — all-apps `app_chain` CTE ⋈ group-owned stateful templates, a precise create/delete diff preserving cron state — run full-live at the route-rebuild chokepoints (apply single+tree, app create/delete in `apps_api`, group reparent in `groups_api`; `AppsState`/`GroupsState` gained a `pool`). The scheduler + `list_active_queue_consumers` + `email_inbound_target` gained `AND t.app_id IS NOT NULL` so a group TEMPLATE is never dispatched directly (nor invocable via its own webhook URL) — only its per-app copies are; a queue copy is skipped-with-warning when the app already fills that queue's slot. **M5.5 email uses the shared-group-secret model:** the template resolves its `inbound_secret_ref` against the **group's own** secret store once at apply and seals it (`resolve_and_seal`/`insert_email_trigger_tx` generalized to a `SecretOwner`/`ScriptOwner`); email secrets are now **v1 AAD-bound to the SEALING OWNER** (Track A M3, migration `0069_email_secret_version` + `seal_email`/`open_email`; the group AAD is stable across rows, so materialization still copies the sealed bytes **verbatim** — the inbound path recovers the sealing group via `materialized_from` and opens under its AAD; a legacy v0 read path remains for pre-M3 rows). All descendant webhooks share the one group HMAC secret; an unset group secret fails apply hard. Loading a group's own set-secret names into `CurrentState` (was hardcoded empty) lets the plan-time email-secret check resolve against the group. Pinned by `tests/stateful_templates.rs` + the `stateful_templates` journey. **Deferred:** the per-app (non-shared) email secret model), and **D1 the `materialized` column** (`pic triggers ls --app` now shows a read-only `materialized` column — a copy of an M5 group stateful template reads `true`, distinct from a hand-authored trigger; a derived `materialized` bool = `materialized_from IS NOT NULL` threaded row→domain→API→CLI mirroring `sealed`/`shared`), and **§11.6 D2 shared TOPICS + shared pub/sub triggers** (a group declares a storeless `topic` shared collection (`group_collections.kind` widened to `topic`+`queue` in `0064`); a `[[triggers.pubsub]] shared = true` handler watches it. `BundleTrigger::Pubsub` gained `shared` (part of the identity); `validate_bundle_for` requires the topic pattern's ROOT segment (`events.*` → `events`) be a declared `kind='topic'` collection, group-only. Scripts publish via the explicit `pubsub::shared_topic("events").publish("created", msg)` handle → `GroupPubsubServiceImpl` resolves the owning group (kind `topic`) from `cx.app_id`'s chain, requires editor+ (`GroupPubsubPublish`, fails closed on anon), and fans out via `PubsubRepo::fan_out_shared_publish` to `shared = true` pubsub triggers on that group, each outbox row stamped the WRITER app_id (M2 model). The per-app `fan_out_publish` gained `AND t.shared = FALSE` — a shared trigger never fires on a per-app publish and vice-versa (the `shared` flag is the namespace boundary); the owning-group chain walk is the isolation boundary. Pinned by `tests/shared_topics.rs` + the `shared_topics` journey. **External SSE subscription shipped** (Track A M6): `GET /realtime/shared/topics/{topic}` streams a shared topic to external clients; `RealtimeAuthority::authorize_subscribe_shared` resolves the owning group from the subscriber app's chain (reads-open — the resolution IS the authorization, a foreign subtree 404s), and `GroupPubsubServiceImpl::with_realtime` bridges publish→broadcast after the durable fan-out. **Deferred:** multi-node broadcast propagation (cluster mode)), and **§11.6 D3 shared durable QUEUES** (a group declares a `queue` shared collection; any subtree app enqueues into ONE group-keyed store (`group_queue_messages`, `0065`, mirrors `queue_messages` keyed by `(group_id, collection)`; CASCADE on group delete) via `queue::shared_collection("name").enqueue(...)` → `GroupQueueServiceImpl` resolves the owning group (kind `queue`), requires editor+ (`GroupQueueEnqueue`, fails closed on anon). **Consumption is by COMPETING CONSUMERS:** a group `[[triggers.queue]] shared = true` consumer (shared threaded through `BundleTrigger::Queue` + identity) **materializes** a consumer copy per descendant app (`materialize` skips the M5 one-consumer-slot check for shared — each descendant intentionally gets a consumer), and all copies claim the SHARED store with `FOR UPDATE SKIP LOCKED` — each message delivered at-most-once across the subtree, scaling horizontally, each handler under its own `cx.app_id`. The dispatcher's queue arm gained a `shared_group` on `ActiveQueueConsumer` (recovered via a LEFT JOIN to the materialized copy's source template) and routes claim/ack/nack/terminal to the group store when set (`q_claim/q_ack/q_nack/q_terminal` helpers; a group claim normalizes to a `ClaimedMessage` under the consuming app); the reclaim task drains both stores. `validate_bundle_for` requires a shared queue on a group to name a declared `kind='queue'` collection; a shared queue on an app is rejected by the app-owner shared guard. Pinned by `tests/group_queue.rs` (competing-consumer at-most-once) + `stateful_templates.rs` + the `shared_queues` journey. **Group dead-letter store shipped** (Track A M2): an exhausted shared-queue message is preserved in `group_dead_letters` (`0068_group_dead_letters`) via `GroupQueueRepo::dead_letter` (atomic INSERT+DELETE) and is operator-visible at read-only `GET /api/v1/admin/groups/{id}/dead-letters` (`GroupKvRead`). **Deferred:** fan-out to a *shared* dead-letter trigger), and **§7 multi-repo ownership M1 — the single-owner claim** (the `owner_project` seam (0047) is now live behind a first-class `projects` table (`0066_projects.sql`, UUID pk + unique slug; `owner_project` FKs it `ON DELETE SET NULL` — un-claim, never cascade-destroy a tree). A `[project]` block (slug + optional name) in the repo's ROOT manifest declares identity (independent of the `[app]`/`[group]` XOR); the first apply with a new slug registers the project and **claims** each group node it touches. The claim runs inside the apply tx under the per-node advisory lock **before the diff**, so a conflict short-circuits with a **409** before any write. Pure policy `decide_group_claim`/`decide_app_owner` in `apply_service` (unit-tested, DB-free): unclaimed→claim, owner→no-op, foreign→conflict unless `--takeover` (which additionally requires `GroupAdmin` — ownership ⟂ RBAC, mapped to 403 vs the 409 conflict), no-project-into-a-claimed-subtree→conflict. **Apps carry no owner** — an app inherits ownership from its **nearest claimed ancestor group** (the ancestor walk, via `groups.ancestors` now carrying `owner_project` through its recursive CTE, is the isolation boundary); an unclaimed subtree stays open, so nothing changes until a repo first declares `[project]` (backward-compatible). `ProjectRepository` (read side) + tx free-fns `upsert_project_tx`/`read_group_owner_tx`/`write_group_owner_tx`; the claim deliberately does **not** bump `structure_version` (not a diff change → won't churn a pending bound plan). Wire: `project`/`takeover` on the apply request (both `#[serde(default)]` → the pre-M1 CLI stays compatible); the CLI surfaces the server's 409 message verbatim (covers `StateMoved` + `OwnershipConflict`). Visibility: `pic groups ls` `owner` column (server `list_with_owner` LEFT JOIN; the shared `Group` deserialize ignores the extra field so `pic groups tree`/dashboard are unaffected) + `pic apply --takeover`. Pinned by `apply_service` unit tests, `tests/projects_repo.rs`, and the `apply_ownership` journey. **M2 shipped — the attach-point ceiling:** a `[project] parent_group = ""` binds the repo under a pre-existing group; `check_within_attach` refuses (422 `OutsideAttachPoint`) any node not strictly within that subtree (a group node must be a *proper* descendant — you can't apply the attach point itself; an app node's group must be at-or-below it), resolved via `groups.ancestors` and enforced read-only before the claim in `apply_owner`/`apply_tree`; absent = instance root = no ceiling. **M3 shipped — plan preview + `pic projects ls`:** `pic plan` now carries the `[project]` and returns an `ownership` preview per node (`claim`/`owned`/`conflict`-owner-named/`unclaimed`, pure `preview_ownership`) plus, for a group node, the cross-repo **blast radius** (descendant apps owned by OTHER projects the change reaches, via `group_blast_radius` — a subtree CTE with per-group memoized nearest-claimed resolution); the attach ceiling is previewed at plan too. Read-only `pic projects ls` (`GET /api/v1/admin/projects`, `list_with_counts`) lists projects + owned-group counts. Pinned by the `apply_ownership` plan-preview case. **With M3 the §7 multi-repo ownership track (M1 claim · M2 attach ceiling · M3 preview + `pic projects ls`) is COMPLETE.** Also shipped: **§6 group-tree Tier 1** (declarative group create via dir-nesting · structural-divergence detection · declarative reparent — `reconcile_group_structure_tx`/`StructureMode`, 422 `StructuralDivergence`) and **§3 M3 the per-env approval gate**, now **server-authoritative** (migration `0067_project_environments`; the gate resolves the governing project from the target node's nearest-claimed ancestor — `governing_env_policy`/`_tree` — so omitting/spoofing `[project]` can't bypass it; an approved gated apply requires AppAdmin/GroupAdmin step-up + audit). -**Track A (v1.2 deferred-gap closeout, migrations 0067–0069) shipped to local main:** M1 hermetic approval gate · M2 shared-queue dead-letter store · M3 email-secret AAD v0→v1 · M4 per-group KV/docs byte quotas · M5 `set_if` compare-and-swap for KV (per-app + shared + Rhai SDK) · M6 shared-topic external SSE. **Audit 2026-07-11 remediation** (migration 0070, admin-session absolute cap) also shipped. **With that, v1.2 _Hierarchies_ is complete.** The **Workflows** track then shipped too (M1–M6: DAG execution + conditional `when`, nested sub-workflows, durable orchestrator, `workflow::start` SDK + admin run API + `pic workflows`, dashboard DAG + run-history; migrations `0071`/`0072`). Next: **§9.4 service interceptors** (the one remaining Workflows-track item) and multi-node cluster mode (the deferred multi-node route/broadcast propagation lives there). +**Track A (v1.2 deferred-gap closeout, migrations 0067–0069) shipped to local main:** M1 hermetic approval gate · M2 shared-queue dead-letter store · M3 email-secret AAD v0→v1 · M4 per-group KV/docs byte quotas · M5 `set_if` compare-and-swap for KV (per-app + shared + Rhai SDK) · M6 shared-topic external SSE. **Audit 2026-07-11 remediation** (migration 0070, admin-session absolute cap) also shipped. **With that, v1.2 _Hierarchies_ is complete.** The **Workflows** track then shipped too (M1–M6: DAG execution + conditional `when`, nested sub-workflows, durable orchestrator, `workflow::start` SDK + admin run API + `pic workflows`, dashboard DAG + run-history; migrations `0071`/`0072`). A **§9.4 service-interceptor** thin slice then shipped (migration `0073_interceptors.sql`): a `[[interceptors]]` block (app or group) binds a script to run before `kv::set`/`delete` and allow/deny it, resolved nearest-owner-wins on the app chain and run via the `invoke()` re-entry path; the rest of §9.4 (data transform, non-kv services, after-hooks, chaining) is deferred. Next: the rest of §9.4 and multi-node cluster mode (the deferred multi-node route/broadcast propagation lives there). **Data-model invariant:** app-owned data-plane tables (KV, docs, files, …) start with `app_id UUID NOT NULL REFERENCES apps(id) ON DELETE CASCADE`; the group-inheritable tables — _config_ (`vars`, `secrets`) and now group-owned _code_ (`scripts`, `0050`) — instead carry a **polymorphic owner**: nullable `group_id` and `app_id` with an exactly-one CHECK and per-owner partial-unique indexes (config is `ON DELETE CASCADE`, scripts `RESTRICT` — code is not data). Inheritance resolves **live** down `apps.group_id → groups.parent_id` via `CHAIN_LEVELS_CTE` (no materialized view); nearest-owner-wins with an app's own row shadowing the inherited one (CoW). Every Rhai SDK call resolves its app from `cx.app_id`, never a script-passed arg, and a group script always runs under the *inheriting* app's `cx.app_id` (the cross-app isolation boundary). @@ -145,6 +145,6 @@ Environment variables consumed by the `picloud` binary: ## Out of MVP -This section captured the original MVP cut. Most of it has since **shipped in v1.1.x**: queue triggers, cron triggers, inbound email (`email:receive`, HMAC-webhook model), KV / docs / email / users / HTTP SDKs, function-to-function `invoke()`, and secrets are all live (blueprint §12 Phase 4 table). The **Workflows** track (DAG + nested workflows) has since **shipped** (migrations `0071`/`0072`). **Still deferred:** **service interceptors** (§9.4), a metrics/observability dashboard, a raw SMTP-listener ingress, and **multi-node cluster mode**. Don't pre-build for them — but don't make decisions that close the door on them either. +This section captured the original MVP cut. Most of it has since **shipped in v1.1.x**: queue triggers, cron triggers, inbound email (`email:receive`, HMAC-webhook model), KV / docs / email / users / HTTP SDKs, function-to-function `invoke()`, and secrets are all live (blueprint §12 Phase 4 table). The **Workflows** track (DAG + nested workflows) has since **shipped** (migrations `0071`/`0072`). A **§9.4 service-interceptor** KV allow/deny slice has since shipped (migration `0073`). **Still deferred:** the rest of §9.4 (data transform, non-kv services, after-hooks, chaining), a metrics/observability dashboard, a raw SMTP-listener ingress, and **multi-node cluster mode**. Don't pre-build for them — but don't make decisions that close the door on them either. **Pulled forward to Phase 3 (pre-v1.1):** admin auth, multi-app scoping. The general cross-app **export/import** sharing model stays at v1.3+; note that v1.2 §11.6 shipped a narrower form — **group-owned shared collections** (KV/docs/files/topics/queues) let apps in one subtree share data through the owning group, with the ancestor-chain walk as the isolation boundary. See blueprint §11.5 + design-doc §11.6. diff --git a/crates/executor-core/src/sdk/interceptor.rs b/crates/executor-core/src/sdk/interceptor.rs new file mode 100644 index 0000000..4f819cb --- /dev/null +++ b/crates/executor-core/src/sdk/interceptor.rs @@ -0,0 +1,111 @@ +//! §9.4 Service Interceptors — the executor-side before-op hook. +//! +//! MVP: `kv::set` / `kv::delete` run an allow/deny interceptor first. The hook +//! (1) resolves the nearest interceptor script name for `(service, op)` on the +//! calling app's chain via the injected `InterceptorService` (a cheap indexed +//! query; `None` = un-hooked → allow, no further work), then (2) resolves that +//! name to a script and runs it through the SAME `invoke()` re-entry path +//! (`run_resolved_blocking`) — no second dispatch mechanism. The interceptor +//! receives the operation context as its request body and returns a map; the +//! op is DENIED iff that map has `allowed == false` (fail-open on any other +//! shape, documented — allow/deny only, no data transform). + +use std::sync::Arc; + +use picloud_shared::{InterceptorService, InvokeService, InvokeTarget, SdkCallCx}; +use rhai::EvalAltResult; +use serde_json::{json, Value as Json}; +use tokio::runtime::Handle as TokioHandle; + +use crate::engine::Engine; +use crate::sandbox::Limits; +use crate::sdk::bridge::runtime_err; +use crate::sdk::invoke::run_resolved_blocking; + +/// Run the before-op interceptor for `(service, op)` if one is registered. +/// Returns `Ok(())` to allow the operation (un-hooked, or the interceptor +/// allowed it) or an `Err` runtime error to deny it (the caller must NOT then +/// perform the write). `value` is the payload being written (`None` for delete). +#[allow(clippy::too_many_arguments)] +pub(super) fn run_before( + interceptors: &Arc, + invoke: &Arc, + self_engine: Option<&Arc>, + cx: &Arc, + limits: Limits, + service: &'static str, + op: &'static str, + collection: &str, + key: &str, + value: Option<&Json>, +) -> Result<(), Box> { + let handle = TokioHandle::try_current() + .map_err(|e| runtime_err(&format!("{service} interceptor: no tokio runtime: {e}")))?; + + // (1) Resolve the nearest interceptor script name. Un-hooked → allow. + let name = { + let interceptors = interceptors.clone(); + let cx = cx.clone(); + handle + .block_on(async move { interceptors.resolve_before(&cx, service, op).await }) + .map_err(|e| runtime_err(&format!("{service}::{op} interceptor resolve: {e}")))? + }; + let Some(name) = name else { return Ok(()) }; + + // An interceptor is registered but the engine back-reference isn't installed + // (a bare test engine without `set_self_weak`): there is no way to run it, so + // skip rather than block a write. Production always installs the back-ref. + let Some(self_engine) = self_engine else { + return Ok(()); + }; + + // Depth bound (shared with invoke / trigger fan-out): an interceptor that + // itself writes and re-enters can't recurse past the ceiling. + if cx.trigger_depth + 1 > limits.trigger_depth_max { + return Err(runtime_err(&format!( + "{service}::{op} interceptor `{name}`: depth limit exceeded (max {})", + limits.trigger_depth_max + ))); + } + + // (2) Resolve the interceptor script by name on the caller's chain and run + // it with the operation context as its request body. + let resolved = { + let invoke = invoke.clone(); + let cx = cx.clone(); + let target = InvokeTarget::Name(name.clone()); + handle + .block_on(async move { invoke.resolve(&cx, target).await }) + .map_err(|e| runtime_err(&format!("{service}::{op} interceptor `{name}`: {e}")))? + }; + let payload = json!({ + "service": service, + "action": op, + "collection": collection, + "key": key, + "value": value, + "caller_script_id": cx.script_id.to_string(), + "caller_execution_id": cx.execution_id.to_string(), + }); + let ret = run_resolved_blocking( + self_engine, + cx, + &resolved, + payload, + &format!("{service}::{op} interceptor `{name}`"), + )?; + + // Deny iff the interceptor returned a map with `allowed == false`. + if let Json::Object(m) = &ret { + if m.get("allowed") == Some(&Json::Bool(false)) { + let reason = m + .get("reason") + .and_then(Json::as_str) + .unwrap_or("denied by interceptor"); + return Err(runtime_err(&format!( + "{service}::{op} denied by interceptor `{name}`: {reason}" + ))); + } + } + Ok(()) +} diff --git a/crates/executor-core/src/sdk/invoke.rs b/crates/executor-core/src/sdk/invoke.rs index 2668410..aeeb379 100644 --- a/crates/executor-core/src/sdk/invoke.rs +++ b/crates/executor-core/src/sdk/invoke.rs @@ -163,9 +163,28 @@ fn invoke_blocking( .into() })?; - let execution_id = ExecutionId::new(); + // The callee's return is the response `body` JSON. Convert back to + // Dynamic for the caller. Status code + headers are dropped; the + // function-call mental model is "return value", not HTTP response. + let body = run_resolved_blocking(self_engine, cx, &resolved, args_json, &target_label)?; + Ok(json_to_dynamic(body)) +} + +/// Synchronous same-engine re-entry: build the callee `ExecRequest` (inheriting +/// the caller's app/principal/root/depth+1 and the callee's lexical owner), +/// compile through the per-Engine AST cache (F-P-004), execute, and return the +/// response `body` JSON. Shared by `invoke()` and the §9.4 interceptor hook — +/// the caller is responsible for the depth check (both do it before resolving). +/// `label` prefixes any compile/execute error. +pub(super) fn run_resolved_blocking( + self_engine: &Arc, + cx: &Arc, + resolved: &picloud_shared::ResolvedScript, + body_json: Json, + label: &str, +) -> Result> { let req = ExecRequest { - execution_id, + execution_id: ExecutionId::new(), request_id: cx.request_id, script_id: resolved.script_id, script_name: resolved.name.clone(), @@ -173,7 +192,7 @@ fn invoke_blocking( path: "/invoke".into(), method: String::new(), headers: BTreeMap::new(), - body: args_json, + body: body_json, params: BTreeMap::new(), query: BTreeMap::new(), rest: String::new(), @@ -183,42 +202,24 @@ fn invoke_blocking( // script's imports resolve from the group even when invoked by an // app. `None` falls back to `App(cx.app_id)` in the engine. script_owner: resolved.owner, - // Same-app invoke is a function call, not a re-auth boundary — - // inherit the caller's principal. + // Same-app re-entry is not a re-auth boundary — inherit the principal. principal: cx.principal.clone(), trigger_depth: cx.trigger_depth + 1, root_execution_id: cx.root_execution_id, is_dead_letter_handler: cx.is_dead_letter_handler, event: None, }; - - // F-P-004: synchronous re-entry — route through the per-Engine - // AST cache so each callee parses once per (script_id, updated_at), - // not once per invoke. Composed workflows multiply parse cost by - // depth; the cache cuts that to constant compile + N executions. let ast = self_engine .compile_for_identity(resolved.script_id, resolved.updated_at, &resolved.source) .map_err(|e| -> Box { - EvalAltResult::ErrorRuntime( - format!("invoke({target_label}): {e}").into(), - rhai::Position::NONE, - ) - .into() + EvalAltResult::ErrorRuntime(format!("{label}: {e}").into(), rhai::Position::NONE).into() })?; let resp = self_engine .execute_ast(&ast, req) .map_err(|e| -> Box { - EvalAltResult::ErrorRuntime( - format!("invoke({target_label}): {e}").into(), - rhai::Position::NONE, - ) - .into() + EvalAltResult::ErrorRuntime(format!("{label}: {e}").into(), rhai::Position::NONE).into() })?; - - // The callee's return is the response `body` JSON. Convert back to - // Dynamic for the caller. Status code + headers are dropped; the - // function-call mental model is "return value", not HTTP response. - Ok(json_to_dynamic(resp.body)) + Ok(resp.body) } /// Accept a string (route path OR script name) or a Rhai script-id diff --git a/crates/executor-core/src/sdk/kv.rs b/crates/executor-core/src/sdk/kv.rs index 67e99bb..42d4f3e 100644 --- a/crates/executor-core/src/sdk/kv.rs +++ b/crates/executor-core/src/sdk/kv.rs @@ -30,18 +30,27 @@ use std::sync::Arc; -use picloud_shared::{GroupKvService, KvService, SdkCallCx, Services}; +use picloud_shared::{ + GroupKvService, InterceptorService, InvokeService, KvService, SdkCallCx, Services, +}; use rhai::{Array, Dynamic, Engine as RhaiEngine, EvalAltResult, Map, Module}; use super::bridge::{block_on, dynamic_to_json, json_to_dynamic}; +use crate::engine::Engine; +use crate::sandbox::Limits; -/// Per-call handle captured by the Rhai SDK. Cheap to clone (two Arcs -/// plus an owned string). +/// Per-call handle captured by the Rhai SDK. Cheap to clone (a few Arcs plus an +/// owned string). Carries the §9.4 interceptor deps so `set`/`delete` can run a +/// before-op allow/deny hook (the resolver + the `invoke()` re-entry engine). #[derive(Clone)] pub struct KvHandle { collection: String, service: Arc, cx: Arc, + interceptors: Arc, + invoke: Arc, + self_engine: Option>, + limits: Limits, } /// §11.6 shared-collection handle, returned by `kv::shared_collection(name)`. A distinct @@ -56,9 +65,17 @@ pub struct GroupKvHandle { cx: Arc, } -pub(super) fn register(engine: &mut RhaiEngine, services: &Services, cx: Arc) { +pub(super) fn register( + engine: &mut RhaiEngine, + services: &Services, + cx: Arc, + limits: Limits, + self_engine: Option>, +) { let kv_service = services.kv.clone(); let group_kv_service = services.group_kv.clone(); + let interceptors = services.interceptors.clone(); + let invoke = services.invoke.clone(); // `kv::collection(name)` / `kv::shared_collection(name)` — both constructors live in // the `kv` static module so the script-visible calls are `kv::collection` @@ -67,6 +84,9 @@ pub(super) fn register(engine: &mut RhaiEngine, services: &Services, cx: Arc Result> { @@ -77,6 +97,10 @@ pub(super) fn register(engine: &mut RhaiEngine, services: &Services, cx: Arc Result<(), Box> { - let h = handle.clone(); let json = dynamic_to_json(&value); + // §9.4 before-op interceptor (allow/deny). A denial errors here and + // the write below never runs. + super::interceptor::run_before( + &handle.interceptors, + &handle.invoke, + handle.self_engine.as_ref(), + &handle.cx, + handle.limits, + "kv", + "set", + &handle.collection, + key, + Some(&json), + )?; + let h = handle.clone(); block_on("kv", async move { h.service.set(&h.cx, &h.collection, key, json).await }) @@ -197,6 +235,18 @@ fn register_delete(engine: &mut RhaiEngine) { engine.register_fn( "delete", |handle: &mut KvHandle, key: &str| -> Result> { + super::interceptor::run_before( + &handle.interceptors, + &handle.invoke, + handle.self_engine.as_ref(), + &handle.cx, + handle.limits, + "kv", + "delete", + &handle.collection, + key, + None, + )?; let h = handle.clone(); block_on("kv", async move { h.service.delete(&h.cx, &h.collection, key).await diff --git a/crates/executor-core/src/sdk/mod.rs b/crates/executor-core/src/sdk/mod.rs index 911a469..9d1f3be 100644 --- a/crates/executor-core/src/sdk/mod.rs +++ b/crates/executor-core/src/sdk/mod.rs @@ -18,6 +18,7 @@ pub mod docs; pub mod email; pub mod files; pub mod http; +pub mod interceptor; pub mod invoke; pub mod kv; pub mod pubsub; @@ -55,7 +56,7 @@ pub fn register_all( limits: Limits, self_engine: Option>, ) { - kv::register(engine, services, cx.clone()); + kv::register(engine, services, cx.clone(), limits, self_engine.clone()); docs::register(engine, services, cx.clone()); dead_letters::register(engine, services, cx.clone()); http::register(engine, services, cx.clone()); diff --git a/crates/manager-core/migrations/0073_interceptors.sql b/crates/manager-core/migrations/0073_interceptors.sql new file mode 100644 index 0000000..847ba60 --- /dev/null +++ b/crates/manager-core/migrations/0073_interceptors.sql @@ -0,0 +1,42 @@ +-- §9.4 Service Interceptors (v1.2) — before-op allow/deny hooks. +-- +-- A marker `(owner, service, op) -> interceptor_script` declares that a script +-- runs BEFORE a data-plane operation and may deny it. MVP scope: `service='kv'`, +-- `op IN ('set','delete')`, allow/deny only (no data transform, no chaining, +-- no after-hooks). The interceptor is itself a script owned by the same node +-- (or an ancestor group); it is resolved + run through the existing `invoke()` +-- re-entry path, so this table holds only the MARKER — pure declaration, +-- structurally like an `extension_points` row (0051): config, not code, so +-- ON DELETE CASCADE. +-- +-- Ownership is polymorphic (mirrors extension_points/vars/secrets/scripts): +-- exactly one of (app_id, group_id) is set. Resolution walks the calling app's +-- chain (app, then nearest ancestor group) and picks the nearest declaration +-- for a (service, op) — nearest-owner-wins, so an app overrides a group's +-- interceptor, the deliberate inverse of a sealed import. + +CREATE TABLE interceptors ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + group_id UUID REFERENCES groups(id) ON DELETE CASCADE, + app_id UUID REFERENCES apps(id) ON DELETE CASCADE, + CONSTRAINT interceptors_owner_exactly_one + CHECK ((group_id IS NULL) <> (app_id IS NULL)), + -- The intercepted operation. MVP: service='kv', op IN ('set','delete'). + service TEXT NOT NULL, + op TEXT NOT NULL, + -- Name of the interceptor script (resolved on the owner's chain at run time). + interceptor_script TEXT NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +-- One marker per (owner, service, op). Partial because the owner is split +-- across two nullable columns. +CREATE UNIQUE INDEX interceptors_group_uidx + ON interceptors (group_id, service, op) WHERE group_id IS NOT NULL; +CREATE UNIQUE INDEX interceptors_app_uidx + ON interceptors (app_id, service, op) WHERE app_id IS NOT NULL; + +-- Lookup indexes for the resolver's chain join + list-by-owner. +CREATE INDEX interceptors_group_id_idx ON interceptors (group_id) WHERE group_id IS NOT NULL; +CREATE INDEX interceptors_app_id_idx ON interceptors (app_id) WHERE app_id IS NOT NULL; diff --git a/crates/manager-core/src/apply_service.rs b/crates/manager-core/src/apply_service.rs index 21d5bf1..60bf7f7 100644 --- a/crates/manager-core/src/apply_service.rs +++ b/crates/manager-core/src/apply_service.rs @@ -107,6 +107,21 @@ pub struct Bundle { /// `[group]` carrying workflows is rejected in `validate_bundle_for`. #[serde(default)] pub workflows: Vec, + /// §9.4 service interceptors: before-op allow/deny hooks. One entry per + /// `(service, op)` guarded, naming the interceptor script. App- OR + /// group-owned; resolved nearest-owner-wins on a descendant app's chain. + #[serde(default)] + pub interceptors: Vec, +} + +/// One §9.4 interceptor marker on the wire: which `(service, op)` it guards and +/// the interceptor script name. The CLI expands a manifest `[[interceptors]]` +/// entry's `ops = [...]` into one of these per op. +#[derive(Debug, Clone, Deserialize)] +pub struct BundleInterceptor { + pub service: String, + pub op: String, + pub script: String, } /// One declared workflow on the wire: a name + its (already-parsed) DAG. @@ -467,6 +482,10 @@ pub struct Plan { /// v1.2 Workflows: workflow definitions, keyed by name. #[serde(default)] pub workflows: Vec, + /// §9.4 interceptor markers, keyed `"{service}/{op}"` with the script as the + /// value (create/update/delete like `vars`). + #[serde(default)] + pub interceptors: Vec, } impl Plan { @@ -483,6 +502,7 @@ impl Plan { .chain(&self.collections) .chain(&self.suppressions) .chain(&self.workflows) + .chain(&self.interceptors) .all(|c| c.op == Op::NoOp) } } @@ -765,6 +785,9 @@ pub struct CurrentState { pub suppressions: Vec<(String, String)>, /// v1.2 Workflows owned directly by this node (app-owned; empty for a group). pub workflows: Vec, + /// §9.4 interceptor markers declared directly at this node, as + /// `(service, op, script)` triples. + pub interceptors: Vec<(String, String, String)>, } /// One row of the read-only extension-point report (§5.5). @@ -1040,6 +1063,7 @@ impl ApplyService { /// `validate_bundle`, plus a group-node guard: a group owns only scripts + /// vars (+ declared secret names), so a `[group]` manifest carrying routes /// or triggers is a 422 — those are app concerns. + #[allow(clippy::too_many_lines)] fn validate_bundle_for( &self, is_group: bool, @@ -1173,6 +1197,35 @@ impl ApplyService { for w in &bundle.workflows { validate_workflow_definition(&w.name, &w.definition).map_err(ApplyError::Invalid)?; } + // §9.4 interceptors (MVP): `service = "kv"`, `op ∈ {set, delete}`, + // authored on app OR group. One marker per `(service, op)` — a duplicate + // would collide on the reconcile key and the DB's partial-unique index. + let mut seen_hooks: HashSet<(String, String)> = HashSet::new(); + for i in &bundle.interceptors { + if i.service != "kv" { + return Err(ApplyError::Invalid(format!( + "interceptor service `{}` is not supported — only `kv` (MVP)", + i.service + ))); + } + if i.op != "set" && i.op != "delete" { + return Err(ApplyError::Invalid(format!( + "interceptor op `{}` is not supported for kv — only `set` / `delete`", + i.op + ))); + } + if i.script.trim().is_empty() { + return Err(ApplyError::Invalid( + "an interceptor must name a `script`".into(), + )); + } + if !seen_hooks.insert((i.service.clone(), i.op.clone())) { + return Err(ApplyError::Invalid(format!( + "duplicate interceptor for `{}/{}`", + i.service, i.op + ))); + } + } self.validate_bundle(bundle, inherited_endpoints) } @@ -1450,6 +1503,35 @@ impl ApplyService { } } + // 3c.5. §9.4 interceptor markers — upsert each Create/Update by + // `(service, op)`. A changed script is an in-place Update; deletes happen + // in prune. + for ch in &plan.interceptors { + if ch.op == Op::Create || ch.op == Op::Update { + let bi = bundle + .interceptors + .iter() + .find(|i| format!("{}/{}", i.service, i.op) == ch.key) + .ok_or_else(|| { + ApplyError::Backend("internal: interceptor plan/bundle mismatch".into()) + })?; + crate::interceptor_repo::insert_interceptor_tx( + &mut *tx, + owner.as_script_owner(), + &bi.service, + &bi.op, + &bi.script, + ) + .await + .map_err(|e| ApplyError::Backend(e.to_string()))?; + if ch.op == Op::Create { + report.interceptors_created += 1; + } else { + report.interceptors_updated += 1; + } + } + } + // 3d. Shared group-collection markers (§11.6) — insert each Create // (idempotent). Name-only identity; deletes happen in prune. Pruning a // marker hides the store but does NOT drop its `group_kv_entries` data @@ -1580,6 +1662,24 @@ impl ApplyService { } } + // §9.4 interceptor markers are prunable config too. The delete keys + // by `(service, op)` (the marker's identity), so a script-change + // Update above is never clobbered here. + for ch in &plan.interceptors { + if ch.op == Op::Delete { + let (service, op) = ch.key.split_once('/').unwrap_or((ch.key.as_str(), "")); + crate::interceptor_repo::delete_interceptor_tx( + &mut *tx, + owner.as_script_owner(), + service, + op, + ) + .await + .map_err(|e| ApplyError::Backend(e.to_string()))?; + report.interceptors_deleted += 1; + } + } + // Shared group-collection markers are prunable config too (§11.6). // Only the marker is removed here — the data survives until the // owning group is deleted. @@ -4048,6 +4148,14 @@ impl ApplyService { } ApplyOwner::Group(_) => Vec::new(), }; + // §9.4 interceptor markers declared directly at this node. + let interceptors = + crate::interceptor_repo::list_for_owner(&self.pool, owner.as_script_owner()) + .await + .map_err(|e| ApplyError::Backend(e.to_string()))? + .into_iter() + .map(|m| (m.service, m.op, m.script)) + .collect(); Ok(CurrentState { scripts, routes, @@ -4058,6 +4166,7 @@ impl ApplyService { collections, suppressions, workflows, + interceptors, }) } @@ -4124,9 +4233,58 @@ fn compute_diff_with_names( collections: diff_collections(current, bundle), suppressions: diff_suppressions(current, bundle), workflows: diff_workflows(current, bundle), + interceptors: diff_interceptors(current, bundle), } } +/// Diff §9.4 interceptor markers by `"{service}/{op}"` key with the script name +/// as the value — create/update/noop/delete like `vars`. A changed script for +/// the same `(service, op)` is an `Update`; a live marker the manifest stops +/// declaring is a `Delete` (applied only under `--prune`). +fn diff_interceptors(current: &CurrentState, bundle: &Bundle) -> Vec { + let key = |service: &str, op: &str| format!("{service}/{op}"); + let live: HashMap = current + .interceptors + .iter() + .map(|(s, o, script)| (key(s, o), script.as_str())) + .collect(); + let mut out = Vec::new(); + for bi in &bundle.interceptors { + let k = key(&bi.service, &bi.op); + match live.get(&k) { + Some(cur) if *cur == bi.script => out.push(ResourceChange { + op: Op::NoOp, + key: k, + detail: None, + }), + Some(_) => out.push(ResourceChange { + op: Op::Update, + key: k, + detail: Some("interceptor script changed".into()), + }), + None => out.push(ResourceChange { + op: Op::Create, + key: k, + detail: Some(bi.script.clone()), + }), + } + } + for (s, o, _) in ¤t.interceptors { + let present = bundle + .interceptors + .iter() + .any(|bi| bi.service == *s && bi.op == *o); + if !present { + out.push(ResourceChange { + op: Op::Delete, + key: key(s, o), + detail: Some("on server, not declared".into()), + }); + } + } + out +} + /// Diff workflows by `lower(name)`. Like scripts, a workflow has a stable name /// identity with a mutable body, so a changed definition (or `enabled`) is an /// `Update`, not a delete+create. Live workflows absent from the manifest are @@ -5472,6 +5630,12 @@ pub struct ApplyReport { pub workflows_updated: u32, #[serde(default)] pub workflows_deleted: u32, + #[serde(default)] + pub interceptors_created: u32, + #[serde(default)] + pub interceptors_updated: u32, + #[serde(default)] + pub interceptors_deleted: u32, #[serde(skip_serializing_if = "Vec::is_empty")] pub warnings: Vec, } @@ -5816,6 +5980,7 @@ mod tests { suppress_triggers: Vec::new(), suppress_routes: Vec::new(), workflows: Vec::new(), + interceptors: Vec::new(), }; let plan = compute_diff(¤t, &bundle); assert_eq!(plan.scripts.len(), 1); @@ -5843,6 +6008,7 @@ mod tests { suppress_triggers: Vec::new(), suppress_routes: Vec::new(), workflows: Vec::new(), + interceptors: Vec::new(), }; let plan = compute_diff(¤t, &bundle); assert!(plan.is_noop(), "expected all no-op, got {plan:?}"); @@ -5865,6 +6031,7 @@ mod tests { suppress_triggers: Vec::new(), suppress_routes: Vec::new(), workflows: Vec::new(), + interceptors: Vec::new(), }; let plan = compute_diff(¤t, &bundle); assert_eq!(plan.scripts[0].op, Op::Update); @@ -5888,6 +6055,7 @@ mod tests { suppress_triggers: Vec::new(), suppress_routes: Vec::new(), workflows: Vec::new(), + interceptors: Vec::new(), }; let plan = compute_diff(¤t, &bundle); assert_eq!(plan.scripts[0].op, Op::Delete); @@ -5942,6 +6110,7 @@ mod tests { suppress_triggers: Vec::new(), suppress_routes: Vec::new(), workflows: Vec::new(), + interceptors: Vec::new(), }; let plan = compute_diff(¤t, &bundle); assert_eq!(plan.routes[0].op, Op::Update); @@ -5959,6 +6128,7 @@ mod tests { suppress_triggers: Vec::new(), suppress_routes: Vec::new(), workflows: Vec::new(), + interceptors: Vec::new(), } } diff --git a/crates/manager-core/src/interceptor_repo.rs b/crates/manager-core/src/interceptor_repo.rs new file mode 100644 index 0000000..d53b76d --- /dev/null +++ b/crates/manager-core/src/interceptor_repo.rs @@ -0,0 +1,184 @@ +//! §9.4 Service-interceptor markers — the `interceptors` table (0073). +//! +//! A marker `(owner, service, op) -> interceptor_script` declares a before-op +//! allow/deny hook. Pure declaration (like `extension_points`, 0051); the +//! interceptor's behaviour is a normal script resolved + run through `invoke()` +//! re-entry. This module holds the read + transactional-write helpers (free +//! functions over `&PgPool` / `&mut Transaction`, keyed by [`ScriptOwner`]), +//! plus the runtime [`resolve_before`] chain walk (nearest-owner-wins). + +use picloud_shared::{AppId, ScriptOwner}; +use sqlx::{PgPool, Postgres, Transaction}; + +use crate::config_resolver::CHAIN_LEVELS_CTE; + +/// One interceptor marker at an owner: which `(service, op)` it guards and the +/// interceptor script name. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct InterceptorMarker { + pub service: String, + pub op: String, + pub script: String, +} + +/// List the markers declared **directly at** `owner` (not inherited), ordered +/// deterministically. Used by `load_current` (apply diff) and `interceptors ls`. +pub async fn list_for_owner( + pool: &PgPool, + owner: ScriptOwner, +) -> Result, sqlx::Error> { + let rows: Vec<(String, String, String)> = match owner { + ScriptOwner::App(a) => { + sqlx::query_as( + "SELECT service, op, interceptor_script FROM interceptors \ + WHERE app_id = $1 ORDER BY service, op", + ) + .bind(a.into_inner()) + .fetch_all(pool) + .await? + } + ScriptOwner::Group(g) => { + sqlx::query_as( + "SELECT service, op, interceptor_script FROM interceptors \ + WHERE group_id = $1 ORDER BY service, op", + ) + .bind(g.into_inner()) + .fetch_all(pool) + .await? + } + }; + Ok(rows + .into_iter() + .map(|(service, op, script)| InterceptorMarker { + service, + op, + script, + }) + .collect()) +} + +/// All markers **visible to an app** — declared at the app or any ancestor +/// group (each `(service, op)` resolved nearest-owner-wins). Used by +/// `interceptors ls --app` and by the resolver's chain view. +pub async fn list_on_app_chain( + pool: &PgPool, + app_id: AppId, +) -> Result, sqlx::Error> { + let rows: Vec<(String, String, String)> = sqlx::query_as(&format!( + "{CHAIN_LEVELS_CTE} \ + SELECT DISTINCT ON (i.service, i.op) i.service, i.op, i.interceptor_script \ + FROM chain c \ + JOIN interceptors i ON (i.app_id = c.app_owner OR i.group_id = c.group_owner) \ + ORDER BY i.service, i.op, c.depth ASC", + )) + .bind(app_id.into_inner()) + .fetch_all(pool) + .await?; + Ok(rows + .into_iter() + .map(|(service, op, script)| InterceptorMarker { + service, + op, + script, + }) + .collect()) +} + +/// Resolve the interceptor script guarding `(service, op)` for a calling app — +/// the nearest declaration on the app's chain (app, then nearest ancestor +/// group). `None` when un-hooked. **This chain walk is the resolution boundary** +/// — a sibling-subtree app never sees another subtree's interceptor. +pub async fn resolve_before( + pool: &PgPool, + app_id: AppId, + service: &str, + op: &str, +) -> Result, sqlx::Error> { + let row: Option<(String,)> = sqlx::query_as(&format!( + "{CHAIN_LEVELS_CTE} \ + SELECT i.interceptor_script FROM chain c \ + JOIN interceptors i ON (i.app_id = c.app_owner OR i.group_id = c.group_owner) \ + WHERE i.service = $2 AND i.op = $3 \ + ORDER BY c.depth ASC LIMIT 1", + )) + .bind(app_id.into_inner()) + .bind(service) + .bind(op) + .fetch_optional(pool) + .await?; + Ok(row.map(|(s,)| s)) +} + +/// Upsert a marker at `owner` in the apply transaction. A re-apply that changes +/// only the interceptor script updates it in place (no version churn on the +/// `(service, op)` identity). +pub async fn insert_interceptor_tx( + tx: &mut Transaction<'_, Postgres>, + owner: ScriptOwner, + service: &str, + op: &str, + script: &str, +) -> Result<(), sqlx::Error> { + match owner { + ScriptOwner::App(a) => { + sqlx::query( + "INSERT INTO interceptors (app_id, service, op, interceptor_script) \ + VALUES ($1, $2, $3, $4) \ + ON CONFLICT (app_id, service, op) WHERE app_id IS NOT NULL \ + DO UPDATE SET interceptor_script = EXCLUDED.interceptor_script, updated_at = NOW()", + ) + .bind(a.into_inner()) + .bind(service) + .bind(op) + .bind(script) + .execute(&mut **tx) + .await?; + } + ScriptOwner::Group(g) => { + sqlx::query( + "INSERT INTO interceptors (group_id, service, op, interceptor_script) \ + VALUES ($1, $2, $3, $4) \ + ON CONFLICT (group_id, service, op) WHERE group_id IS NOT NULL \ + DO UPDATE SET interceptor_script = EXCLUDED.interceptor_script, updated_at = NOW()", + ) + .bind(g.into_inner()) + .bind(service) + .bind(op) + .bind(script) + .execute(&mut **tx) + .await?; + } + } + Ok(()) +} + +/// Delete a marker at `owner` (by `(service, op)`), in the apply transaction. +/// Used by `--prune` when the manifest stops declaring it. +pub async fn delete_interceptor_tx( + tx: &mut Transaction<'_, Postgres>, + owner: ScriptOwner, + service: &str, + op: &str, +) -> Result<(), sqlx::Error> { + match owner { + ScriptOwner::App(a) => { + sqlx::query("DELETE FROM interceptors WHERE app_id = $1 AND service = $2 AND op = $3") + .bind(a.into_inner()) + .bind(service) + .bind(op) + .execute(&mut **tx) + .await?; + } + ScriptOwner::Group(g) => { + sqlx::query( + "DELETE FROM interceptors WHERE group_id = $1 AND service = $2 AND op = $3", + ) + .bind(g.into_inner()) + .bind(service) + .bind(op) + .execute(&mut **tx) + .await?; + } + } + Ok(()) +} diff --git a/crates/manager-core/src/interceptor_service.rs b/crates/manager-core/src/interceptor_service.rs new file mode 100644 index 0000000..fde9997 --- /dev/null +++ b/crates/manager-core/src/interceptor_service.rs @@ -0,0 +1,37 @@ +//! `InterceptorServiceImpl` — the Postgres-backed §9.4 interceptor resolver +//! injected into `Services`. Resolve-only: it maps `(cx.app_id, service, op)` +//! to the nearest interceptor script name on the calling app's chain (via +//! [`crate::interceptor_repo::resolve_before`]). Running that script is the +//! executor's job (the `invoke()` re-entry path), which keeps `executor-core` +//! Postgres-free. + +use async_trait::async_trait; +use picloud_shared::{InterceptorService, SdkCallCx}; +use sqlx::PgPool; + +pub struct InterceptorServiceImpl { + pool: PgPool, +} + +impl InterceptorServiceImpl { + #[must_use] + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +#[async_trait] +impl InterceptorService for InterceptorServiceImpl { + async fn resolve_before( + &self, + cx: &SdkCallCx, + service: &str, + op: &str, + ) -> Result, String> { + // `app_id` derives from `cx` (never a script arg) — the isolation + // boundary; the chain walk then scopes resolution to this app's subtree. + crate::interceptor_repo::resolve_before(&self.pool, cx.app_id, service, op) + .await + .map_err(|e| e.to_string()) + } +} diff --git a/crates/manager-core/src/lib.rs b/crates/manager-core/src/lib.rs index 7e60635..be52121 100644 --- a/crates/manager-core/src/lib.rs +++ b/crates/manager-core/src/lib.rs @@ -68,6 +68,8 @@ pub mod group_repo; pub mod group_scripts_api; pub mod groups_api; pub mod http_service; +pub mod interceptor_repo; +pub mod interceptor_service; pub mod invoke_service; pub mod kv_api; pub mod kv_repo; diff --git a/crates/manager-core/tests/expected_schema.txt b/crates/manager-core/tests/expected_schema.txt index 1cc8f8c..ec9a544 100644 --- a/crates/manager-core/tests/expected_schema.txt +++ b/crates/manager-core/tests/expected_schema.txt @@ -308,6 +308,16 @@ table: groups created_at: timestamp with time zone NOT NULL default=now() updated_at: timestamp with time zone NOT NULL default=now() +table: interceptors + id: uuid NOT NULL default=gen_random_uuid() + group_id: uuid NULL + app_id: uuid NULL + service: text NOT NULL + op: text NOT NULL + interceptor_script: text NOT NULL + created_at: timestamp with time zone NOT NULL default=now() + updated_at: timestamp with time zone NOT NULL default=now() + table: kv_entries app_id: uuid NOT NULL collection: text NOT NULL @@ -663,6 +673,13 @@ indexes on groups: groups_pkey: public.groups USING btree (id) groups_slug_key: public.groups USING btree (slug) +indexes on interceptors: + interceptors_app_id_idx: public.interceptors USING btree (app_id) WHERE (app_id IS NOT NULL) + interceptors_app_uidx: public.interceptors USING btree (app_id, service, op) WHERE (app_id IS NOT NULL) + interceptors_group_id_idx: public.interceptors USING btree (group_id) WHERE (group_id IS NOT NULL) + interceptors_group_uidx: public.interceptors USING btree (group_id, service, op) WHERE (group_id IS NOT NULL) + interceptors_pkey: public.interceptors USING btree (id) + indexes on kv_entries: idx_kv_entries_app_collection: public.kv_entries USING btree (app_id, collection) kv_entries_pkey: public.kv_entries USING btree (app_id, collection, key) @@ -933,6 +950,12 @@ constraints on groups: [PRIMARY KEY] groups_pkey: PRIMARY KEY (id) [UNIQUE] groups_slug_key: UNIQUE (slug) +constraints on interceptors: + [CHECK] interceptors_owner_exactly_one: CHECK (((group_id IS NULL) <> (app_id IS NULL))) + [FOREIGN KEY] interceptors_app_id_fkey: FOREIGN KEY (app_id) REFERENCES apps(id) ON DELETE CASCADE + [FOREIGN KEY] interceptors_group_id_fkey: FOREIGN KEY (group_id) REFERENCES groups(id) ON DELETE CASCADE + [PRIMARY KEY] interceptors_pkey: PRIMARY KEY (id) + constraints on kv_entries: [FOREIGN KEY] kv_entries_app_id_fkey: FOREIGN KEY (app_id) REFERENCES apps(id) ON DELETE CASCADE [PRIMARY KEY] kv_entries_pkey: PRIMARY KEY (app_id, collection, key) @@ -1121,3 +1144,4 @@ constraints on workflows: 0070: admin session absolute expiry 0071: workflows 0072: execution source workflow + 0073: interceptors diff --git a/crates/picloud-cli/src/client.rs b/crates/picloud-cli/src/client.rs index 5ac726d..8ae40dd 100644 --- a/crates/picloud-cli/src/client.rs +++ b/crates/picloud-cli/src/client.rs @@ -1738,6 +1738,8 @@ pub struct PlanDto { pub suppressions: Vec, #[serde(default)] pub workflows: Vec, + #[serde(default)] + pub interceptors: Vec, /// Fingerprint of the live state this plan was computed against; carried /// in `.picloud/` and replayed to `apply` for the bound-plan check. #[serde(default)] @@ -1815,6 +1817,8 @@ pub struct NodePlanDto { pub suppressions: Vec, #[serde(default)] pub workflows: Vec, + #[serde(default)] + pub interceptors: Vec, /// §7 M3: this node's ownership outcome under the tree's `[project]`. #[serde(default)] pub ownership: Option, @@ -1880,6 +1884,12 @@ pub struct ApplyReportDto { #[serde(default)] pub workflows_deleted: u32, #[serde(default)] + pub interceptors_created: u32, + #[serde(default)] + pub interceptors_updated: u32, + #[serde(default)] + pub interceptors_deleted: u32, + #[serde(default)] pub warnings: Vec, } diff --git a/crates/picloud-cli/src/cmds/apply.rs b/crates/picloud-cli/src/cmds/apply.rs index d684891..8a80100 100644 --- a/crates/picloud-cli/src/cmds/apply.rs +++ b/crates/picloud-cli/src/cmds/apply.rs @@ -155,6 +155,15 @@ pub async fn run( "+{} ~{} -{}", report.workflows_created, report.workflows_updated, report.workflows_deleted ), + ) + .field( + "interceptors", + format!( + "+{} ~{} -{}", + report.interceptors_created, + report.interceptors_updated, + report.interceptors_deleted + ), ); for w in &report.warnings { block.field("warning", w.clone()); @@ -283,6 +292,15 @@ pub async fn run_tree( "+{} ~{} -{}", report.workflows_created, report.workflows_updated, report.workflows_deleted ), + ) + .field( + "interceptors", + format!( + "+{} ~{} -{}", + report.interceptors_created, + report.interceptors_updated, + report.interceptors_deleted + ), ); for w in &report.warnings { block.field("warning", w.clone()); diff --git a/crates/picloud-cli/src/cmds/init.rs b/crates/picloud-cli/src/cmds/init.rs index d3f0b24..5198c3a 100644 --- a/crates/picloud-cli/src/cmds/init.rs +++ b/crates/picloud-cli/src/cmds/init.rs @@ -144,6 +144,7 @@ fn scaffold_manifest(slug: &str, name: &str) -> Manifest { vars: std::collections::BTreeMap::new(), suppress: crate::manifest::ManifestSuppress::default(), workflows: Vec::new(), + interceptors: Vec::new(), } } diff --git a/crates/picloud-cli/src/cmds/plan.rs b/crates/picloud-cli/src/cmds/plan.rs index 44a8738..f42bddf 100644 --- a/crates/picloud-cli/src/cmds/plan.rs +++ b/crates/picloud-cli/src/cmds/plan.rs @@ -106,7 +106,7 @@ fn render_tree(plan: &TreePlanDto, mode: OutputMode) { ]); } } - let groups: [(&str, &Vec); 9] = [ + let groups: [(&str, &Vec); 10] = [ ("script", &n.scripts), ("route", &n.routes), ("trigger", &n.triggers), @@ -116,6 +116,7 @@ fn render_tree(plan: &TreePlanDto, mode: OutputMode) { ("collection", &n.collections), ("suppression", &n.suppressions), ("workflow", &n.workflows), + ("interceptor", &n.interceptors), ]; for (rk, changes) in groups { for c in changes { @@ -252,6 +253,17 @@ pub fn build_bundle(manifest: &Manifest, base_dir: &Path) -> Result { .iter() .map(workflow_to_wire) .collect::>>()?, + // §9.4: expand each `[[interceptors]]` entry's `ops` into one wire + // marker per (service, op). + "interceptors": manifest + .interceptors + .iter() + .flat_map(|i| { + i.ops.iter().map(move |op| { + json!({ "service": i.service, "op": op, "script": i.script }) + }) + }) + .collect::>(), })) } @@ -295,7 +307,7 @@ fn render(plan: &PlanDto, mode: OutputMode) { ]); } } - let groups: [(&str, &Vec); 9] = [ + let groups: [(&str, &Vec); 10] = [ ("script", &plan.scripts), ("route", &plan.routes), ("trigger", &plan.triggers), @@ -305,6 +317,7 @@ fn render(plan: &PlanDto, mode: OutputMode) { ("collection", &plan.collections), ("suppression", &plan.suppressions), ("workflow", &plan.workflows), + ("interceptor", &plan.interceptors), ]; for (kind, changes) in groups { for c in changes { diff --git a/crates/picloud-cli/src/cmds/pull.rs b/crates/picloud-cli/src/cmds/pull.rs index ba75f4d..4b4f8ea 100644 --- a/crates/picloud-cli/src/cmds/pull.rs +++ b/crates/picloud-cli/src/cmds/pull.rs @@ -268,6 +268,9 @@ pub async fn run(app_ident: &str, dir: &Path, force: bool, mode: OutputMode) -> vars: manifest_vars, suppress: crate::manifest::ManifestSuppress::default(), workflows: workflows.iter().map(wire_workflow_to_manifest).collect(), + // `pull` does not round-trip interceptor markers yet (read surface TBD); + // an interceptor authored in the manifest survives re-apply regardless. + interceptors: Vec::new(), }; std::fs::write(&manifest_path, manifest.to_toml()?) diff --git a/crates/picloud-cli/src/manifest.rs b/crates/picloud-cli/src/manifest.rs index aab8669..aff9b28 100644 --- a/crates/picloud-cli/src/manifest.rs +++ b/crates/picloud-cli/src/manifest.rs @@ -64,6 +64,29 @@ pub struct Manifest { /// App-owned; a `[group]` carrying them is rejected server-side. #[serde(default, skip_serializing_if = "Vec::is_empty")] pub workflows: Vec, + /// `[[interceptors]]` (§9.4) — before-op allow/deny hooks. Each names a + /// `script` and the `ops` it guards on a `service` (MVP: `service = "kv"`, + /// `ops ⊆ [set, delete]`). Authored on an app OR group node. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub interceptors: Vec, +} + +/// One `[[interceptors]]` entry — an interceptor script and the operations it +/// guards. `ops` expands to one wire marker per `(service, op)` at apply. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ManifestInterceptor { + /// Name of the interceptor script (resolved on the node's chain at run time). + pub script: String, + /// The guarded service. MVP: only `"kv"`. + #[serde(default = "default_interceptor_service")] + pub service: String, + /// The guarded operations, e.g. `["set", "delete"]`. + pub ops: Vec, +} + +fn default_interceptor_service() -> String { + "kv".to_string() } impl Manifest { @@ -783,6 +806,7 @@ mod tests { ]), suppress: ManifestSuppress::default(), workflows: Vec::new(), + interceptors: Vec::new(), } } @@ -890,6 +914,7 @@ mod tests { vars: BTreeMap::new(), suppress: ManifestSuppress::default(), workflows: Vec::new(), + interceptors: Vec::new(), }; let text = m.to_toml().unwrap(); assert!(!text.contains("[[scripts]]"), "got:\n{text}"); diff --git a/crates/picloud-cli/tests/cli.rs b/crates/picloud-cli/tests/cli.rs index 7a2d017..6cc50cd 100644 --- a/crates/picloud-cli/tests/cli.rs +++ b/crates/picloud-cli/tests/cli.rs @@ -35,6 +35,7 @@ mod group_secrets; mod group_triggers; mod groups; mod init; +mod interceptors; mod invoke; mod logs; mod output; diff --git a/crates/picloud-cli/tests/interceptors.rs b/crates/picloud-cli/tests/interceptors.rs new file mode 100644 index 0000000..e0ce498 --- /dev/null +++ b/crates/picloud-cli/tests/interceptors.rs @@ -0,0 +1,193 @@ +//! §9.4 Service Interceptors, end to end via `pic` — the KV allow/deny slice. +//! +//! A `[[interceptors]]` marker binds a script to run BEFORE `kv::set`/`delete`. +//! The interceptor reads the operation context (`ctx.request.body` — service, +//! action, collection, key, value) and returns `#{ allowed: bool, reason }`; a +//! `false` denies the op (the write never happens, the caller gets an error). +//! Registration is nearest-owner-wins on the calling app's chain, so a group's +//! interceptor is inherited by a descendant app (and an app can override it). + +use std::fs; + +use tempfile::TempDir; + +use crate::common; +use crate::common::cleanup::{AppGuard, GroupGuard}; + +fn manifest_dir() -> TempDir { + let dir = TempDir::new().expect("tempdir"); + fs::create_dir_all(dir.path().join("scripts")).expect("scripts dir"); + dir +} + +fn app_script_id(env: &common::TestEnv, app: &str, name: &str) -> String { + let ls = common::pic_as(env) + .args(["scripts", "ls", "--app", app]) + .output() + .expect("scripts ls"); + let table = String::from_utf8(ls.stdout).unwrap(); + table + .lines() + .map(common::cells) + .find(|c| c.get(2) == Some(&name)) + .and_then(|c| c.first().map(|s| (*s).to_string())) + .unwrap_or_else(|| panic!("script `{name}` not found:\n{table}")) +} + +fn invoke_body(env: &common::TestEnv, id: &str) -> serde_json::Value { + let out = common::pic_as(env) + .args(["scripts", "invoke", id]) + .output() + .expect("scripts invoke"); + assert!( + out.status.success(), + "invoke failed: {}", + String::from_utf8_lossy(&out.stderr) + ); + serde_json::from_slice(&out.stdout).expect("invoke body is JSON") +} + +/// A guard that denies a `kv::set` of the key `"secret"` and allows everything +/// else, reading the op context from `ctx.request.body`. +const GUARD: &str = r#" +let op = ctx.request.body; +if op.action == "set" && op.key == "secret" { + #{ allowed: false, reason: "the `secret` key is protected" } +} else { + #{ allowed: true } +} +"#; + +/// A writer that tries a denied set (caught), then a permitted set. +const WRITER: &str = r#" +let denied = false; +try { kv::collection("c").set("secret", 1); } catch(e) { denied = true; } +kv::collection("c").set("ok", 2); +#{ denied: denied } +"#; + +#[ignore = "needs DATABASE_URL pointing at a running Postgres"] +#[test] +fn app_interceptor_denies_a_guarded_kv_write_and_allows_others() { + let Some(fx) = common::fixture_or_skip() else { + return; + }; + let env = common::admin_env(fx); + let app = common::unique_slug("ic-app"); + let _a = AppGuard::new(&env.url, &env.token, &app); + common::pic_as(&env) + .args(["apps", "create", &app]) + .assert() + .success(); + + // One apply deploys the guard + writer scripts AND registers the marker. + let dir = manifest_dir(); + fs::write(dir.path().join("scripts/guard.rhai"), GUARD).unwrap(); + fs::write(dir.path().join("scripts/writer.rhai"), WRITER).unwrap(); + fs::write( + dir.path().join("picloud.toml"), + format!( + "[app]\nslug = \"{app}\"\nname = \"IC\"\n\n\ + [[scripts]]\nname = \"guard\"\nfile = \"scripts/guard.rhai\"\n\n\ + [[scripts]]\nname = \"writer\"\nfile = \"scripts/writer.rhai\"\n\n\ + [[interceptors]]\nscript = \"guard\"\nops = [\"set\"]\n" + ), + ) + .unwrap(); + common::pic_as(&env) + .args(["apply", "--file"]) + .arg(dir.path().join("picloud.toml")) + .assert() + .success(); + + // The writer's guarded set is denied (caught), the permitted set goes through. + let body = invoke_body(&env, &app_script_id(&env, &app, "writer")); + assert_eq!( + body, + serde_json::json!({ "denied": true }), + "the `secret` set must be denied by the interceptor" + ); + + // `ok` was written; `secret` was not (the write never ran). + let ok = String::from_utf8( + common::pic_as(&env) + .args(["kv", "get", "--app", &app, "--collection", "c", "ok"]) + .output() + .unwrap() + .stdout, + ) + .unwrap(); + assert!(ok.contains('2'), "the allowed write must persist:\n{ok}"); + let secret = common::pic_as(&env) + .args(["kv", "get", "--app", &app, "--collection", "c", "secret"]) + .output() + .unwrap(); + let secret_out = String::from_utf8_lossy(&secret.stdout); + assert!( + secret_out.trim() == "null" || secret_out.trim() == "()" || secret_out.trim().is_empty(), + "the denied write must NOT persist, got: {secret_out}" + ); +} + +#[ignore = "needs DATABASE_URL pointing at a running Postgres"] +#[test] +fn a_group_interceptor_is_inherited_by_a_descendant_app() { + let Some(fx) = common::fixture_or_skip() else { + return; + }; + let env = common::admin_env(fx); + let group = common::unique_slug("icg-grp"); + let app = common::unique_slug("icg-app"); + let _g = GroupGuard::new(&env.url, &env.token, &group); + common::pic_as(&env) + .args(["groups", "create", &group]) + .assert() + .success(); + + // The GROUP owns the guard script + the interceptor marker. + let dir = manifest_dir(); + fs::write(dir.path().join("scripts/guard.rhai"), GUARD).unwrap(); + fs::write( + dir.path().join("group.toml"), + format!( + "[group]\nslug = \"{group}\"\nname = \"ICG\"\n\n\ + [[scripts]]\nname = \"guard\"\nfile = \"scripts/guard.rhai\"\n\n\ + [[interceptors]]\nscript = \"guard\"\nops = [\"set\"]\n" + ), + ) + .unwrap(); + common::pic_as(&env) + .args(["apply", "--file"]) + .arg(dir.path().join("group.toml")) + .assert() + .success(); + + // An app UNDER the group inherits the interceptor (nearest-owner-wins), even + // though the marker + guard live on the group. + let _a = AppGuard::new(&env.url, &env.token, &app); + common::pic_as(&env) + .args(["apps", "create", &app, "--group", &group]) + .assert() + .success(); + fs::write(dir.path().join("scripts/writer.rhai"), WRITER).unwrap(); + fs::write( + dir.path().join("app.toml"), + format!( + "[app]\nslug = \"{app}\"\nname = \"ICApp\"\n\n\ + [[scripts]]\nname = \"writer\"\nfile = \"scripts/writer.rhai\"\n" + ), + ) + .unwrap(); + common::pic_as(&env) + .args(["apply", "--file"]) + .arg(dir.path().join("app.toml")) + .assert() + .success(); + + let body = invoke_body(&env, &app_script_id(&env, &app, "writer")); + assert_eq!( + body, + serde_json::json!({ "denied": true }), + "a descendant app must inherit the group's interceptor" + ); +} diff --git a/crates/picloud/src/lib.rs b/crates/picloud/src/lib.rs index 9413267..fe9a114 100644 --- a/crates/picloud/src/lib.rs +++ b/crates/picloud/src/lib.rs @@ -416,6 +416,11 @@ pub async fn build_app( let workflow: Arc = Arc::new( picloud_manager_core::WorkflowServiceImpl::new(pool.clone()).with_authz(authz.clone()), ); + // §9.4 interceptors: resolve-only (which script guards a kv op); running it + // reuses the invoke() re-entry path. + let interceptors: Arc = Arc::new( + picloud_manager_core::interceptor_service::InterceptorServiceImpl::new(pool.clone()), + ); let services = Services::new( kv, docs, @@ -437,7 +442,8 @@ pub async fn build_app( .with_group_files(group_files) .with_group_pubsub(group_pubsub) .with_group_queue(group_queue) - .with_workflow(workflow); + .with_workflow(workflow) + .with_interceptors(interceptors); // 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 { diff --git a/crates/shared/src/interceptor.rs b/crates/shared/src/interceptor.rs new file mode 100644 index 0000000..3486332 --- /dev/null +++ b/crates/shared/src/interceptor.rs @@ -0,0 +1,51 @@ +//! §9.4 Service Interceptors — the injected resolver seam. +//! +//! An interceptor is a script registered (declaratively, per app/group) to run +//! **before** a data-plane operation and allow or deny it. This trait is the +//! narrow, executor-facing seam: it only RESOLVES which interceptor script (by +//! name) applies to a `(service, op)` on the calling app's chain — nearest-owner +//! wins, like extension points. Running the resolved script reuses the existing +//! `invoke()` re-entry path in `executor-core` (resolve the name → compile → +//! execute), so `executor-core` stays Postgres-free and there is no second +//! script-dispatch mechanism. +//! +//! MVP scope (v1.2): `service = "kv"`, `op ∈ {set, delete}`, allow/deny only — +//! no data transform, no chaining, no `after_*` hooks. See blueprint §9.4. + +use async_trait::async_trait; + +use crate::SdkCallCx; + +#[async_trait] +pub trait InterceptorService: Send + Sync { + /// The NAME of the interceptor script registered for `(service, op)` at the + /// nearest owner on `cx.app_id`'s chain (app, else nearest ancestor group), + /// or `None` when the operation is un-hooked. The executor resolves that + /// name to a script and runs it. `Err` is a backend failure (fail-closed: + /// the caller turns it into an operation error rather than silently + /// allowing). + async fn resolve_before( + &self, + cx: &SdkCallCx, + service: &str, + op: &str, + ) -> Result, String>; +} + +/// Default: nothing is ever intercepted. The shape every non-picloud `Services` +/// (tests, cluster skeletons) gets for free — an un-hooked write pays exactly +/// one `Ok(None)` here, no I/O. +#[derive(Debug, Default, Clone, Copy)] +pub struct NoopInterceptorService; + +#[async_trait] +impl InterceptorService for NoopInterceptorService { + async fn resolve_before( + &self, + _cx: &SdkCallCx, + _service: &str, + _op: &str, + ) -> Result, String> { + Ok(None) + } +} diff --git a/crates/shared/src/lib.rs b/crates/shared/src/lib.rs index 4845fcf..ba9b967 100644 --- a/crates/shared/src/lib.rs +++ b/crates/shared/src/lib.rs @@ -24,6 +24,7 @@ pub mod group_queue; pub mod http; pub mod ids; pub mod inbox; +pub mod interceptor; pub mod invoke; pub mod kv; pub mod log_sink; @@ -86,6 +87,7 @@ pub use ids::{ pub use inbox::{ InboxDeliveryOutcome, InboxFailureKind, InboxResolver, InboxResult, NoopInboxResolver, }; +pub use interceptor::{InterceptorService, NoopInterceptorService}; pub use invoke::{InvokeError, InvokeService, InvokeTarget, NoopInvokeService, ResolvedScript}; pub use kv::{KvError, KvListPage, KvService, NoopKvService}; pub use log_sink::{ExecutionLogSink, LogSinkError}; diff --git a/crates/shared/src/services.rs b/crates/shared/src/services.rs index e66de8b..2c94164 100644 --- a/crates/shared/src/services.rs +++ b/crates/shared/src/services.rs @@ -22,13 +22,13 @@ use std::sync::Arc; use crate::{ DeadLetterService, DocsService, EmailService, FilesService, GroupDocsService, GroupFilesService, GroupKvService, GroupPubsubService, GroupQueueService, HttpService, - InvokeService, KvService, ModuleSource, NoopDeadLetterService, NoopDocsService, - NoopEmailService, NoopEventEmitter, NoopFilesService, NoopGroupDocsService, + InterceptorService, InvokeService, KvService, ModuleSource, NoopDeadLetterService, + NoopDocsService, NoopEmailService, NoopEventEmitter, NoopFilesService, NoopGroupDocsService, NoopGroupFilesService, NoopGroupKvService, NoopGroupPubsubService, NoopGroupQueueService, - NoopHttpService, NoopInvokeService, NoopKvService, NoopModuleSource, NoopPubsubService, - NoopQueueService, NoopSecretsService, NoopUsersService, NoopVarsService, NoopWorkflowService, - PubsubService, QueueService, SecretsService, ServiceEventEmitter, UsersService, VarsService, - WorkflowService, + NoopHttpService, NoopInterceptorService, NoopInvokeService, NoopKvService, NoopModuleSource, + NoopPubsubService, NoopQueueService, NoopSecretsService, NoopUsersService, NoopVarsService, + NoopWorkflowService, PubsubService, QueueService, SecretsService, ServiceEventEmitter, + UsersService, VarsService, WorkflowService, }; /// SDK service bundle. See module docs for the lifecycle and the v1.1.x @@ -152,6 +152,12 @@ pub struct Services { /// run of a named workflow in the caller's app. Wired via /// [`Services::with_workflow`]; defaults to `NoopWorkflowService`. pub workflow: Arc, + + /// §9.4 Service Interceptors — resolves which (if any) interceptor script + /// guards a `(service, op)` before it runs. Wired via + /// [`Services::with_interceptors`]; defaults to `NoopInterceptorService` + /// (nothing intercepted). + pub interceptors: Arc, } impl Services { @@ -200,9 +206,18 @@ impl Services { group_pubsub: Arc::new(NoopGroupPubsubService), group_queue: Arc::new(NoopGroupQueueService), workflow: Arc::new(NoopWorkflowService), + interceptors: Arc::new(NoopInterceptorService), } } + /// Set the §9.4 interceptor resolver (picloud binary wires the + /// Postgres-backed impl; tests leave the noop default = nothing hooked). + #[must_use] + pub fn with_interceptors(mut self, interceptors: Arc) -> Self { + self.interceptors = interceptors; + self + } + /// Set the v1.2 Workflows service (picloud binary wires the Postgres-backed /// impl; tests leave the noop default). #[must_use] diff --git a/serverless_cloud_blueprint.md b/serverless_cloud_blueprint.md index b58cb52..f4d6cad 100644 --- a/serverless_cloud_blueprint.md +++ b/serverless_cloud_blueprint.md @@ -1,6 +1,6 @@ # Project Blueprint: Lightweight Event-Based Serverless Cloud -**Status**: v1.1 shipped (SDK + services) · v1.2 *Hierarchies* track **complete** · v1.2 *Workflows* track **shipped** (M1–M6: DAG execution, conditional branching, nested sub-workflows, `workflow::start` SDK, dashboard DAG + run-history); §9.4 service interceptors + v1.3 cluster mode are next +**Status**: v1.1 shipped (SDK + services) · v1.2 *Hierarchies* track **complete** · v1.2 *Workflows* track **shipped** (M1–M6: DAG execution, conditional branching, nested sub-workflows, `workflow::start` SDK, dashboard DAG + run-history); a §9.4 service-interceptor KV allow/deny slice has shipped (migration 0073); the rest of §9.4 + v1.3 cluster mode are next **Last Updated**: 2026-07-12 (reconciled to shipped code — CLAUDE.md is the live source of truth) **Audience**: Solo developer (DIY self-hosted) @@ -171,8 +171,8 @@ rhai_executor --script $SCRIPT_PATH --request "$REQUEST_JSON" ### 3.4 PostgreSQL Database **Schema (MVP sketch — NOT authoritative):** the block below is the original MVP shape. The **authoritative -schema is the migration set** in `crates/manager-core/migrations/` (through `0072` as of v1.2, incl. the -Workflows tables), which has +schema is the migration set** in `crates/manager-core/migrations/` (through `0073` as of v1.2, incl. the +Workflows tables + the §9.4 interceptor markers), which has since added apps/domains, RBAC (`admin_users`/`app_members`/`api_keys`), the v1.1 data-plane services (KV/docs/files/queues/…), and the v1.2 groups/collections/templates/projects tables. Treat this as illustration only. @@ -1697,6 +1697,16 @@ CREATE INDEX idx_execution_parent ON execution_logs(parent_execution_id); ### 9.4 Service Interceptors & Middleware (v1.2+) +> **Status — thin KV allow/deny slice SHIPPED** (migration `0073_interceptors.sql`). A `[[interceptors]]` +> manifest block (app OR group) binds a script to run BEFORE `kv::set` / `kv::delete`; it reads the operation +> context (`ctx.request.body`: service, action, collection, key, value, caller ids) and returns +> `#{ allowed, reason }` — `allowed == false` denies the op (the write never runs). Registration is a marker +> `(owner, service, op) → script`, resolved **nearest-owner-wins** on the calling app's chain (an app overrides +> a group's), mirroring extension points (§5.5); the interceptor script itself is resolved + run through the +> `invoke()` re-entry path (shared depth bound, AST cache). **Deferred (the rest of the spec below):** the +> `data` transform return, services other than `kv`, `after_*` hooks, interceptor chaining + circular-dependency +> guard, and the timeout policy. The remaining subsections describe that full design. + **Concept**: A script can act as middleware to intercept and validate/transform service operations before they execute. **Use Cases:**