diff --git a/app/common/src/main/java/stirling/software/common/model/ApplicationProperties.java b/app/common/src/main/java/stirling/software/common/model/ApplicationProperties.java index d2f0472510..3807e57637 100644 --- a/app/common/src/main/java/stirling/software/common/model/ApplicationProperties.java +++ b/app/common/src/main/java/stirling/software/common/model/ApplicationProperties.java @@ -254,6 +254,13 @@ public class ApplicationProperties { * in-network object store. */ private boolean allowPrivateS3Endpoints = false; + + /** + * Largest inbound document (bytes) the webhook source receiver will accept in one delivery. + * A larger body is rejected with 413 before anything is written to the spool, so an + * unauthenticated caller cannot fill the disk. Default 100 MB. + */ + private long webhookMaxBytes = 104857600L; } @Data diff --git a/app/common/src/main/java/stirling/software/common/util/RequestUriUtils.java b/app/common/src/main/java/stirling/software/common/util/RequestUriUtils.java index 4db9f118ec..9a116a4b25 100644 --- a/app/common/src/main/java/stirling/software/common/util/RequestUriUtils.java +++ b/app/common/src/main/java/stirling/software/common/util/RequestUriUtils.java @@ -198,6 +198,9 @@ public class RequestUriUtils { || trimmedUri.startsWith("/readiness") || trimmedUri.startsWith( "/api/v1/mobile-scanner/") // Mobile scanner endpoints (no auth) + // Policy webhook source receiver - authenticated per-request by its HMAC signature, + // not a login session, so it must bypass the session auth wall. + || trimmedUri.startsWith("/api/v1/webhooks/") || trimmedUri.startsWith("/v1/api-docs") // Workflow participant endpoints - access controlled by share tokens, not login || trimmedUri.startsWith("/api/v1/workflow/participant/") diff --git a/app/common/src/test/java/stirling/software/common/util/RequestUriUtilsTest.java b/app/common/src/test/java/stirling/software/common/util/RequestUriUtilsTest.java index a6fefd1091..1912f3808b 100644 --- a/app/common/src/test/java/stirling/software/common/util/RequestUriUtilsTest.java +++ b/app/common/src/test/java/stirling/software/common/util/RequestUriUtilsTest.java @@ -176,6 +176,13 @@ class RequestUriUtilsTest { assertFalse(RequestUriUtils.isPublicAuthEndpoint("/api/v1/convert", "")); } + @Test + void testIsPublicAuthEndpoint_webhookReceiver() { + // The webhook source receiver authenticates each delivery by HMAC signature, not a session. + assertTrue(RequestUriUtils.isPublicAuthEndpoint("/api/v1/webhooks/whk_abc123", "")); + assertTrue(RequestUriUtils.isPublicAuthEndpoint("/app/api/v1/webhooks/whk_abc123", "/app")); + } + @Test void testIsPublicAuthEndpoint_withContextPath() { assertTrue(RequestUriUtils.isPublicAuthEndpoint("/app/login", "/app")); diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/input/InputSource.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/input/InputSource.java index d32c2fc546..aa3863fc1d 100644 --- a/app/proprietary/src/main/java/stirling/software/proprietary/policy/input/InputSource.java +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/input/InputSource.java @@ -3,6 +3,7 @@ package stirling.software.proprietary.policy.input; import java.io.IOException; import java.nio.file.Path; import java.util.List; +import java.util.Map; import stirling.software.proprietary.policy.model.InputSpec; @@ -22,6 +23,18 @@ public interface InputSource { /** Throws {@link IllegalArgumentException} on bad config. Called on save to fail fast. */ default void validate(InputSpec spec) {} + /** + * Normalise a source's options just before it is persisted, so a source type can populate + * server-owned config the client neither supplies nor controls. {@code isCreate} is true only + * for a brand-new source (blank id). Runs before {@link #validate}. The default returns the + * options unchanged; the webhook source overrides it to mint its routing id and signing secret + * on create. Must not mutate the argument. + */ + default Map prepareOptionsForSave( + Map options, boolean isCreate) { + return options; + } + /** * Resolve the spec into zero or more units of work, each carrying one run's files and a * completion hook. Empty list means nothing to run right now. Discovery is read-only - files diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/input/WebhookInputSource.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/input/WebhookInputSource.java new file mode 100644 index 0000000000..73887ad9cc --- /dev/null +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/input/WebhookInputSource.java @@ -0,0 +1,190 @@ +package stirling.software.proprietary.policy.input; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Stream; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty; +import org.springframework.core.io.FileSystemResource; +import org.springframework.core.io.Resource; +import org.springframework.stereotype.Service; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; + +import stirling.software.common.util.FileReadinessChecker; +import stirling.software.proprietary.policy.ledger.FolderIdentities; +import stirling.software.proprietary.policy.model.InputSpec; +import stirling.software.proprietary.policy.model.PolicyInputs; +import stirling.software.proprietary.policy.webhook.WebhookConfig; +import stirling.software.proprietary.policy.webhook.WebhookIds; +import stirling.software.proprietary.policy.webhook.WebhookSpool; + +/** + * A push source: documents arrive by signed HTTP POST to {@code /api/v1/webhooks/{webhookId}} and + * are staged by {@link WebhookSpool}. This source then reads that per-webhook spool exactly as + * {@link FolderInputSource} reads a watched directory - each spooled file is a unit of work claimed + * through the {@link ResolveContext} ledger and tracked in place - so the whole claim/consume model + * is reused rather than reinvented. Options (see {@link WebhookConfig}): "mode" is "consume" + * (default: a spooled file is removed once every policy that claimed it has settled successfully) + * or "snapshot" (stateless, every run re-reads whatever is currently spooled). "webhookId" and + * "signingSecret" are minted server-side on create, never supplied by the client. An empty or + * not-yet-created spool means nothing has been delivered, which reads as "verifiably no files". + */ +@Slf4j +@Service +@RequiredArgsConstructor +@ConditionalOnBooleanProperty(name = "policies.enabled") +public class WebhookInputSource implements InputSource { + + static final String TYPE = "webhook"; + + private final WebhookSpool spool; + private final FileReadinessChecker readinessChecker; + + @Override + public String type() { + return TYPE; + } + + @Override + public boolean supports(InputSpec spec) { + return spec != null && TYPE.equals(spec.type()); + } + + @Override + public void validate(InputSpec spec) { + WebhookConfig.from(spec.options()); + } + + /** + * Mint the routing id and signing secret on create, so the delivery URL and its HMAC key are + * always server-generated. An edit (or any options that already carry a webhookId) is left + * untouched, so the URL a sender is already configured against never changes underneath them. + */ + @Override + public Map prepareOptionsForSave( + Map options, boolean isCreate) { + boolean hasId = + options.get(WebhookConfig.WEBHOOK_ID_OPTION) != null + && !options.get(WebhookConfig.WEBHOOK_ID_OPTION).toString().isBlank(); + if (!isCreate && hasId) { + return options; + } + Map prepared = new LinkedHashMap<>(options); + if (!hasId) { + prepared.put(WebhookConfig.WEBHOOK_ID_OPTION, WebhookIds.newWebhookId()); + } + Object secret = prepared.get(WebhookConfig.SIGNING_SECRET_OPTION); + if (secret == null || secret.toString().isBlank()) { + prepared.put(WebhookConfig.SIGNING_SECRET_OPTION, WebhookIds.newSigningSecret()); + } + return prepared; + } + + @Override + public List resolve(InputSpec spec, ResolveContext ctx) throws IOException { + WebhookConfig config = WebhookConfig.from(spec.options()); + Path dir = spool.dirFor(config.webhookId()); + if (!Files.isDirectory(dir)) { + // No deliveries yet (or none since the last consume): a verifiably empty source, so the + // sweep may prune ledger rows for files that are gone. Unlike the folder source a + // missing directory is normal here, not an unmounted-drive error. + ctx.reportPresent(List.of()); + return List.of(); + } + Path canonicalDir = FolderIdentities.canonicalDir(dir); + List present = listFiles(dir); + + if (config.snapshot()) { + List work = new ArrayList<>(); + for (Path file : present) { + if (readinessChecker.isReady(file)) { + work.add(ResolvedInput.of(PolicyInputs.of(List.of(fileResource(file))))); + } + } + return work; + } + + ctx.reportPresent( + present.stream() + .map(file -> FolderIdentities.identity(canonicalDir, dir, file)) + .toList()); + + List work = new ArrayList<>(); + for (Path file : present) { + if (!readinessChecker.isReady(file)) { + continue; + } + String identity = FolderIdentities.identity(canonicalDir, dir, file); + String gate; + boolean claimed; + try { + gate = FolderIdentities.statGate(file); + claimed = ctx.claim(identity, gate, null); + } catch (IOException | UncheckedIOException e) { + log.debug("Could not read {} for its version: {}", file, e.getMessage()); + continue; // vanished or unreadable mid-sweep; the next sweep sees the truth + } + if (!claimed) { + continue; + } + work.add( + new ResolvedInput( + PolicyInputs.of(List.of(fileResource(file))), + success -> completeConsumed(ctx, identity, file, gate, success))); + } + return work; + } + + /** + * Settle at the claimed version, then remove the spooled file only when it is still that + * version and every policy that claimed it has settled DONE - the same consensus delete the + * folder and S3 sources use, so a shared webhook feeding several policies keeps a delivery + * until all are done and one failure parks it for everyone. A failed run settles ERROR and + * never deletes. + */ + private static void completeConsumed( + ResolveContext ctx, String identity, Path file, String claimGate, boolean success) { + ctx.settle(identity, claimGate, null, success); + if (!success) { + return; + } + try { + if (FolderIdentities.statGate(file).equals(claimGate) && ctx.allSettledDone(identity)) { + Files.deleteIfExists(file); + } + } catch (java.nio.file.NoSuchFileException alreadyGone) { + // Removed by the user or a co-watching policy's own consensus delete: nothing to do. + } catch (IOException e) { + log.warn("Could not remove consumed webhook delivery {}: {}", file, e.getMessage()); + } + } + + /** Every non-hidden regular file currently spooled for the webhook. */ + private static List listFiles(Path dir) throws IOException { + List files = new ArrayList<>(); + try (Stream entries = Files.list(dir)) { + entries.filter(Files::isRegularFile) + .filter(file -> !file.getFileName().toString().startsWith(".")) + .forEach(files::add); + } + return files; + } + + private static Resource fileResource(Path path) { + String name = WebhookSpool.displayName(path.getFileName().toString()); + return new FileSystemResource(path.toFile()) { + @Override + public String getFilename() { + return name; + } + }; + } +} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceController.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceController.java index 524b8c369c..1011cb4e71 100644 --- a/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceController.java +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceController.java @@ -2,6 +2,7 @@ package stirling.software.proprietary.policy.source; import java.util.List; import java.util.Map; +import java.util.Optional; import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty; import org.springframework.http.HttpStatus; @@ -46,6 +47,8 @@ import stirling.software.proprietary.util.SecretMasker; @ConditionalOnBooleanProperty(name = "policies.enabled") public class SourceController { + private static final String WEBHOOK_TYPE = "webhook"; + private final SourceStore sourceStore; private final SourceAccessGuard sourceAccessGuard; private final SourceOverviewService overviewService; @@ -109,7 +112,8 @@ public class SourceController { public ResponseEntity save(@RequestBody Source source) { requireSourceEditingAllowed(); requireNotEditor(source.id(), source.type()); - Source owned = withStoredSecrets(resolveOwnership(source)); + boolean isCreate = source.id() == null || source.id().isBlank(); + Source owned = withPreparedOptions(withStoredSecrets(resolveOwnership(source)), isCreate); try { validateConfig(owned); } catch (IllegalArgumentException e) { @@ -119,7 +123,7 @@ public class SourceController { // An edited folder source can change which directory needs watching, so re-sync trigger // registrations now instead of waiting for the next reconcile. policyTriggerManager.notifyPoliciesChanged(); - return ResponseEntity.ok(withMaskedSecrets(saved)); + return ResponseEntity.ok(revealOnCreate(saved, isCreate)); } @DeleteMapping("/{sourceId}") @@ -221,14 +225,45 @@ public class SourceController { /** Validate the config against the bean that handles the source's type, as the engine will. */ private void validateConfig(Source source) { InputSpec spec = source.toInputSpec(); - inputSources.stream() - .filter(inputSource -> inputSource.supports(spec)) - .findFirst() + inputSourceFor(spec) .orElseThrow( () -> new IllegalArgumentException("unknown source type: " + source.type())) .validate(spec); } + /** + * Let the matching source type populate server-owned config just before persistence (a webhook + * mints its routing id and signing secret on create). An unknown type is left untouched; {@link + * #validateConfig} then rejects it with a clear message. + */ + private Source withPreparedOptions(Source source, boolean isCreate) { + InputSpec spec = source.toInputSpec(); + InputSource input = inputSourceFor(spec).orElse(null); + if (input == null) { + return source; + } + Map prepared = input.prepareOptionsForSave(source.options(), isCreate); + // A source type must return the (possibly augmented) options; if it returns null, keep the + // originals rather than let the Source record normalise null to an empty map and wipe + // config. + return prepared == null ? source : withOptions(source, prepared); + } + + /** + * Reveal server-minted secrets (a webhook's signing secret) once, on the create response only, + * so the operator can copy them; every other read - including an edit - is masked. + */ + private static Source revealOnCreate(Source saved, boolean isCreate) { + if (isCreate && WEBHOOK_TYPE.equals(saved.type())) { + return saved; + } + return withMaskedSecrets(saved); + } + + private Optional inputSourceFor(InputSpec spec) { + return inputSources.stream().filter(input -> input.supports(spec)).findFirst(); + } + /** * Editing sources requires the editor role for the caller's team (a team leader on SaaS), the * same rule as policies. Single-user deployments (login disabled) trust the local operator. diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceOverviewService.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceOverviewService.java index f202a30e23..9ab4d2d0d4 100644 --- a/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceOverviewService.java +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceOverviewService.java @@ -102,7 +102,8 @@ public class SourceOverviewService { List.of(), docs.total(), docs.last24h(), - docs.last30d()); + docs.last30d(), + null); } /** @@ -143,7 +144,17 @@ public class SourceOverviewService { configRows(source), docs.total(), docs.last24h(), - docs.last30d()); + docs.last30d(), + webhookPath(source)); + } + + /** The server-relative delivery path for a webhook source, else null. Never a secret. */ + private static String webhookPath(Source source) { + if (!"webhook".equals(source.type())) { + return null; + } + Object webhookId = source.options().get("webhookId"); + return webhookId == null ? null : "/api/v1/webhooks/" + webhookId; } /** A disabled (paused) source reads as "disabled"; an unreferenced one reads as "unused". */ diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceView.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceView.java index 6c6a9c4c6d..a043518f2c 100644 --- a/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceView.java +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/source/SourceView.java @@ -5,7 +5,9 @@ import java.util.List; /** * One row in the Sources overview: a persisted input connection shown exactly once, with how many * policies reference it (and which) and how many documents it has fed into runs ({@code docsTotal} - * lifetime plus the trailing 24-hour and 30-day windows). + * lifetime plus the trailing 24-hour and 30-day windows). {@code webhookPath} is set only for a + * webhook source - the server-relative delivery path senders POST to - and is null for every other + * type; the client turns it into an absolute URL. */ public record SourceView( String id, @@ -17,7 +19,8 @@ public record SourceView( List config, long docsTotal, long docs24h, - long docs30d) { + long docs30d, + String webhookPath) { /** A policy that references this source. */ public record PolicyRef(String id, String name) {} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/trigger/WebhookTrigger.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/trigger/WebhookTrigger.java new file mode 100644 index 0000000000..304ecaaede --- /dev/null +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/trigger/WebhookTrigger.java @@ -0,0 +1,152 @@ +package stirling.software.proprietary.policy.trigger; + +import java.util.Set; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty; +import org.springframework.stereotype.Service; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; + +import stirling.software.common.model.ApplicationProperties; +import stirling.software.proprietary.policy.engine.PolicyRunner; +import stirling.software.proprietary.policy.engine.SweepKind; +import stirling.software.proprietary.policy.model.Policy; +import stirling.software.proprietary.policy.source.Source; +import stirling.software.proprietary.policy.source.SourceStore; +import stirling.software.proprietary.policy.store.PolicyStore; +import stirling.software.proprietary.policy.webhook.WebhookConfig; + +/** + * Fires policies when a document is delivered to one of their webhook sources. Delivery is push, so + * the receiver calls {@link #fireForWebhook} the moment it spools a file and the referencing + * policies run immediately. As with folder-watch, the instant fire is a latency optimisation, not + * the source of truth: a periodic reconcile ({@code watchReconcileSeconds}) re-runs every webhook + * policy to cover deliveries that landed while the process was down and any fire that was missed. + * Redundant runs are harmless because the input source does the claiming. + * + *

The spool is on local disk, so - like folder-watch - this assumes a single node. + */ +@Slf4j +@Service +@RequiredArgsConstructor +@ConditionalOnBooleanProperty(name = "policies.enabled") +public class WebhookTrigger implements PolicyTrigger { + + static final String TYPE = "webhook"; + // The webhook source type; equal to the trigger type but a distinct concept (source vs + // trigger). + private static final String WEBHOOK_SOURCE_TYPE = "webhook"; + + private final PolicyStore policyStore; + private final PolicyRunner policyRunner; + private final SourceStore sourceStore; + private final ApplicationProperties applicationProperties; + + private volatile ScheduledExecutorService reconciler; + + @Override + public String type() { + return TYPE; + } + + @Override + public boolean requiresSource() { + return true; + } + + @Override + public Set supportedSourceTypes() { + return Set.of(WEBHOOK_SOURCE_TYPE); + } + + @Override + public void validate(Policy policy) { + boolean hasWebhookSource = + policy.sourceIds().stream() + .map(sourceStore::get) + .flatMap(java.util.Optional::stream) + .anyMatch(source -> WEBHOOK_SOURCE_TYPE.equals(source.type())); + if (!hasWebhookSource) { + throw new IllegalArgumentException( + "webhook trigger requires at least one webhook input source"); + } + } + + @Override + public synchronized void start() { + if (reconciler != null) { + return; + } + long reconcileSeconds = applicationProperties.getPolicies().getWatchReconcileSeconds(); + reconciler = + Executors.newSingleThreadScheduledExecutor( + Thread.ofVirtual().name("policy-webhook-reconcile-", 0).factory()); + // First reconcile runs immediately so deliveries spooled before startup are picked up. + reconciler.scheduleAtFixedRate(this::safeReconcile, 0, reconcileSeconds, TimeUnit.SECONDS); + log.info("Webhook trigger started (reconcile every {}s)", reconcileSeconds); + } + + @Override + public synchronized void stop() { + if (reconciler != null) { + reconciler.shutdownNow(); + reconciler = null; + } + } + + /** + * Run every webhook policy that references the delivered-to source. Called by the receiver + * right after it spools a document, so processing starts without waiting for the next + * reconcile. A LIGHT sweep: the periodic reconcile does the full-listing pass. Best-effort - a + * run failure is logged and the others still fire. + */ + public void fireForWebhook(String webhookId) { + for (Policy policy : policyStore.findByTriggerType(TYPE)) { + if (!referencesWebhook(policy, webhookId)) { + continue; + } + try { + log.debug("Webhook policy {} ({}) saw a delivery", policy.id(), policy.name()); + policyRunner.run(policy, SweepKind.LIGHT); + } catch (RuntimeException e) { + log.warn("Webhook run failed for policy {}: {}", policy.id(), e.getMessage()); + } + } + } + + private void safeReconcile() { + try { + for (Policy policy : policyStore.findByTriggerType(TYPE)) { + try { + policyRunner.run(policy); + } catch (RuntimeException e) { + log.warn( + "Webhook reconcile run failed for policy {}: {}", + policy.id(), + e.getMessage()); + } + } + } catch (RuntimeException e) { + log.error("Webhook reconcile failed: {}", e.getMessage(), e); + } + } + + /** Whether the policy references a webhook source whose routing id is {@code webhookId}. */ + private boolean referencesWebhook(Policy policy, String webhookId) { + for (String sourceId : policy.sourceIds()) { + Source source = sourceStore.get(sourceId).orElse(null); + if (source == null || !WEBHOOK_SOURCE_TYPE.equals(source.type())) { + continue; + } + Object configured = source.options().get(WebhookConfig.WEBHOOK_ID_OPTION); + if (configured != null && configured.toString().equals(webhookId)) { + return true; + } + } + return false; + } +} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookConfig.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookConfig.java new file mode 100644 index 0000000000..26d9521351 --- /dev/null +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookConfig.java @@ -0,0 +1,56 @@ +package stirling.software.proprietary.policy.webhook; + +import java.util.Map; + +/** + * Config for a webhook input source, parsed from a spec's options map. A webhook is a push source: + * external systems POST documents to a signed URL keyed by {@code webhookId}, and each delivery is + * verified against {@code signingSecret} before it is spooled for the referencing policies. Both + * {@code webhookId} (the public URL token) and {@code signingSecret} (the HMAC key) are generated + * server-side on create - a client never supplies them - so the receiver's identity and the + * sender's proof of authenticity are always Stirling's own. {@code mode} is "consume" (default: a + * spooled document is deleted once every policy that claimed it has settled successfully) or + * "snapshot" (stateless, every run re-reads whatever is currently spooled). + */ +public record WebhookConfig(String webhookId, String signingSecret, boolean snapshot) { + + public static final String WEBHOOK_ID_OPTION = "webhookId"; + public static final String SIGNING_SECRET_OPTION = "signingSecret"; + private static final String MODE_OPTION = "mode"; + private static final String MODE_CONSUME = "consume"; + private static final String MODE_SNAPSHOT = "snapshot"; + + public static WebhookConfig from(Map options) { + String webhookId = trimmed(options.get(WEBHOOK_ID_OPTION)); + if (webhookId == null) { + throw new IllegalArgumentException("webhook config requires a 'webhookId' option"); + } + if (!WebhookIds.isValidId(webhookId)) { + throw new IllegalArgumentException("webhook config 'webhookId' has an invalid format"); + } + String signingSecret = trimmed(options.get(SIGNING_SECRET_OPTION)); + if (signingSecret == null) { + throw new IllegalArgumentException("webhook config requires a 'signingSecret' option"); + } + String mode = trimmed(options.get(MODE_OPTION)); + if (mode != null && !MODE_CONSUME.equals(mode) && !MODE_SNAPSHOT.equals(mode)) { + throw new IllegalArgumentException( + "webhook config 'mode' must be 'consume' or 'snapshot'"); + } + return new WebhookConfig(webhookId, signingSecret, MODE_SNAPSHOT.equals(mode)); + } + + private static String trimmed(Object value) { + if (value == null) { + return null; + } + String text = value.toString().trim(); + return text.isEmpty() ? null : text; + } + + /** Never prints the signing secret, so an accidental log line cannot leak it. */ + @Override + public String toString() { + return "WebhookConfig[webhookId=" + webhookId + ", snapshot=" + snapshot + "]"; + } +} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookIds.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookIds.java new file mode 100644 index 0000000000..ed36fbba18 --- /dev/null +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookIds.java @@ -0,0 +1,49 @@ +package stirling.software.proprietary.policy.webhook; + +import java.security.SecureRandom; +import java.util.Base64; +import java.util.regex.Pattern; + +/** + * Generates and validates the two server-side tokens a webhook source carries: the public {@code + * webhookId} that routes a delivery URL to a source, and the {@code signingSecret} that + * authenticates each delivery. Both are URL-safe base64 of {@link SecureRandom} bytes, so both are + * unguessable; the id is the routing capability and the secret is the HMAC key. The id's character + * set is constrained so it can be used verbatim as a single path segment and as a spool directory + * name without traversal risk. + */ +public final class WebhookIds { + + /** URL-safe base64 without padding: exactly the characters a single path segment allows. */ + private static final Pattern VALID_ID = Pattern.compile("^[A-Za-z0-9_-]{16,128}$"); + + private static final SecureRandom RANDOM = new SecureRandom(); + private static final Base64.Encoder ENCODER = Base64.getUrlEncoder().withoutPadding(); + + private WebhookIds() {} + + /** + * A fresh routing token (~24 chars). Not secret, but unguessable so it cannot be enumerated. + */ + public static String newWebhookId() { + return randomToken(18); + } + + /** A fresh HMAC signing key (~43 chars) revealed to the operator once at creation. */ + public static String newSigningSecret() { + return randomToken(32); + } + + /** + * Whether {@code id} is a well-formed webhook id, safe as a path segment and directory name. + */ + public static boolean isValidId(String id) { + return id != null && VALID_ID.matcher(id).matches(); + } + + private static String randomToken(int bytes) { + byte[] buffer = new byte[bytes]; + RANDOM.nextBytes(buffer); + return ENCODER.encodeToString(buffer); + } +} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookReceiverController.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookReceiverController.java new file mode 100644 index 0000000000..55904eba32 --- /dev/null +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookReceiverController.java @@ -0,0 +1,167 @@ +package stirling.software.proprietary.policy.webhook; + +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.InputStream; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty; +import org.springframework.http.HttpStatus; +import org.springframework.http.MediaType; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestHeader; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; +import org.springframework.web.server.ResponseStatusException; + +import io.swagger.v3.oas.annotations.Hidden; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.tags.Tag; + +import jakarta.servlet.http.HttpServletRequest; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; + +import stirling.software.common.model.ApplicationProperties; +import stirling.software.proprietary.policy.source.Source; +import stirling.software.proprietary.policy.source.SourceStore; +import stirling.software.proprietary.policy.trigger.WebhookTrigger; + +/** + * Public receiver for webhook input sources. External systems POST a document to {@code + * /api/v1/webhooks/{webhookId}}; the request is authenticated not by a login session but by an + * HMAC-SHA256 signature of the exact body under the source's signing secret, so this endpoint sits + * outside the session auth wall (see {@code RequestUriUtils.isPublicAuthEndpoint}). A verified + * delivery is spooled and the referencing policies are fired immediately. + * + *

The request body is the raw document bytes, sent with a binary content type - {@code + * application/pdf} or {@code application/octet-stream} (e.g. {@code curl --data-binary @file.pdf -H + * 'Content-Type: application/pdf'}). A form content type ({@code + * application/x-www-form-urlencoded}) is rejected upstream because the container tries to parse the + * body as form fields. The signature header is {@code sha256=}. Delivery is rejected before + * anything is written if the id is unknown (404), the signature is missing or wrong (401), the + * source is paused (403), or the body is empty (400) or larger than {@code + * policies.webhookMaxBytes} (413). + */ +@Slf4j +@RestController +@RequestMapping("/api/v1/webhooks") +@Hidden +@RequiredArgsConstructor +@Tag(name = "Webhooks", description = "Inbound webhook source receiver") +@ConditionalOnBooleanProperty(name = "policies.enabled") +public class WebhookReceiverController { + + static final String SIGNATURE_HEADER = "X-Stirling-Signature"; + static final String FILENAME_HEADER = "X-Stirling-Filename"; + private static final String WEBHOOK_TYPE = "webhook"; + + private final SourceStore sourceStore; + private final WebhookSpool spool; + private final WebhookTrigger webhookTrigger; + private final ApplicationProperties applicationProperties; + + @PostMapping("/{webhookId}") + @Operation( + summary = "Deliver a document to a webhook source", + description = + "The body is the raw document; sign it with the source's secret and present" + + " 'sha256=' in the X-Stirling-Signature header. Returns 202 once" + + " the document is spooled for the referencing policies.") + public ResponseEntity receive( + @PathVariable String webhookId, + @RequestHeader(value = SIGNATURE_HEADER, required = false) String signature, + @RequestHeader(value = FILENAME_HEADER, required = false) String filename, + HttpServletRequest request) { + if (!WebhookIds.isValidId(webhookId)) { + throw new ResponseStatusException(HttpStatus.NOT_FOUND, "No such webhook"); + } + Source source = findWebhookSource(webhookId); + if (source == null) { + throw new ResponseStatusException(HttpStatus.NOT_FOUND, "No such webhook"); + } + + WebhookConfig config = WebhookConfig.from(source.options()); + byte[] body = readBoundedBody(request); + if (!WebhookSignatures.verify(config.signingSecret(), body, signature)) { + // Same 401 whether the header was absent or wrong: never confirm a guess. + throw new ResponseStatusException(HttpStatus.UNAUTHORIZED, "Invalid signature"); + } + if (!source.enabled()) { + throw new ResponseStatusException( + HttpStatus.FORBIDDEN, "Webhook source is paused; deliveries are not accepted"); + } + if (body.length == 0) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Empty request body"); + } + + String storedName; + try { + storedName = + WebhookSpool.displayName( + spool.store(webhookId, filename, body).getFileName().toString()); + } catch (IOException e) { + log.error("Could not spool webhook delivery for {}: {}", webhookId, e.getMessage()); + throw new ResponseStatusException( + HttpStatus.INTERNAL_SERVER_ERROR, "Could not store delivery"); + } + + // Fire the referencing policies now; the trigger's reconcile is the safety net. + webhookTrigger.fireForWebhook(webhookId); + log.info( + "Accepted webhook delivery '{}' ({} bytes) for {}", + storedName, + body.length, + webhookId); + return ResponseEntity.accepted() + .contentType(MediaType.APPLICATION_JSON) + .body(new WebhookDeliveryResponse(true, storedName, body.length)); + } + + /** The enabled-or-not webhook source whose routing id matches, or null if there is none. */ + private Source findWebhookSource(String webhookId) { + for (Source source : sourceStore.all()) { + if (!WEBHOOK_TYPE.equals(source.type())) { + continue; + } + Object configured = source.options().get(WebhookConfig.WEBHOOK_ID_OPTION); + if (configured != null && configured.toString().equals(webhookId)) { + return source; + } + } + return null; + } + + /** + * Read the body into memory, capped at {@code policies.webhookMaxBytes}. Reading one byte past + * the limit is enough to reject an over-sized (or unbounded, chunked) delivery with 413 before + * it can fill memory or the disk. + */ + private byte[] readBoundedBody(HttpServletRequest request) { + long maxBytes = applicationProperties.getPolicies().getWebhookMaxBytes(); + ByteArrayOutputStream buffer = new ByteArrayOutputStream(); + byte[] chunk = new byte[8192]; + long total = 0; + try (InputStream in = request.getInputStream()) { + int read; + while ((read = in.read(chunk)) != -1) { + total += read; + if (total > maxBytes) { + throw new ResponseStatusException( + HttpStatus.PAYLOAD_TOO_LARGE, + "Delivery exceeds the " + maxBytes + "-byte limit"); + } + buffer.write(chunk, 0, read); + } + } catch (IOException e) { + throw new ResponseStatusException( + HttpStatus.BAD_REQUEST, "Could not read request body"); + } + return buffer.toByteArray(); + } + + /** The 202 body: the stored (display) name and byte count of an accepted delivery. */ + public record WebhookDeliveryResponse(boolean accepted, String filename, int bytes) {} +} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookSignatures.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookSignatures.java new file mode 100644 index 0000000000..9f4f9ae221 --- /dev/null +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookSignatures.java @@ -0,0 +1,65 @@ +package stirling.software.proprietary.policy.webhook; + +import java.nio.charset.StandardCharsets; +import java.security.InvalidKeyException; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.util.HexFormat; + +import javax.crypto.Mac; +import javax.crypto.spec.SecretKeySpec; + +/** + * HMAC-SHA256 signing scheme for webhook deliveries. The sender signs the exact request body bytes + * with the source's {@code signingSecret} and presents the result as {@code sha256=} in the + * signature header; the receiver recomputes it and compares in constant time. Signing the raw body + * (rather than form fields) keeps the contract unambiguous - what is signed is exactly what is + * delivered - and matches how established webhook providers (Stripe, GitHub) sign payloads. + */ +public final class WebhookSignatures { + + private static final String ALGORITHM = "HmacSHA256"; + private static final String PREFIX = "sha256="; + + private WebhookSignatures() {} + + /** + * The header value a sender should present for {@code body}: {@code sha256=}. + */ + public static String sign(String signingSecret, byte[] body) { + return PREFIX + HexFormat.of().formatHex(hmac(signingSecret, body)); + } + + /** + * Whether {@code presented} (a {@code sha256=} header value, or a bare hex string) is a + * valid signature of {@code body} under {@code signingSecret}. Constant-time in the compared + * bytes; a malformed or missing header is simply false, never an exception. + */ + public static boolean verify(String signingSecret, byte[] body, String presented) { + if (signingSecret == null || presented == null || body == null) { + return false; + } + String hex = presented.trim(); + if (hex.regionMatches(true, 0, PREFIX, 0, PREFIX.length())) { + hex = hex.substring(PREFIX.length()); + } + byte[] presentedBytes; + try { + presentedBytes = HexFormat.of().parseHex(hex); + } catch (IllegalArgumentException notHex) { + return false; + } + return MessageDigest.isEqual(hmac(signingSecret, body), presentedBytes); + } + + private static byte[] hmac(String signingSecret, byte[] body) { + try { + Mac mac = Mac.getInstance(ALGORITHM); + mac.init(new SecretKeySpec(signingSecret.getBytes(StandardCharsets.UTF_8), ALGORITHM)); + return mac.doFinal(body); + } catch (NoSuchAlgorithmException | InvalidKeyException e) { + // HmacSHA256 is a required JCE algorithm and the key is always non-empty here. + throw new IllegalStateException("HMAC-SHA256 unavailable", e); + } + } +} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookSpool.java b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookSpool.java new file mode 100644 index 0000000000..182943ee8c --- /dev/null +++ b/app/proprietary/src/main/java/stirling/software/proprietary/policy/webhook/WebhookSpool.java @@ -0,0 +1,110 @@ +package stirling.software.proprietary.policy.webhook; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.StandardCopyOption; +import java.util.UUID; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty; +import org.springframework.stereotype.Component; + +import stirling.software.common.configuration.InstallationPathConfig; + +/** + * Server-owned staging area for webhook deliveries. Each webhook source gets its own directory, + * named by its (unguessable, path-segment-safe) {@code webhookId}, under a fixed root inside + * Stirling's install path - never a user-supplied location, so this is not subject to the folder + * source's allow-listing. The receiver writes each accepted document here atomically; {@link + * stirling.software.proprietary.policy.input.WebhookInputSource} then reads it exactly as the + * folder source reads a watched directory, reusing the whole claim/consume ledger. A delivered file + * keeps its original name for downstream output naming, prefixed with a unique token so two + * deliveries of the same name never collide. + */ +@Component +@ConditionalOnBooleanProperty(name = "policies.enabled") +public class WebhookSpool { + + private static final String SPOOL_DIR = "policy-webhook-spool"; + private static final String TEMP_SUFFIX = ".part"; + private static final String DEFAULT_NAME = "document.pdf"; + private static final int UNIQUE_LEN = 32; // UUID hex without dashes + + private final Path spoolRoot; + + public WebhookSpool() { + this(Path.of(InstallationPathConfig.getPath(), SPOOL_DIR)); + } + + // Lets a caller (and tests) root the spool at a chosen directory; Spring uses the no-arg one. + public WebhookSpool(Path spoolRoot) { + this.spoolRoot = spoolRoot.toAbsolutePath().normalize(); + } + + /** The staging directory for one webhook. Never escapes the spool root; may not yet exist. */ + public Path dirFor(String webhookId) { + if (!WebhookIds.isValidId(webhookId)) { + throw new IllegalArgumentException("invalid webhookId"); + } + Path dir = spoolRoot.resolve(webhookId).normalize(); + if (!dir.getParent().equals(spoolRoot)) { + // A validated id is a single safe segment; this only trips on a bug, never user input. + throw new IllegalArgumentException("invalid webhookId"); + } + return dir; + } + + /** + * Write one delivered document into the webhook's spool, atomically so a partial write is never + * picked up: staged under a hidden {@code .part} name and moved into place. Returns the final + * path. The original filename (sanitised to a bare, safe basename) is preserved after a unique + * prefix. + */ + public Path store(String webhookId, String filename, byte[] content) throws IOException { + Path dir = dirFor(webhookId); + Files.createDirectories(dir); + String finalName = spoolName(filename); + Path target = dir.resolve(finalName); + Path temp = dir.resolve("." + finalName + TEMP_SUFFIX); + Files.write(temp, content); + try { + Files.move(temp, target, StandardCopyOption.ATOMIC_MOVE); + } catch (IOException atomicUnsupported) { + Files.move(temp, target, StandardCopyOption.REPLACE_EXISTING); + } + return target; + } + + /** The spool file name for a delivery: a unique prefix plus the sanitised original name. */ + static String spoolName(String filename) { + return UUID.randomUUID().toString().replace("-", "") + "-" + sanitize(filename); + } + + /** The original name recovered from a spool file name, for the resolved input's filename. */ + public static String displayName(String spoolFileName) { + int dash = spoolFileName.indexOf('-'); + // The unique prefix is fixed-length hex with no dashes, so the first dash is the separator. + if (dash == UNIQUE_LEN && dash + 1 < spoolFileName.length()) { + return spoolFileName.substring(dash + 1); + } + return spoolFileName; + } + + /** Reduce a client-supplied filename to a safe bare basename; fall back to a default. */ + private static String sanitize(String filename) { + if (filename == null) { + return DEFAULT_NAME; + } + String base = filename.replace('\\', '/'); + int slash = base.lastIndexOf('/'); + if (slash >= 0) { + base = base.substring(slash + 1); + } + base = base.replaceAll("[^A-Za-z0-9._-]", "_").trim(); + // Strip leading dots so the result is never hidden (which resolve would skip) or empty. + while (base.startsWith(".")) { + base = base.substring(1); + } + return base.isEmpty() ? DEFAULT_NAME : base; + } +} diff --git a/app/proprietary/src/main/java/stirling/software/proprietary/util/SecretMasker.java b/app/proprietary/src/main/java/stirling/software/proprietary/util/SecretMasker.java index 5a9619c2dd..3402b89bb3 100644 --- a/app/proprietary/src/main/java/stirling/software/proprietary/util/SecretMasker.java +++ b/app/proprietary/src/main/java/stirling/software/proprietary/util/SecretMasker.java @@ -20,9 +20,10 @@ public final class SecretMasker { private static final Pattern SENSITIVE = RegexPatternUtils.getInstance() .getPattern( - // secret[_-]?access[_-]?key precedes plain secret so camelCase keys - // like secretAccessKey (no word boundary after "secret") still match. - "(?i)\\b(password|token|secret[_-]?access[_-]?key|secret|api[_-]?key|authorization|auth|jwt|cred|cert)\\b"); + // secret[_-]?access[_-]?key and signing[_-]?secret precede plain secret + // so camelCase keys like secretAccessKey and signingSecret (whose inner + // "secret" has no word boundary before/after it) still match. + "(?i)\\b(password|token|secret[_-]?access[_-]?key|signing[_-]?secret|secret|api[_-]?key|authorization|auth|jwt|cred|cert)\\b"); private SecretMasker() {} diff --git a/app/proprietary/src/test/java/stirling/software/proprietary/policy/input/WebhookInputSourceTest.java b/app/proprietary/src/test/java/stirling/software/proprietary/policy/input/WebhookInputSourceTest.java new file mode 100644 index 0000000000..8cf2ec7d5c --- /dev/null +++ b/app/proprietary/src/test/java/stirling/software/proprietary/policy/input/WebhookInputSourceTest.java @@ -0,0 +1,175 @@ +package stirling.software.proprietary.policy.input; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.lenient; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.function.Supplier; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.api.io.TempDir; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import stirling.software.common.util.FileReadinessChecker; +import stirling.software.proprietary.policy.ledger.InProcessProcessedLedger; +import stirling.software.proprietary.policy.model.InputSpec; +import stirling.software.proprietary.policy.webhook.WebhookConfig; +import stirling.software.proprietary.policy.webhook.WebhookSpool; + +/** + * Tests for {@link WebhookInputSource}: spooled deliveries are read through the ledger like folder + * files (consume removes them once processed, snapshot stays stateless), a not-yet-created spool is + * an empty source rather than an error, and options gain a routing id + signing secret on create. + */ +@ExtendWith(MockitoExtension.class) +class WebhookInputSourceTest { + + private static final String POLICY = "p1"; + private static final String WEBHOOK_ID = "testwebhookid1234"; + + @Mock private FileReadinessChecker readinessChecker; + + @TempDir Path tempDir; + + private WebhookSpool spool; + private WebhookInputSource source; + private InProcessProcessedLedger ledger; + private RecordingContext ctx; + + @BeforeEach + void setUp() { + spool = new WebhookSpool(tempDir.resolve("spool")); + source = new WebhookInputSource(spool, readinessChecker); + ledger = new InProcessProcessedLedger(); + ctx = new RecordingContext(); + lenient().when(readinessChecker.isReady(any())).thenReturn(true); + } + + private static InputSpec spec(String mode) { + return new InputSpec( + "webhook", + Map.of("webhookId", WEBHOOK_ID, "signingSecret", "secret", "mode", mode)); + } + + @Test + void consumeRemovesTheDeliveryOnceProcessed() throws IOException { + Path delivered = spool.store(WEBHOOK_ID, "doc.pdf", "data".getBytes()); + + List work = source.resolve(spec("consume"), ctx); + + assertEquals(1, work.size()); + assertEquals("doc.pdf", work.get(0).inputs().primary().get(0).getFilename()); + // In flight: still spooled, but a second sweep does not pick it up again. + assertTrue(Files.exists(delivered)); + assertTrue(source.resolve(spec("consume"), ctx).isEmpty()); + + work.get(0).onComplete().accept(true); + assertTrue(Files.notExists(delivered)); + assertTrue(source.resolve(spec("consume"), ctx).isEmpty()); + } + + @Test + void aFailedRunLeavesTheDeliveryInPlace() throws IOException { + Path delivered = spool.store(WEBHOOK_ID, "doc.pdf", "data".getBytes()); + + List work = source.resolve(spec("consume"), ctx); + work.get(0).onComplete().accept(false); + + assertTrue(Files.exists(delivered)); + } + + @Test + void snapshotReReadsEveryRunAndNeverDeletes() throws IOException { + Path delivered = spool.store(WEBHOOK_ID, "doc.pdf", "data".getBytes()); + + assertEquals(1, source.resolve(spec("snapshot"), ctx).size()); + List second = source.resolve(spec("snapshot"), ctx); + assertEquals(1, second.size()); + second.get(0).onComplete().accept(true); + assertTrue(Files.exists(delivered)); + } + + @Test + void nothingDeliveredIsAnEmptySourceNotAnError() throws IOException { + List work = source.resolve(spec("consume"), ctx); + assertTrue(work.isEmpty()); + assertTrue(ctx.present.isEmpty()); + } + + @Test + void validateRejectsMissingIdOrSecret() { + assertThrows( + IllegalArgumentException.class, + () -> source.validate(new InputSpec("webhook", Map.of("signingSecret", "s")))); + assertThrows( + IllegalArgumentException.class, + () -> source.validate(new InputSpec("webhook", Map.of("webhookId", WEBHOOK_ID)))); + } + + @Test + void prepareMintsIdAndSecretOnCreate() { + Map prepared = + source.prepareOptionsForSave(Map.of("mode", "consume"), true); + + String id = prepared.get(WebhookConfig.WEBHOOK_ID_OPTION).toString(); + String secret = prepared.get(WebhookConfig.SIGNING_SECRET_OPTION).toString(); + assertFalse(id.isBlank()); + assertFalse(secret.isBlank()); + assertEquals("consume", prepared.get("mode")); + // Two creates never collide. + Map other = source.prepareOptionsForSave(Map.of(), true); + assertNotEquals(id, other.get(WebhookConfig.WEBHOOK_ID_OPTION).toString()); + } + + @Test + void prepareLeavesAnExistingWebhookUntouchedOnEdit() { + Map existing = + Map.of("webhookId", WEBHOOK_ID, "signingSecret", "keepme", "mode", "snapshot"); + + Map prepared = source.prepareOptionsForSave(existing, false); + + assertEquals(WEBHOOK_ID, prepared.get("webhookId")); + assertEquals("keepme", prepared.get("signingSecret")); + } + + /** Policy-scoped context backed by the in-process ledger, recording presence reports. */ + private class RecordingContext implements ResolveContext { + + private final List present = new ArrayList<>(); + + @Override + public boolean claim(String identity, String gate, Supplier contentHash) { + return ledger.claim(POLICY, identity, gate, contentHash); + } + + @Override + public void settle( + String identity, String finalGate, String finalContentHash, boolean success) { + ledger.settle(POLICY, identity, finalGate, finalContentHash, success); + } + + @Override + public boolean allSettledDone(String identity) { + return ledger.allSettledDone(identity); + } + + @Override + public void reportPresent(Collection identities) { + present.addAll(identities); + } + } +} diff --git a/app/proprietary/src/test/java/stirling/software/proprietary/policy/source/SourceControllerTest.java b/app/proprietary/src/test/java/stirling/software/proprietary/policy/source/SourceControllerTest.java index 295b4785b8..136c871af9 100644 --- a/app/proprietary/src/test/java/stirling/software/proprietary/policy/source/SourceControllerTest.java +++ b/app/proprietary/src/test/java/stirling/software/proprietary/policy/source/SourceControllerTest.java @@ -1,33 +1,41 @@ package stirling.software.proprietary.policy.source; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import java.nio.file.Path; import java.util.List; import java.util.Map; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; import org.springframework.http.ResponseEntity; import org.springframework.web.server.ResponseStatusException; import stirling.software.common.model.ApplicationProperties; import stirling.software.common.service.UserServiceInterface; +import stirling.software.common.util.FileReadinessChecker; import stirling.software.proprietary.policy.config.PolicyAccessGuard; import stirling.software.proprietary.policy.config.PolicyManagementAuthority; import stirling.software.proprietary.policy.input.InputSource; +import stirling.software.proprietary.policy.input.WebhookInputSource; import stirling.software.proprietary.policy.model.OutputSpec; import stirling.software.proprietary.policy.model.PipelineStep; import stirling.software.proprietary.policy.model.Policy; import stirling.software.proprietary.policy.store.InProcessPolicyStore; import stirling.software.proprietary.policy.store.PolicyStore; import stirling.software.proprietary.policy.trigger.PolicyTriggerManager; +import stirling.software.proprietary.policy.webhook.WebhookSpool; import stirling.software.proprietary.util.SecretMasker; /** @@ -41,6 +49,9 @@ class SourceControllerTest { private final PolicyStore policyStore = new InProcessPolicyStore(); private PolicyTriggerManager triggerManager; private SourceController controller; + private SourceController webhookController; + + @TempDir Path tempDir; @BeforeEach void setUp() { @@ -58,9 +69,13 @@ class SourceControllerTest { policyGuard, new InProcessSourceDocCounter()); triggerManager = mock(PolicyTriggerManager.class); - // A permissive input source so config validation passes and save can be exercised. + // A permissive input source so config validation passes and save can be exercised. Model a + // real source's prepareOptionsForSave as a pass-through (a bare mock would return an empty + // map for the Map-typed method and wipe the config). InputSource folderInput = mock(InputSource.class); when(folderInput.supports(any())).thenReturn(true); + when(folderInput.prepareOptionsForSave(any(), anyBoolean())) + .thenAnswer(invocation -> invocation.getArgument(0)); controller = new SourceController( sourceStore, @@ -72,6 +87,47 @@ class SourceControllerTest { triggerManager, properties, List.of(folderInput)); + WebhookInputSource webhookInput = + new WebhookInputSource(new WebhookSpool(tempDir), mock(FileReadinessChecker.class)); + webhookController = + new SourceController( + sourceStore, + sourceGuard, + overviewService, + policyStore, + policyGuard, + authority, + triggerManager, + properties, + List.of(webhookInput)); + } + + @Test + void creatingAWebhookRevealsItsSecretOnceThenMasks() { + Source created = + webhookController + .save( + new Source( + null, + "Partner uploads", + "webhook", + Map.of("mode", "consume"), + true, + null, + null)) + .getBody(); + + // The create response reveals the server-minted secret + routing id exactly once. + String secret = String.valueOf(created.options().get("signingSecret")); + String webhookId = String.valueOf(created.options().get("webhookId")); + assertNotEquals(SecretMasker.REDACTED, secret); + assertFalse(secret.isBlank()); + assertFalse(webhookId.isBlank()); + + // Every later read masks the secret but keeps the (non-secret) routing id. + Source read = webhookController.get(created.id()).getBody(); + assertEquals(SecretMasker.REDACTED, read.options().get("signingSecret")); + assertEquals(webhookId, read.options().get("webhookId")); } @Test diff --git a/app/proprietary/src/test/java/stirling/software/proprietary/policy/trigger/WebhookTriggerTest.java b/app/proprietary/src/test/java/stirling/software/proprietary/policy/trigger/WebhookTriggerTest.java new file mode 100644 index 0000000000..a21968b638 --- /dev/null +++ b/app/proprietary/src/test/java/stirling/software/proprietary/policy/trigger/WebhookTriggerTest.java @@ -0,0 +1,121 @@ +package stirling.software.proprietary.policy.trigger; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import stirling.software.common.model.ApplicationProperties; +import stirling.software.proprietary.policy.engine.PolicyRunner; +import stirling.software.proprietary.policy.engine.SweepKind; +import stirling.software.proprietary.policy.model.OutputSpec; +import stirling.software.proprietary.policy.model.PipelineStep; +import stirling.software.proprietary.policy.model.Policy; +import stirling.software.proprietary.policy.model.TriggerConfig; +import stirling.software.proprietary.policy.source.InProcessSourceStore; +import stirling.software.proprietary.policy.source.Source; +import stirling.software.proprietary.policy.source.SourceStore; +import stirling.software.proprietary.policy.store.PolicyStore; + +/** + * Tests for {@link WebhookTrigger}'s dispatch: a delivery fires only the policies that reference + * the delivered-to webhook, and validation requires at least one webhook source. + */ +@ExtendWith(MockitoExtension.class) +class WebhookTriggerTest { + + private static final String TYPE = "webhook"; + + @Mock private PolicyStore policyStore; + @Mock private PolicyRunner policyRunner; + + private final SourceStore sourceStore = new InProcessSourceStore(); + private WebhookTrigger trigger; + + @BeforeEach + void setUp() { + trigger = + new WebhookTrigger( + policyStore, policyRunner, sourceStore, new ApplicationProperties()); + } + + @Test + void firesOnlyPoliciesReferencingTheDeliveredWebhook() { + Policy matching = webhookPolicy("a", "whkA"); + Policy other = webhookPolicy("b", "whkB"); + when(policyStore.findByTriggerType(TYPE)).thenReturn(List.of(matching, other)); + + trigger.fireForWebhook("whkA"); + + verify(policyRunner).run(matching, SweepKind.LIGHT); + verify(policyRunner, never()).run(other, SweepKind.LIGHT); + } + + @Test + void ignoresADeliveryForAnUnknownWebhookId() { + Policy policy = webhookPolicy("a", "whkA"); + when(policyStore.findByTriggerType(TYPE)).thenReturn(List.of(policy)); + + trigger.fireForWebhook("whkZ"); + + verify(policyRunner, never()).run(any(), any(SweepKind.class)); + } + + @Test + void validateRequiresAWebhookSource() { + assertThrows( + IllegalArgumentException.class, + () -> trigger.validate(policy("p", webhookTriggerConfig(), List.of()))); + // A policy that references a webhook source validates. + trigger.validate(webhookPolicy("p", "whkA")); + } + + private static TriggerConfig webhookTriggerConfig() { + return new TriggerConfig(TYPE, Map.of()); + } + + /** Persist a webhook source with the given routing id and return a policy referencing it. */ + private Policy webhookPolicy(String id, String webhookId) { + String sourceId = + sourceStore + .save( + new Source( + null, + "hook", + "webhook", + Map.of( + "webhookId", + webhookId, + "signingSecret", + "s", + "mode", + "consume"), + true, + "owner", + null)) + .id(); + return policy(id, webhookTriggerConfig(), List.of(sourceId)); + } + + private static Policy policy(String id, TriggerConfig trigger, List sourceIds) { + return new Policy( + id, + "hook", + "owner", + true, + trigger, + sourceIds, + List.of(new PipelineStep("/api/v1/misc/compress-pdf", Map.of())), + OutputSpec.inline()); + } +} diff --git a/app/proprietary/src/test/java/stirling/software/proprietary/policy/webhook/WebhookReceiverControllerTest.java b/app/proprietary/src/test/java/stirling/software/proprietary/policy/webhook/WebhookReceiverControllerTest.java new file mode 100644 index 0000000000..64ad8cfad9 --- /dev/null +++ b/app/proprietary/src/test/java/stirling/software/proprietary/policy/webhook/WebhookReceiverControllerTest.java @@ -0,0 +1,180 @@ +package stirling.software.proprietary.policy.webhook; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; +import java.util.Map; +import java.util.stream.Stream; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.springframework.http.ResponseEntity; +import org.springframework.mock.web.MockHttpServletRequest; +import org.springframework.web.server.ResponseStatusException; + +import stirling.software.common.model.ApplicationProperties; +import stirling.software.proprietary.policy.source.InProcessSourceStore; +import stirling.software.proprietary.policy.source.Source; +import stirling.software.proprietary.policy.source.SourceStore; +import stirling.software.proprietary.policy.trigger.WebhookTrigger; +import stirling.software.proprietary.policy.webhook.WebhookReceiverController.WebhookDeliveryResponse; + +/** + * Tests for the public webhook receiver: a signature-valid delivery is spooled and fires the + * trigger; an unknown id, wrong signature, paused source, empty body, or over-sized body is + * rejected before anything is stored. + */ +class WebhookReceiverControllerTest { + + private static final String WEBHOOK_ID = "receivertestid12"; + private static final String SECRET = "topsecret"; + private static final byte[] BODY = "a pdf".getBytes(StandardCharsets.UTF_8); + + @TempDir Path tempDir; + + private SourceStore sourceStore; + private WebhookSpool spool; + private WebhookTrigger trigger; + private ApplicationProperties properties; + private WebhookReceiverController controller; + + @BeforeEach + void setUp() { + sourceStore = new InProcessSourceStore(); + sourceStore.save(webhookSource(true)); + spool = new WebhookSpool(tempDir.resolve("spool")); + trigger = mock(WebhookTrigger.class); + properties = new ApplicationProperties(); + controller = new WebhookReceiverController(sourceStore, spool, trigger, properties); + } + + private static Source webhookSource(boolean enabled) { + return new Source( + "s1", + "Partner uploads", + "webhook", + Map.of("webhookId", WEBHOOK_ID, "signingSecret", SECRET, "mode", "consume"), + enabled, + "owner", + null); + } + + private static MockHttpServletRequest request(byte[] body) { + MockHttpServletRequest req = + new MockHttpServletRequest("POST", "/api/v1/webhooks/" + WEBHOOK_ID); + req.setContent(body); + return req; + } + + @Test + void aValidDeliveryIsSpooledAndFiresTheTrigger() throws IOException { + String signature = WebhookSignatures.sign(SECRET, BODY); + + ResponseEntity response = + controller.receive(WEBHOOK_ID, signature, "invoice.pdf", request(BODY)); + + assertEquals(202, response.getStatusCode().value()); + assertTrue(response.getBody().accepted()); + assertEquals("invoice.pdf", response.getBody().filename()); + assertEquals(1, spooledFiles().size()); + verify(trigger).fireForWebhook(WEBHOOK_ID); + } + + @Test + void aWrongSignatureIsRejectedAndStoresNothing() { + ResponseStatusException ex = + assertThrows( + ResponseStatusException.class, + () -> + controller.receive( + WEBHOOK_ID, "sha256=deadbeef", "x.pdf", request(BODY))); + + assertEquals(401, ex.getStatusCode().value()); + assertTrue(spooledFiles().isEmpty()); + verify(trigger, never()).fireForWebhook(WEBHOOK_ID); + } + + @Test + void anUnknownWebhookIsNotFound() { + ResponseStatusException ex = + assertThrows( + ResponseStatusException.class, + () -> + controller.receive( + "unknownwebhookid", + WebhookSignatures.sign(SECRET, BODY), + "x.pdf", + request(BODY))); + + assertEquals(404, ex.getStatusCode().value()); + } + + @Test + void aPausedSourceRejectsDeliveries() { + sourceStore.save(webhookSource(false)); + String signature = WebhookSignatures.sign(SECRET, BODY); + + ResponseStatusException ex = + assertThrows( + ResponseStatusException.class, + () -> controller.receive(WEBHOOK_ID, signature, "x.pdf", request(BODY))); + + assertEquals(403, ex.getStatusCode().value()); + assertTrue(spooledFiles().isEmpty()); + } + + @Test + void anEmptyBodyIsRejected() { + byte[] empty = new byte[0]; + String signature = WebhookSignatures.sign(SECRET, empty); + + ResponseStatusException ex = + assertThrows( + ResponseStatusException.class, + () -> controller.receive(WEBHOOK_ID, signature, null, request(empty))); + + assertEquals(400, ex.getStatusCode().value()); + } + + @Test + void anOversizeDeliveryIsRejectedBeforeStoring() { + properties.getPolicies().setWebhookMaxBytes(2); + + ResponseStatusException ex = + assertThrows( + ResponseStatusException.class, + () -> + controller.receive( + WEBHOOK_ID, + WebhookSignatures.sign(SECRET, BODY), + "x.pdf", + request(BODY))); + + assertEquals(413, ex.getStatusCode().value()); + assertTrue(spooledFiles().isEmpty()); + } + + private List spooledFiles() { + Path dir = spool.dirFor(WEBHOOK_ID); + if (!Files.isDirectory(dir)) { + return List.of(); + } + try (Stream entries = Files.list(dir)) { + return entries.filter(Files::isRegularFile) + .filter(p -> !p.getFileName().toString().startsWith(".")) + .toList(); + } catch (IOException e) { + throw new RuntimeException(e); + } + } +} diff --git a/app/proprietary/src/test/java/stirling/software/proprietary/policy/webhook/WebhookSignaturesTest.java b/app/proprietary/src/test/java/stirling/software/proprietary/policy/webhook/WebhookSignaturesTest.java new file mode 100644 index 0000000000..4e1e4cdcf5 --- /dev/null +++ b/app/proprietary/src/test/java/stirling/software/proprietary/policy/webhook/WebhookSignaturesTest.java @@ -0,0 +1,49 @@ +package stirling.software.proprietary.policy.webhook; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.nio.charset.StandardCharsets; + +import org.junit.jupiter.api.Test; + +/** The HMAC-SHA256 webhook signing scheme: correct signatures verify, tampered ones do not. */ +class WebhookSignaturesTest { + + private static final String SECRET = "whsec_test_secret"; + private static final byte[] BODY = "the document bytes".getBytes(StandardCharsets.UTF_8); + + @Test + void aFreshlySignedBodyVerifies() { + String header = WebhookSignatures.sign(SECRET, BODY); + assertTrue(header.startsWith("sha256=")); + assertTrue(WebhookSignatures.verify(SECRET, BODY, header)); + } + + @Test + void abarehexSignatureVerifiesToo() { + String header = WebhookSignatures.sign(SECRET, BODY); + String bareHex = header.substring("sha256=".length()); + assertTrue(WebhookSignatures.verify(SECRET, BODY, bareHex)); + } + + @Test + void aWrongSecretDoesNotVerify() { + String header = WebhookSignatures.sign(SECRET, BODY); + assertFalse(WebhookSignatures.verify("other-secret", BODY, header)); + } + + @Test + void atamperedBodyDoesNotVerify() { + String header = WebhookSignatures.sign(SECRET, BODY); + byte[] tampered = "the document byteS".getBytes(StandardCharsets.UTF_8); + assertFalse(WebhookSignatures.verify(SECRET, tampered, header)); + } + + @Test + void aMissingOrMalformedHeaderIsFalseNotAnError() { + assertFalse(WebhookSignatures.verify(SECRET, BODY, null)); + assertFalse(WebhookSignatures.verify(SECRET, BODY, "sha256=not-hex")); + assertFalse(WebhookSignatures.verify(SECRET, BODY, "")); + } +} diff --git a/app/proprietary/src/test/java/stirling/software/proprietary/util/SecretMaskerTest.java b/app/proprietary/src/test/java/stirling/software/proprietary/util/SecretMaskerTest.java index 415117a446..ef5571824d 100644 --- a/app/proprietary/src/test/java/stirling/software/proprietary/util/SecretMaskerTest.java +++ b/app/proprietary/src/test/java/stirling/software/proprietary/util/SecretMaskerTest.java @@ -103,6 +103,18 @@ class SecretMaskerTest { assertEquals("AKIAEXAMPLE", result.get("accessKeyId")); } + @Test + @DisplayName("should mask camelCase signingSecret despite no word boundary") + void shouldMaskCamelCaseSigningSecret() { + Map input = Map.of("signingSecret", "shh", "webhookId", "whk_abc"); + + Map result = SecretMasker.mask(input); + + assertEquals(SecretMasker.REDACTED, result.get("signingSecret")); + // The routing id is a public URL token, not a secret. + assertEquals("whk_abc", result.get("webhookId")); + } + @Test @DisplayName("should mask nested map sensitive keys") void shouldMaskNestedMapSensitiveKeys() { diff --git a/frontend/editor/public/locales/en-US/translation.toml b/frontend/editor/public/locales/en-US/translation.toml index 9663478000..95c9c9bf89 100644 --- a/frontend/editor/public/locales/en-US/translation.toml +++ b/frontend/editor/public/locales/en-US/translation.toml @@ -8105,6 +8105,32 @@ placeholder = "us-east-1" [portal.sources.types.s3.fields.secretAccessKey] label = "Secret access key" +[portal.sources.types.webhook] +description = "Receive documents by signed HTTP POST from any external system." +label = "Webhook" + +[portal.sources.types.webhook.detail] +deliveryUrl = "Delivery URL" +secretNote = "The signing secret is shown once, when the webhook is created. Recreate the source to roll it." + +[portal.sources.types.webhook.fields.mode] +helperText = "Consume removes each delivered document once every policy has processed it." +label = "Read mode" + +[portal.sources.types.webhook.fields.mode.options] +consume = "Consume: process each delivery once" +snapshot = "Snapshot: re-read the spool every run" + +[portal.sources.types.webhook.reveal] +copy = "Copy" +done = "Done" +secret = "Signing secret" +secretHelp = "Sign each delivery's raw body with this key (HMAC-SHA256) and send it as the X-Stirling-Signature header." +secretWarning = "Copy the signing secret now. For your security it is shown only once and cannot be retrieved later." +title = "Webhook created" +url = "Delivery URL" +usage = "POST each document as the raw request body to the delivery URL with a binary content type (application/pdf or application/octet-stream). Referencing policies run automatically on arrival." + [portal.sources.types.unknown] label = "Source" diff --git a/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.test.ts b/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.test.ts index 1eec2ac959..83354cd44d 100644 --- a/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.test.ts +++ b/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.test.ts @@ -3,8 +3,10 @@ import { describe, expect, it } from "vitest"; import { creatableSourceTypes } from "@portal/components/sources/creatableSourceTypes"; describe("creatableSourceTypes (SaaS)", () => { - it("never offers folder sources: hosted deployments do not read the server filesystem", () => { - expect(creatableSourceTypes().map((t) => t.type)).not.toContain("folder"); + it("never offers server-local-disk sources (folder, webhook) in hosted deployments", () => { + const offered = creatableSourceTypes().map((t) => t.type); + expect(offered).not.toContain("folder"); + expect(offered).not.toContain("webhook"); }); it("still offers the cloud source types", () => { diff --git a/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.ts b/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.ts index 9fa73be0ff..7415b20fa9 100644 --- a/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.ts +++ b/frontend/editor/src/portal-saas/components/sources/creatableSourceTypes.ts @@ -4,10 +4,16 @@ import { } from "@portal/components/sources/sourceTypes"; /** - * Hosted deployments never read the server's filesystem (the backend's - * FolderAccessGuard denies it outright), so folder connections are not offered - * in the connect wizard. + * Hosted deployments don't rely on the server's local filesystem, so the two + * server-local-disk source types are not offered in the connect wizard: folder + * (denied outright by the backend's FolderAccessGuard) and webhook (whose + * delivery spool is node-local, so it can't be relied on across a multi-node + * fleet). Cloud sources such as S3 remain. */ +const SERVER_LOCAL_TYPES = new Set(["folder", "webhook"]); + export function creatableSourceTypes(): CreatableSourceType[] { - return CREATABLE_SOURCE_TYPES.filter((type) => type.type !== "folder"); + return CREATABLE_SOURCE_TYPES.filter( + (type) => !SERVER_LOCAL_TYPES.has(type.type), + ); } diff --git a/frontend/editor/src/portal/api/sources.ts b/frontend/editor/src/portal/api/sources.ts index 9a66aeef2d..6252010738 100644 --- a/frontend/editor/src/portal/api/sources.ts +++ b/frontend/editor/src/portal/api/sources.ts @@ -32,6 +32,8 @@ export interface SourceView { docsTotal: number; docs24h: number; docs30d: number; + /** Webhook sources only: the server-relative delivery path senders POST to. Null otherwise. */ + webhookPath?: string | null; } export interface SourceKpi { diff --git a/frontend/editor/src/portal/components/sources/ConnectWizard.test.tsx b/frontend/editor/src/portal/components/sources/ConnectWizard.test.tsx index 7723b447ca..325d6096cc 100644 --- a/frontend/editor/src/portal/components/sources/ConnectWizard.test.tsx +++ b/frontend/editor/src/portal/components/sources/ConnectWizard.test.tsx @@ -171,6 +171,57 @@ describe("ConnectWizard", () => { }); }); + it("reveals the delivery URL and signing secret once after creating a webhook", async () => { + createSource.mockResolvedValue({ + id: "w1", + options: { webhookId: "whk_abc123", signingSecret: "whsec_topsecret" }, + }); + const onCreated = vi.fn(); + const onClose = vi.fn(); + + renderWithMantine( + , + ); + + // Step 0: pick the webhook type card, continue. + fireEvent.click(screen.getByText("portal.sources.types.webhook.label")); + fireEvent.click(screen.getByText("portal.sources.wizard.continue")); + + // Step 1: only a name is required (mode defaults to consume). + fireEvent.change(screen.getAllByRole("textbox")[0], { + target: { value: "Partner uploads" }, + }); + fireEvent.click(screen.getByText("portal.sources.wizard.continue")); + + // Step 2: submit. + fireEvent.click(screen.getByText("portal.sources.actions.connectSource")); + + await waitFor(() => { + expect(createSource).toHaveBeenCalledWith({ + name: "Partner uploads", + type: "webhook", + options: { mode: "consume" }, + enabled: true, + }); + }); + + // The reveal step shows the secret and the delivery URL; the modal stays open + // (onClose not called) until the operator clicks Done. + expect( + await screen.findByDisplayValue("whsec_topsecret"), + ).toBeInTheDocument(); + expect( + screen.getByDisplayValue(/\/api\/v1\/webhooks\/whk_abc123$/), + ).toBeInTheDocument(); + expect(onCreated).toHaveBeenCalledTimes(1); + expect(onClose).not.toHaveBeenCalled(); + + fireEvent.click( + screen.getByText("portal.sources.types.webhook.reveal.done"), + ); + expect(onClose).toHaveBeenCalledTimes(1); + }); + it("renders the inline error message when create fails", async () => { createSource.mockRejectedValue( new HttpError(400, "Bad Request", { diff --git a/frontend/editor/src/portal/components/sources/ConnectWizard.tsx b/frontend/editor/src/portal/components/sources/ConnectWizard.tsx index c284c24655..dcd08136b1 100644 --- a/frontend/editor/src/portal/components/sources/ConnectWizard.tsx +++ b/frontend/editor/src/portal/components/sources/ConnectWizard.tsx @@ -15,6 +15,7 @@ import { creatableSourceTypes } from "@portal/components/sources/creatableSource import { defaultOptions, sourceTypeMeta, + WEBHOOK_SOURCE_TYPE, type CreatableSourceType, } from "@portal/components/sources/sourceTypes"; import "@portal/views/Sources.css"; @@ -78,6 +79,11 @@ export function ConnectWizard({ ); const [submitting, setSubmitting] = useState(false); const [error, setError] = useState(null); + // Set after a webhook is created: its one-time delivery id + signing secret, shown before close. + const [reveal, setReveal] = useState<{ + webhookId: string; + secret: string; + } | null>(null); // Re-seed the form whenever the wizard opens (or its target source changes) so // editing prefills the current config and a reopened create starts clean. @@ -90,6 +96,7 @@ export function ConnectWizard({ setOptions(optionsFor(ct, source?.options)); setSubmitting(false); setError(null); + setReveal(null); }, [open, source]); const stepId = steps[stepIndex]; @@ -120,8 +127,20 @@ export function ConnectWizard({ options, enabled: source?.enabled ?? true, }; - await createSource(isEdit ? { ...fields, id: source.id } : fields); + const saved = await createSource( + isEdit ? { ...fields, id: source.id } : fields, + ); onCreated(); + // A new webhook returns its server-minted routing id + signing secret once; reveal them + // (with the delivery URL) before closing so the operator can copy the secret. + if (!isEdit && type.type === WEBHOOK_SOURCE_TYPE) { + const webhookId = String(saved.options?.webhookId ?? ""); + const secret = String(saved.options?.signingSecret ?? ""); + if (webhookId && secret) { + setReveal({ webhookId, secret }); + return; + } + } onClose(); } catch (e) { setError(errorMessage(e)); @@ -130,6 +149,15 @@ export function ConnectWizard({ } } + function copy(text: string) { + void navigator.clipboard?.writeText(text); + } + + const webhookUrl = reveal + ? `${window.location.origin}/api/v1/webhooks/${reveal.webhookId}` + : ""; + const revealSecret = reveal ? reveal.secret : ""; + const stepLabels: Record = { type: t("portal.sources.wizard.steps.chooseType"), configure: t("portal.sources.wizard.steps.configure"), @@ -142,157 +170,233 @@ export function ConnectWizard({ onClose={onClose} width="lg" title={ - isEdit - ? t("portal.sources.wizard.editTitle") - : t("portal.sources.wizard.title") + reveal + ? t("portal.sources.types.webhook.reveal.title") + : isEdit + ? t("portal.sources.wizard.editTitle") + : t("portal.sources.wizard.title") + } + subtitle={ + reveal + ? undefined + : t("portal.sources.wizard.subtitle", { + current: stepIndex + 1, + total: steps.length, + label: stepLabels[stepId], + }) } - subtitle={t("portal.sources.wizard.subtitle", { - current: stepIndex + 1, - total: steps.length, - label: stepLabels[stepId], - })} footer={ -

- - -
+ reveal ? ( +
+ +
+ ) : ( +
+ + +
+ ) } > -
    - {steps.map((id, i) => ( -
  1. - - {i < stepIndex ? "✓" : i + 1} - - {stepLabels[id]} -
  2. - ))} -
- - {stepId === "type" && ( -
- {OFFERED_TYPES.map((ct) => ( - - ))} -
- )} - - {stepId === "configure" && ( + {reveal && (
- - setName(e.target.value)} - /> - - {type.fields.map((field) => ( - - {field.control === "select" ? ( - - setOptions((o) => ({ ...o, [field.key]: e.target.value })) - } - /> - )} - - ))} -
- )} - - {stepId === "review" && ( -
-
- - - {type.fields.map((field) => ( - + +
+ e.currentTarget.select()} /> - ))} -
- {error && } + +
+ + +
+ e.currentTarget.select()} + /> + +
+
+

+ {t("portal.sources.types.webhook.reveal.usage")} +

)} + + {!reveal && ( + <> +
    + {steps.map((id, i) => ( +
  1. + + {i < stepIndex ? "✓" : i + 1} + + {stepLabels[id]} +
  2. + ))} +
+ + {stepId === "type" && ( +
+ {OFFERED_TYPES.map((ct) => ( + + ))} +
+ )} + + {stepId === "configure" && ( +
+ + setName(e.target.value)} + /> + + {type.fields.map((field) => ( + + {field.control === "select" ? ( + + setOptions((o) => ({ + ...o, + [field.key]: e.target.value, + })) + } + /> + )} + + ))} +
+ )} + + {stepId === "review" && ( +
+
+ + + {type.fields.map((field) => ( + + ))} +
+ {error && } +
+ )} + + )} ); } diff --git a/frontend/editor/src/portal/components/sources/SourceDetailPanel.tsx b/frontend/editor/src/portal/components/sources/SourceDetailPanel.tsx index 82fd8a5b1b..13ed075067 100644 --- a/frontend/editor/src/portal/components/sources/SourceDetailPanel.tsx +++ b/frontend/editor/src/portal/components/sources/SourceDetailPanel.tsx @@ -1,8 +1,11 @@ import { useTranslation } from "react-i18next"; -import { Chip, StatTile } from "@app/ui"; +import { Button, Chip, StatTile } from "@app/ui"; import type { SourceView } from "@portal/api/sources"; import { Sparkline } from "@portal/components/sources/Sparkline"; -import { EDITOR_SOURCE_TYPE } from "@portal/components/sources/sourceTypes"; +import { + EDITOR_SOURCE_TYPE, + WEBHOOK_SOURCE_TYPE, +} from "@portal/components/sources/sourceTypes"; import "@portal/views/Sources.css"; interface SourceDetailPanelProps { @@ -22,6 +25,10 @@ export function SourceDetailPanel({ }: SourceDetailPanelProps) { const { t } = useTranslation(); const isEditor = source.type === EDITOR_SOURCE_TYPE; + const isWebhook = source.type === WEBHOOK_SOURCE_TYPE; + const webhookUrl = source.webhookPath + ? `${window.location.origin}${source.webhookPath}` + : null; return (
{isEditor && ( @@ -30,6 +37,29 @@ export function SourceDetailPanel({

)} + {isWebhook && webhookUrl && ( +
+ + {t("portal.sources.types.webhook.detail.deliveryUrl")} + +
+ {webhookUrl} + +
+

+ {t("portal.sources.types.webhook.detail.secretNote")} +

+
+ )} + {source.config.length > 0 && (
{source.config.map((row) => ( diff --git a/frontend/editor/src/portal/components/sources/sourceTypes.ts b/frontend/editor/src/portal/components/sources/sourceTypes.ts index a3f19b6b36..16e52a6637 100644 --- a/frontend/editor/src/portal/components/sources/sourceTypes.ts +++ b/frontend/editor/src/portal/components/sources/sourceTypes.ts @@ -22,6 +22,9 @@ export interface SourceTypeMeta { */ export const EDITOR_SOURCE_TYPE = "editor"; +/** The webhook source type. Its delivery URL + signing secret are minted server-side on create. */ +export const WEBHOOK_SOURCE_TYPE = "webhook"; + const SOURCE_TYPE_META: Record = { folder: { labelKey: "portal.sources.types.folder.label", @@ -38,6 +41,11 @@ const SOURCE_TYPE_META: Record = { icon: "☁", accent: "brand", }, + webhook: { + labelKey: "portal.sources.types.webhook.label", + icon: "↯", + accent: "warning", + }, }; const UNKNOWN_TYPE_META: SourceTypeMeta = { @@ -206,6 +214,34 @@ export const CREATABLE_SOURCE_TYPES: CreatableSourceType[] = [ }, ], }, + { + // The delivery URL and signing secret are generated server-side on create and revealed once, + // so the only user-facing config is how deliveries are consumed. + type: WEBHOOK_SOURCE_TYPE, + labelKey: "portal.sources.types.webhook.label", + descriptionKey: "portal.sources.types.webhook.description", + fields: [ + { + key: "mode", + labelKey: "portal.sources.types.webhook.fields.mode.label", + control: "select", + defaultValue: "consume", + helperTextKey: "portal.sources.types.webhook.fields.mode.helperText", + options: [ + { + value: "consume", + labelKey: + "portal.sources.types.webhook.fields.mode.options.consume", + }, + { + value: "snapshot", + labelKey: + "portal.sources.types.webhook.fields.mode.options.snapshot", + }, + ], + }, + ], + }, ]; /** Default option values for a type's create form. */ diff --git a/frontend/editor/src/portal/mocks/handlers/sources.ts b/frontend/editor/src/portal/mocks/handlers/sources.ts index 4f4dd0e807..f5679be653 100644 --- a/frontend/editor/src/portal/mocks/handlers/sources.ts +++ b/frontend/editor/src/portal/mocks/handlers/sources.ts @@ -55,6 +55,18 @@ function seedSources(): StoredSource[] { enabled: false, owner: "data-eng@acme.com", }, + { + id: "src-webhook", + name: "Partner uploads", + type: "webhook", + options: { + webhookId: "whk_demo_5f3a9c21b7", + signingSecret: "whsec_demo_2b8e1d47a9f60c35", + mode: "consume", + }, + enabled: true, + owner: "you@acme.com", + }, ]; } @@ -76,6 +88,7 @@ const docCounts: Record< "src-contracts": { total: 12840, last24h: 96, last30d: 2310 }, "src-archive": { total: 1180, last24h: 0, last30d: 0 }, "src-legacy": { total: 48600, last24h: 0, last30d: 0 }, + "src-webhook": { total: 3120, last24h: 24, last30d: 640 }, }; function docsFor(id: string): { @@ -96,14 +109,23 @@ function nextId(): string { return `src_${Date.now().toString(36)}_${idCounter}`; } +/** A demo token for a newly created webhook's server-generated id/secret (mock only). */ +function randomToken(): string { + return ( + Math.random().toString(36).slice(2) + Math.random().toString(36).slice(2) + ); +} + function refsFor(id: string): SourcePolicyRef[] { return references[id] ?? []; } +/** Mirror the backend's secret redaction so config rows never surface a secret in the overview. */ function configRows(options: Record) { + const secret = /secret|password|token/i; return Object.entries(options).map(([key, value]) => ({ label: key.charAt(0).toUpperCase() + key.slice(1), - value: String(value), + value: secret.test(key) ? "********" : String(value), })); } @@ -120,6 +142,7 @@ function toSourceView( refs: SourcePolicyRef[], ): SourceView { const docs = docsFor(source.id); + const webhookId = source.options.webhookId; return { id: source.id, name: source.name, @@ -131,6 +154,10 @@ function toSourceView( docsTotal: docs.total, docs24h: docs.last24h, docs30d: docs.last30d, + webhookPath: + source.type === "webhook" && typeof webhookId === "string" + ? `/api/v1/webhooks/${webhookId}` + : null, }; } @@ -207,8 +234,19 @@ export const sourcesHandlers = [ ? store.find((s) => s.id === incoming.id) : undefined; const id = existing?.id ?? nextId(); + // A new webhook has its routing id + signing secret minted server-side, then revealed once on + // this create response (the store keeps them; later reads mask the secret). + let options = incoming.options ?? {}; + if (incoming.type === "webhook" && !existing && !options.webhookId) { + options = { + ...options, + webhookId: `whk_${randomToken()}`, + signingSecret: `whsec_${randomToken()}`, + }; + } const saved: StoredSource = { ...incoming, + options, id, owner: existing?.owner ?? "you@acme.com", }; diff --git a/frontend/editor/src/portal/views/Sources.css b/frontend/editor/src/portal/views/Sources.css index d8cc6897af..3ee0078a8d 100644 --- a/frontend/editor/src/portal/views/Sources.css +++ b/frontend/editor/src/portal/views/Sources.css @@ -419,6 +419,28 @@ gap: 1rem; } +/* Read-only value + copy button, e.g. the webhook delivery URL and signing secret. */ +.portal-sources__copy-row { + display: flex; + gap: 0.5rem; + align-items: center; +} + +.portal-sources__copy-row > :first-child { + flex: 1 1 auto; + min-width: 0; +} + +.portal-sources__webhook-url { + overflow-x: auto; + white-space: nowrap; + padding: 0.4rem 0.6rem; + border-radius: var(--radius-sm, 0.375rem); + background: var(--color-surface-2, rgba(127, 127, 127, 0.12)); + font-size: 0.8125rem; + line-height: 1.4; +} + .portal-sources__wizard-lead { margin: 0; font-size: 0.875rem;