mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-07-25 10:21:47 +00:00
* feat(webhooks): HTTP webhooks to trigger agents with HMAC auth and durable callbacks
Add multi-tenant HTTP webhook endpoints for agent triggering:
- /v1/webhooks/message: send messages to channels
- /v1/webhooks/llm: sync/async LLM prompts with HMAC-signed callbacks
- HMAC-256 + bearer token authentication
- Rate limiting and tenant isolation
- Durable callback worker with exponential backoff
- PG 000056 + SQLite schema v25 migrations
- Unit + integration tests, P0 tenant isolation invariants
- Channel media capability helpers for attachment routing
- Comprehensive webhook documentation and i18n strings
* fix(webhooks): address post-review findings (K1-K10)
Comprehensive post-merge fixes addressing 10 blocking code review issues
and 2 adversarial re-audit findings in webhook-agent-triggering feature:
K1: Fix auth middleware tenant context lookup sequencing — move
tenant context injection before authenticate() call to prevent
unscoped secret lookups.
K2: Canonicalize JSON payload format for jsonb compatibility across
PostgreSQL and SQLite — ensure consistent serialization without
whitespace variance to prevent hash mismatches.
K3: Add fail-closed JSON parsing in body hash extraction with explicit
error handling for malformed payloads before HMAC verification.
K4: Fix worker queue wedge by properly draining slot reservations
when delivery succeeds, preventing permanent slot occupancy.
K5: Implement lease-token optimistic concurrency control to prevent
duplicate webhook delivery under high concurrency or retry storms.
K6: Add AES-256-GCM encrypted secret storage at rest with fail-fast
skip-mount when GOCLAW_ENCRYPTION_KEY environment variable unset.
K7: Implement IP allowlist enforcement supporting both CIDR ranges
and exact IP matching with proper X-Forwarded-For parsing.
K8: Add HMAC replay nonce cache (5min expiry, non-blocking async flush)
to prevent request replay attacks on webhook handler.
K9: Fix invariant test schema selection — replace hardcoded assumption
with explicit schema name from config to support multi-schema testing.
K10: Consolidate rate limiters into single shared instance to prevent
per-endpoint limiter starvation and ensure fair rate limiting.
New database migrations:
- 000057: webhook_calls.lease_token for optimistic concurrency
- 000058: webhooks.encrypted_secret_key for AES-256-GCM encryption
New i18n keys: MsgWebhookIPDenied, MsgWebhookEncryptionUnavailable
(with English, Vietnamese, Chinese translations).
New modules:
- internal/http/webhooks_payload.go: JSON canonicalization + body hash
- internal/http/webhooks_nonce.go: Replay nonce cache implementation
- internal/http/webhooks_idempotency_test.go: Integration tests
Documentation updates:
- docs/webhooks.md: §13-14 security sections, encryption flow
- docs/00-architecture-overview.md: webhook subsystem security overview
- docs/codebase-summary.md: webhook security patterns
- docs/project-changelog.md: webhook fixes changelog
Test coverage: 53 webhook tests + 4 P0 invariant tests all passing.
No tenant isolation violations. All security gates enforced.
* docs(journals): webhook feature ship + fix cycle entries
* fix(webhooks): address Claude review findings
- webhooks_llm.go: remove misleading ptr() helper; use &completedAt
pattern for error-path audit rows (matches success path)
- webhooks_auth.go: wrap TouchLastUsed context in WithoutCancel so
background DB update isn't cancelled when HTTP response completes
- store GetByIDUnscoped (PG+SQLite): add NOT revoked / revoked = 0
filter for defense-in-depth parity with GetByHashUnscoped
- webhooks/sign.go: fix package doc — HMAC key is raw plaintext
secret bytes, not hex-decoded SHA-256
- webhooks_admin.go: check auth before encKey guard to avoid leaking
config state to unauthenticated callers
- webhooks_ratelimit.go: two-phase Load→LoadOrStore to avoid per-call
entry allocation on the hot path
* docs(webhooks): fix Sign() function doc to match actual key input
Function-level comment still referenced hex-decoded SecretHash after
the package-level doc was corrected. Align with actual caller usage
([]byte(rawSecret)).
* fix(webhooks): use WithoutCancel for worker execute DB updates
Terminal status writes in execute() ran through the worker main-loop
ctx, which is cancelled on graceful shutdown. If the outbound send
completed but the status update raced with shutdown, the row stayed
in 'running' and got re-delivered via reclaimStale. WithoutCancel
lets the DB write survive worker cancellation while preserving
propagated values (tenant ID, etc.).
* fix(webhooks): move tctx init before panic defer in worker execute
Panic recovery called updateRetry with raw ctx (no tenant ID), making
requireTenantID fail and the reset-to-retry DB write silently drop.
Row stayed 'running' until reclaimStale (~90s delay). Init tctx first
so defer closure captures tenant-scoped non-cancellable context.
* fix(webhooks): pass tenant-scoped tctx to invokeAgent in worker
execute() was passing the raw worker-loop ctx (no tenant ID) to
invokeAgent → router.Get → PGAgentStore.GetByID. GetByID reads
TenantIDFromContext which returned uuid.Nil, making every lookup
return 'agent not found'. Async LLM webhook calls silently failed
all retries. Pass tctx (already tenant-scoped + WithoutCancel) so
the router resolves the agent correctly.
* fix(tests): resolve integration test compile errors
- Remove duplicate contains() in mcp_grant_revoke_test.go (already
defined in tts_gemini_live_test.go)
- Update webhooks_admin_test.go RotateSecret call to match current
5-arg signature (newSecretHash, newPrefix, newEncryptedSecret)
* fix(webhooks): default nil scopes/ip_allowlist to empty slice in Create
PG columns are NOT NULL DEFAULT '{}'. Explicit NULL from pqStringArray(nil)
violated the constraint, breaking TestWebhookAdminCRUD/TenantIsolation.
Coerce nil slices to empty []string{} so the default applies at the DB layer.
* chore: trigger CI on digitopvn/goclaw fork
* ci: retrigger workflows
* fix(webhooks): renumber migrations to 000059-000061 for merge train
136 lines
4.3 KiB
Go
136 lines
4.3 KiB
Go
package http
|
|
|
|
import (
|
|
"fmt"
|
|
"net"
|
|
"net/http"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nextlevelbuilder/goclaw/internal/security"
|
|
)
|
|
|
|
const (
|
|
// webhookMediaMaxBytes is the maximum allowed media file size (25 MB).
|
|
webhookMediaMaxBytes = 25 * 1024 * 1024
|
|
|
|
// webhookMediaProbeTimeout is the deadline for the HEAD probe request.
|
|
webhookMediaProbeTimeout = 15 * time.Second
|
|
)
|
|
|
|
// allowedMediaMIMETypes is the set of Content-Type values accepted for media attachments.
|
|
// Must be lowercase prefix-matched against the probed value.
|
|
var allowedMediaMIMETypes = map[string]bool{
|
|
"image/jpeg": true,
|
|
"image/png": true,
|
|
"image/gif": true,
|
|
"image/webp": true,
|
|
"video/mp4": true,
|
|
"audio/mpeg": true,
|
|
"audio/ogg": true,
|
|
"application/pdf": true,
|
|
}
|
|
|
|
// mediaProbeResult is returned by probeMediaURL on success.
|
|
type mediaProbeResult struct {
|
|
// ContentType is the canonical MIME type from the HEAD response (trimmed of params).
|
|
ContentType string
|
|
// PinnedIP is the resolved IP from SSRF validation — callers may store for logging.
|
|
PinnedIP net.IP
|
|
}
|
|
|
|
// mediaValidateError categories (callers map these to HTTP status codes).
|
|
type mediaValidateError struct {
|
|
code string // "ssrf" | "too_large" | "mime_denied"
|
|
message string
|
|
}
|
|
|
|
func (e *mediaValidateError) Error() string { return e.message }
|
|
|
|
// probeMediaURL performs SSRF validation, DNS pinning, and a HEAD request to
|
|
// verify the media URL is reachable and within size + MIME constraints.
|
|
//
|
|
// Workflow:
|
|
// 1. security.Validate(rawURL) — rejects private/loopback ranges.
|
|
// 2. Build SafeClient with pinned IP via WithPinnedIP context.
|
|
// 3. HEAD request — parse Content-Length (≤25 MB) and Content-Type (allowlist).
|
|
//
|
|
// Returns (result, nil) on success, or (*mediaValidateError, error) on failure.
|
|
// On error, the returned error is always *mediaValidateError so callers can
|
|
// switch on .code for status-code selection.
|
|
func probeMediaURL(rawURL string) (*mediaProbeResult, error) {
|
|
// Step 1: SSRF validation — resolve DNS and reject blocked CIDRs.
|
|
_, pinnedIP, err := security.Validate(rawURL)
|
|
if err != nil {
|
|
return nil, &mediaValidateError{
|
|
code: "ssrf",
|
|
message: fmt.Sprintf("media URL blocked by SSRF policy: %v", err),
|
|
}
|
|
}
|
|
|
|
// Step 2: Build SSRF-safe client with pinned IP.
|
|
client := security.NewSafeClient(webhookMediaProbeTimeout)
|
|
|
|
// Create HEAD request. Context carries the pinned IP for the safe dialer.
|
|
// We use context.Background here; the caller's request context is not passed
|
|
// to avoid cancellation from the response write path racing with the probe.
|
|
// This is acceptable — the probe has its own 15s timeout via NewSafeClient.
|
|
req, err := http.NewRequest(http.MethodHead, rawURL, nil)
|
|
if err != nil {
|
|
return nil, &mediaValidateError{
|
|
code: "ssrf",
|
|
message: fmt.Sprintf("media URL parse error: %v", err),
|
|
}
|
|
}
|
|
// Inject pinned IP into request context so SafeClient can use it.
|
|
req = req.WithContext(security.WithPinnedIP(req.Context(), pinnedIP))
|
|
|
|
// Step 3: Execute HEAD request.
|
|
resp, err := client.Do(req)
|
|
if err != nil {
|
|
return nil, &mediaValidateError{
|
|
code: "ssrf",
|
|
message: fmt.Sprintf("media HEAD probe failed: %v", err),
|
|
}
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
// Step 4: Validate Content-Length if present.
|
|
if clStr := resp.Header.Get("Content-Length"); clStr != "" {
|
|
cl, parseErr := strconv.ParseInt(clStr, 10, 64)
|
|
if parseErr == nil && cl > webhookMediaMaxBytes {
|
|
return nil, &mediaValidateError{
|
|
code: "too_large",
|
|
message: fmt.Sprintf("media file exceeds size limit (%d bytes > %d)", cl, webhookMediaMaxBytes),
|
|
}
|
|
}
|
|
}
|
|
|
|
// Step 5: Validate Content-Type against allowlist.
|
|
rawCT := resp.Header.Get("Content-Type")
|
|
mimeType := parseMIMEType(rawCT)
|
|
if !allowedMediaMIMETypes[mimeType] {
|
|
return nil, &mediaValidateError{
|
|
code: "mime_denied",
|
|
message: fmt.Sprintf("media MIME type %q is not allowed", mimeType),
|
|
}
|
|
}
|
|
|
|
return &mediaProbeResult{
|
|
ContentType: mimeType,
|
|
PinnedIP: pinnedIP,
|
|
}, nil
|
|
}
|
|
|
|
// parseMIMEType strips parameters from a Content-Type header value and returns
|
|
// the lowercase base type (e.g. "image/jpeg; charset=utf-8" → "image/jpeg").
|
|
func parseMIMEType(ct string) string {
|
|
if ct == "" {
|
|
return ""
|
|
}
|
|
// Split on ";" and take the first part.
|
|
parts := strings.SplitN(ct, ";", 2)
|
|
return strings.ToLower(strings.TrimSpace(parts[0]))
|
|
}
|