Add webhook policy source

This commit is contained in:
Anthony Stirling
2026-07-13 16:40:48 +01:00
parent a84b375f5d
commit a9f7add87d
31 changed files with 1945 additions and 166 deletions
@@ -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
@@ -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/")
@@ -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"));
@@ -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<String, Object> prepareOptionsForSave(
Map<String, Object> 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
@@ -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<String, Object> prepareOptionsForSave(
Map<String, Object> 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<String, Object> 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<ResolvedInput> 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<Path> present = listFiles(dir);
if (config.snapshot()) {
List<ResolvedInput> 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<ResolvedInput> 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<Path> listFiles(Path dir) throws IOException {
List<Path> files = new ArrayList<>();
try (Stream<Path> 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;
}
};
}
}
@@ -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<Source> 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<String, Object> 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<InputSource> 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.
@@ -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". */
@@ -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<DetailRow> config,
long docsTotal,
long docs24h,
long docs30d) {
long docs30d,
String webhookPath) {
/** A policy that references this source. */
public record PolicyRef(String id, String name) {}
@@ -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.
*
* <p>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<String> 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;
}
}
@@ -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<String, Object> 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 + "]";
}
}
@@ -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);
}
}
@@ -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.
*
* <p>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=<hex>}. 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=<hex>' in the X-Stirling-Signature header. Returns 202 once"
+ " the document is spooled for the referencing policies.")
public ResponseEntity<WebhookDeliveryResponse> 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) {}
}
@@ -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=<hex>} 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=<lowercase hex>}.
*/
public static String sign(String signingSecret, byte[] body) {
return PREFIX + HexFormat.of().formatHex(hmac(signingSecret, body));
}
/**
* Whether {@code presented} (a {@code sha256=<hex>} 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);
}
}
}
@@ -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;
}
}
@@ -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() {}
@@ -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<ResolvedInput> 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<ResolvedInput> 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<ResolvedInput> 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<ResolvedInput> 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<String, Object> 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<String, Object> other = source.prepareOptionsForSave(Map.of(), true);
assertNotEquals(id, other.get(WebhookConfig.WEBHOOK_ID_OPTION).toString());
}
@Test
void prepareLeavesAnExistingWebhookUntouchedOnEdit() {
Map<String, Object> existing =
Map.of("webhookId", WEBHOOK_ID, "signingSecret", "keepme", "mode", "snapshot");
Map<String, Object> 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<String> present = new ArrayList<>();
@Override
public boolean claim(String identity, String gate, Supplier<String> 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<String> identities) {
present.addAll(identities);
}
}
}
@@ -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
@@ -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<String> sourceIds) {
return new Policy(
id,
"hook",
"owner",
true,
trigger,
sourceIds,
List.of(new PipelineStep("/api/v1/misc/compress-pdf", Map.of())),
OutputSpec.inline());
}
}
@@ -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<WebhookDeliveryResponse> 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<Path> spooledFiles() {
Path dir = spool.dirFor(WEBHOOK_ID);
if (!Files.isDirectory(dir)) {
return List.of();
}
try (Stream<Path> entries = Files.list(dir)) {
return entries.filter(Files::isRegularFile)
.filter(p -> !p.getFileName().toString().startsWith("."))
.toList();
} catch (IOException e) {
throw new RuntimeException(e);
}
}
}
@@ -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, ""));
}
}
@@ -103,6 +103,18 @@ class SecretMaskerTest {
assertEquals("AKIAEXAMPLE", result.get("accessKeyId"));
}
@Test
@DisplayName("should mask camelCase signingSecret despite no word boundary")
void shouldMaskCamelCaseSigningSecret() {
Map<String, Object> input = Map.of("signingSecret", "shh", "webhookId", "whk_abc");
Map<String, Object> 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() {
@@ -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"
@@ -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", () => {
@@ -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),
);
}
@@ -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 {
@@ -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(
<ConnectWizard open onClose={onClose} onCreated={onCreated} />,
);
// 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", {
@@ -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<string | null>(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<StepId, string> = {
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={
<div className="portal-sources__wizard-footer">
<Button
variant="tertiary"
size="sm"
disabled={submitting}
onClick={() =>
stepIndex === 0 ? onClose() : setStepIndex((i) => i - 1)
}
>
{stepIndex === 0
? t("portal.sources.wizard.cancel")
: t("portal.sources.wizard.back")}
</Button>
<Button
size="sm"
onClick={advance}
loading={submitting}
disabled={!canContinue}
rightSection={!isLast ? <span aria-hidden>→</span> : undefined}
>
{!isLast
? t("portal.sources.wizard.continue")
: isEdit
? t("portal.sources.wizard.save")
: t("portal.sources.actions.connectSource")}
</Button>
</div>
reveal ? (
<div className="portal-sources__wizard-footer">
<Button size="sm" onClick={onClose}>
{t("portal.sources.types.webhook.reveal.done")}
</Button>
</div>
) : (
<div className="portal-sources__wizard-footer">
<Button
variant="tertiary"
size="sm"
disabled={submitting}
onClick={() =>
stepIndex === 0 ? onClose() : setStepIndex((i) => i - 1)
}
>
{stepIndex === 0
? t("portal.sources.wizard.cancel")
: t("portal.sources.wizard.back")}
</Button>
<Button
size="sm"
onClick={advance}
loading={submitting}
disabled={!canContinue}
rightSection={!isLast ? <span aria-hidden>→</span> : undefined}
>
{!isLast
? t("portal.sources.wizard.continue")
: isEdit
? t("portal.sources.wizard.save")
: t("portal.sources.actions.connectSource")}
</Button>
</div>
)
}
>
<ol className="portal-sources__steps" aria-hidden>
{steps.map((id, i) => (
<li
key={id}
className={
"portal-sources__step" +
(i === stepIndex ? " is-active" : i < stepIndex ? " is-done" : "")
}
>
<span className="portal-sources__step-mark">
{i < stepIndex ? "✓" : i + 1}
</span>
{stepLabels[id]}
</li>
))}
</ol>
{stepId === "type" && (
<div className="portal-sources__type-grid">
{OFFERED_TYPES.map((ct) => (
<Button
key={ct.type}
variant="tertiary"
className={
"portal-sources__type-card" +
(type.type === ct.type ? " is-selected" : "")
}
onClick={() => chooseType(ct)}
>
<span className="portal-sources__type-icon" aria-hidden>
{sourceTypeMeta(ct.type).icon}
</span>
<span className="portal-sources__type-name">
{t(ct.labelKey)}
</span>
</Button>
))}
</div>
)}
{stepId === "configure" && (
{reveal && (
<div className="portal-sources__wizard-body">
<FormField label={t("portal.sources.wizard.name")} required>
<Input
value={name}
placeholder={t("portal.sources.wizard.namePlaceholder")}
onChange={(e) => setName(e.target.value)}
/>
</FormField>
{type.fields.map((field) => (
<FormField
key={field.key}
label={t(field.labelKey)}
helperText={
field.helperTextKey ? t(field.helperTextKey) : undefined
}
required={field.required}
>
{field.control === "select" ? (
<Select
value={options[field.key] ?? ""}
options={(field.options ?? []).map((o) => ({
value: o.value,
label: t(o.labelKey),
}))}
onChange={(value) =>
setOptions((o) => ({ ...o, [field.key]: value ?? "" }))
}
/>
) : (
<Input
type={field.control === "password" ? "password" : undefined}
value={options[field.key] ?? ""}
placeholder={
field.placeholderKey ? t(field.placeholderKey) : undefined
}
onChange={(e) =>
setOptions((o) => ({ ...o, [field.key]: e.target.value }))
}
/>
)}
</FormField>
))}
</div>
)}
{stepId === "review" && (
<div className="portal-sources__wizard-body">
<div className="portal-sources__stat-grid">
<StatTile
label={t("portal.sources.wizard.name")}
value={name || "—"}
/>
<StatTile
label={t("portal.sources.wizard.type")}
value={t(type.labelKey)}
/>
{type.fields.map((field) => (
<StatTile
key={field.key}
label={t(field.labelKey)}
value={
field.control === "password" && options[field.key]
? "********"
: options[field.key] || "—"
}
<Banner
tone="warning"
description={t("portal.sources.types.webhook.reveal.secretWarning")}
/>
<FormField label={t("portal.sources.types.webhook.reveal.url")}>
<div className="portal-sources__copy-row">
<Input
value={webhookUrl}
readOnly
onFocus={(e) => e.currentTarget.select()}
/>
))}
</div>
{error && <Banner tone="danger" description={error} />}
<Button
size="sm"
variant="tertiary"
onClick={() => copy(webhookUrl)}
>
{t("portal.sources.types.webhook.reveal.copy")}
</Button>
</div>
</FormField>
<FormField
label={t("portal.sources.types.webhook.reveal.secret")}
helperText={t("portal.sources.types.webhook.reveal.secretHelp")}
>
<div className="portal-sources__copy-row">
<Input
value={revealSecret}
readOnly
onFocus={(e) => e.currentTarget.select()}
/>
<Button
size="sm"
variant="tertiary"
onClick={() => copy(revealSecret)}
>
{t("portal.sources.types.webhook.reveal.copy")}
</Button>
</div>
</FormField>
<p className="portal-sources__muted">
{t("portal.sources.types.webhook.reveal.usage")}
</p>
</div>
)}
{!reveal && (
<>
<ol className="portal-sources__steps" aria-hidden>
{steps.map((id, i) => (
<li
key={id}
className={
"portal-sources__step" +
(i === stepIndex
? " is-active"
: i < stepIndex
? " is-done"
: "")
}
>
<span className="portal-sources__step-mark">
{i < stepIndex ? "✓" : i + 1}
</span>
{stepLabels[id]}
</li>
))}
</ol>
{stepId === "type" && (
<div className="portal-sources__type-grid">
{OFFERED_TYPES.map((ct) => (
<Button
key={ct.type}
variant="tertiary"
className={
"portal-sources__type-card" +
(type.type === ct.type ? " is-selected" : "")
}
onClick={() => chooseType(ct)}
>
<span className="portal-sources__type-icon" aria-hidden>
{sourceTypeMeta(ct.type).icon}
</span>
<span className="portal-sources__type-name">
{t(ct.labelKey)}
</span>
</Button>
))}
</div>
)}
{stepId === "configure" && (
<div className="portal-sources__wizard-body">
<FormField label={t("portal.sources.wizard.name")} required>
<Input
value={name}
placeholder={t("portal.sources.wizard.namePlaceholder")}
onChange={(e) => setName(e.target.value)}
/>
</FormField>
{type.fields.map((field) => (
<FormField
key={field.key}
label={t(field.labelKey)}
helperText={
field.helperTextKey ? t(field.helperTextKey) : undefined
}
required={field.required}
>
{field.control === "select" ? (
<Select
value={options[field.key] ?? ""}
options={(field.options ?? []).map((o) => ({
value: o.value,
label: t(o.labelKey),
}))}
onChange={(value) =>
setOptions((o) => ({ ...o, [field.key]: value ?? "" }))
}
/>
) : (
<Input
type={
field.control === "password" ? "password" : undefined
}
value={options[field.key] ?? ""}
placeholder={
field.placeholderKey
? t(field.placeholderKey)
: undefined
}
onChange={(e) =>
setOptions((o) => ({
...o,
[field.key]: e.target.value,
}))
}
/>
)}
</FormField>
))}
</div>
)}
{stepId === "review" && (
<div className="portal-sources__wizard-body">
<div className="portal-sources__stat-grid">
<StatTile
label={t("portal.sources.wizard.name")}
value={name || "—"}
/>
<StatTile
label={t("portal.sources.wizard.type")}
value={t(type.labelKey)}
/>
{type.fields.map((field) => (
<StatTile
key={field.key}
label={t(field.labelKey)}
value={
field.control === "password" && options[field.key]
? "********"
: options[field.key] || "—"
}
/>
))}
</div>
{error && <Banner tone="danger" description={error} />}
</div>
)}
</>
)}
</Modal>
);
}
@@ -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 (
<div className="portal-sources__detail">
{isEditor && (
@@ -30,6 +37,29 @@ export function SourceDetailPanel({
</p>
)}
{isWebhook && webhookUrl && (
<div className="portal-sources__detail-section">
<span className="portal-sources__detail-heading">
{t("portal.sources.types.webhook.detail.deliveryUrl")}
</span>
<div className="portal-sources__copy-row">
<code className="portal-sources__webhook-url">{webhookUrl}</code>
<Button
size="sm"
variant="tertiary"
onClick={() =>
webhookUrl && void navigator.clipboard?.writeText(webhookUrl)
}
>
{t("portal.sources.types.webhook.reveal.copy")}
</Button>
</div>
<p className="portal-sources__muted">
{t("portal.sources.types.webhook.detail.secretNote")}
</p>
</div>
)}
{source.config.length > 0 && (
<div className="portal-sources__stat-grid">
{source.config.map((row) => (
@@ -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<string, SourceTypeMeta> = {
folder: {
labelKey: "portal.sources.types.folder.label",
@@ -38,6 +41,11 @@ const SOURCE_TYPE_META: Record<string, SourceTypeMeta> = {
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. */
@@ -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<string, unknown>) {
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",
};
@@ -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;