From e459775740993452ea2175b0f3ec719c27ebc13b Mon Sep 17 00:00:00 2001 From: Hein Puth Date: Wed, 15 Jul 2026 04:24:23 +0200 Subject: [PATCH] feat: add webhook thought ingestion --- .gitea/workflows/ci.yml | 14 +++ README.md | 32 +++++ internal/app/app.go | 1 + internal/app/webhooks.go | 214 +++++++++++++++++++++++++++++++++ internal/app/webhooks_test.go | 86 +++++++++++++ internal/metadata/normalize.go | 22 ++++ internal/store/thoughts.go | 16 +++ internal/types/thought.go | 34 ++++-- 8 files changed, 406 insertions(+), 13 deletions(-) create mode 100644 internal/app/webhooks.go create mode 100644 internal/app/webhooks_test.go diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml index 0ce86d1..cc0862e 100644 --- a/.gitea/workflows/ci.yml +++ b/.gitea/workflows/ci.yml @@ -34,6 +34,20 @@ jobs: - name: Tidy modules run: go mod tidy + - name: Set up Node + uses: actions/setup-node@v4 + with: + node-version: 'lts/*' + + - name: Install pnpm + run: npm install -g pnpm + + - name: Build UI + run: | + cd ui + pnpm install --frozen-lockfile + pnpm run build + - name: Run tests run: go test ./... diff --git a/README.md b/README.md index 580e1af..f6c3288 100644 --- a/README.md +++ b/README.md @@ -74,6 +74,38 @@ The AMCS directory is used to store configuration and code for the Avalon Memory | `describe_tools` | List all available MCP tools with names, descriptions, categories, and model-authored usage notes; call this at the start of a session to orient yourself | | `annotate_tool` | Persist your own usage notes for a specific tool; notes are returned by `describe_tools` in future sessions | +## Webhook ingestion + +External automation can create thoughts without speaking MCP by posting JSON to `POST /webhooks/thoughts`. The endpoint is protected by the same AMCS authentication middleware as MCP and file uploads, so pass one configured API key via `x-brain-key`, an authorization bearer token header, or another enabled auth method. + +Example: + +```bash +curl -X POST http://localhost:8080/webhooks/thoughts \ + -H 'Content-Type: application/json' \ + -H 'x-brain-key: ' \ + -H 'Idempotency-Key: n8n-run-123' \ + -d '{ + "content": "External system observed build failure on main", + "project": "amcs", + "source": "n8n", + "type": "task", + "topics": ["ci", "webhook"], + "metadata": {"workflow": "ci-monitor", "run_id": "123"} + }' +``` + +Payload fields: + +- `content` is required and becomes the thought content. +- `project` is optional; when present it must match an existing AMCS project. +- `source`, `type`, `topics`, `people`, `action_items`, and `dates_mentioned` are normalized into the standard thought metadata schema. Unknown `type` values fall back to `observation`. +- `metadata` or `source_metadata` may contain safe source-specific JSON values; unsupported values and overly deep objects are dropped rather than persisted. +- `idempotency_key` or the `Idempotency-Key` header can be supplied to make repeated webhook deliveries return the existing thought with `duplicate: true`. +- `external_id` is stored under `metadata.webhook.external_id` for source-side traceability. + +Successful new ingestion returns `201` with the created thought. Duplicate idempotency-key delivery returns `200` and the previously created thought. Invalid JSON, missing content, missing/unknown projects, or unauthenticated requests are rejected before persistence. Metadata and embedding enrichment are queued after the thought is stored. + ## Learnings Learnings are curated, structured memory records for durable insights you want to keep distinct from raw thoughts. Use them for normalized lessons, decisions, and evidence-backed findings that should be easy to retrieve and review over time. diff --git a/internal/app/app.go b/internal/app/app.go index 8a65214..ca2fa65 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -238,6 +238,7 @@ func routes(logger *slog.Logger, cfg *config.Config, info buildinfo.Info, db *st } mux.Handle("/files", authMiddleware(fileHandler(filesTool))) mux.Handle("/files/{id}", authMiddleware(fileHandler(filesTool))) + mux.Handle("/webhooks/thoughts", authMiddleware(newWebhookThoughtHandler(db, embeddings, cfg.Capture, enrichmentRetryer, backfillTool))) mux.HandleFunc("/.well-known/oauth-authorization-server", oauthMetadataHandler()) mux.HandleFunc("/api/oauth/register", oauthRegisterHandler(dynClients, logger)) mux.HandleFunc("/api/oauth/authorize", oauthAuthorizeHandler(dynClients, authCodes, logger)) diff --git a/internal/app/webhooks.go b/internal/app/webhooks.go new file mode 100644 index 0000000..309f014 --- /dev/null +++ b/internal/app/webhooks.go @@ -0,0 +1,214 @@ +package app + +import ( + "encoding/json" + "errors" + "net/http" + "strings" + "time" + + "git.warky.dev/wdevs/amcs/internal/ai" + "git.warky.dev/wdevs/amcs/internal/config" + "git.warky.dev/wdevs/amcs/internal/metadata" + "git.warky.dev/wdevs/amcs/internal/store" + "git.warky.dev/wdevs/amcs/internal/tools" + thoughttypes "git.warky.dev/wdevs/amcs/internal/types" +) + +const maxWebhookBodyBytes = 1 << 20 + +type webhookThoughtRequest struct { + Content string `json:"content"` + Project string `json:"project,omitempty"` + Source string `json:"source,omitempty"` + Type string `json:"type,omitempty"` + Topics []string `json:"topics,omitempty"` + People []string `json:"people,omitempty"` + ActionItems []string `json:"action_items,omitempty"` + DatesMentioned []string `json:"dates_mentioned,omitempty"` + Metadata map[string]any `json:"metadata,omitempty"` + SourceMetadata map[string]any `json:"source_metadata,omitempty"` + IDempotencyKey string `json:"idempotency_key,omitempty"` + ExternalID string `json:"external_id,omitempty"` +} + +type webhookThoughtResponse struct { + Thought thoughttypes.Thought `json:"thought"` + Duplicate bool `json:"duplicate"` + WebhookMeta thoughttypes.WebhookMetadata `json:"webhook"` +} + +type webhookThoughtHandler struct { + store *store.DB + embeddings *ai.EmbeddingRunner + capture config.CaptureConfig + retryer tools.MetadataQueuer + embedRetryer tools.EmbeddingQueuer +} + +func newWebhookThoughtHandler(db *store.DB, embeddings *ai.EmbeddingRunner, capture config.CaptureConfig, retryer tools.MetadataQueuer, embedRetryer tools.EmbeddingQueuer) http.Handler { + return &webhookThoughtHandler{store: db, embeddings: embeddings, capture: capture, retryer: retryer, embedRetryer: embedRetryer} +} + +func (h *webhookThoughtHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/webhooks/thoughts" { + http.NotFound(w, r) + return + } + if r.Method != http.MethodPost { + w.Header().Set("Allow", http.MethodPost) + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + return + } + + r.Body = http.MaxBytesReader(w, r.Body, maxWebhookBodyBytes) + in, err := parseWebhookThoughtRequest(r) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + webhookMeta := buildWebhookMetadata(in, r.Header.Get("Idempotency-Key"), time.Now().UTC()) + if webhookMeta.IDempotencyKey != "" { + if existing, err := h.store.GetThoughtByWebhookIDempotencyKey(r.Context(), webhookMeta.IDempotencyKey); err == nil { + writeWebhookThoughtResponse(w, http.StatusOK, webhookThoughtResponse{Thought: existing, Duplicate: true, WebhookMeta: webhookMeta}) + return + } + } + + projectID, err := h.resolveWebhookProject(r, in.Project) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + thought := thoughttypes.Thought{ + Content: strings.TrimSpace(in.Content), + Metadata: normalizeWebhookThoughtMetadata(in, webhookMeta, h.capture), + ProjectID: projectID, + } + created, err := h.store.InsertThought(r.Context(), thought, h.embeddings.PrimaryModel()) + if err != nil { + http.Error(w, "insert thought: "+err.Error(), http.StatusInternalServerError) + return + } + if projectID != nil { + _ = h.store.TouchProject(r.Context(), *projectID) + } + if h.retryer != nil { + h.retryer.QueueThought(created.ID) + } + if h.embedRetryer != nil { + h.embedRetryer.QueueThought(r.Context(), created.ID, created.Content) + } + + writeWebhookThoughtResponse(w, http.StatusCreated, webhookThoughtResponse{Thought: created, WebhookMeta: webhookMeta}) +} + +func parseWebhookThoughtRequest(r *http.Request) (webhookThoughtRequest, error) { + if !strings.Contains(r.Header.Get("Content-Type"), "application/json") { + return webhookThoughtRequest{}, errors.New("webhook requires application/json") + } + defer r.Body.Close() + decoder := json.NewDecoder(r.Body) + decoder.DisallowUnknownFields() + var in webhookThoughtRequest + if err := decoder.Decode(&in); err != nil { + return webhookThoughtRequest{}, err + } + if strings.TrimSpace(in.Content) == "" { + return webhookThoughtRequest{}, errors.New("content is required") + } + return in, nil +} + +func (h *webhookThoughtHandler) resolveWebhookProject(r *http.Request, projectName string) (*int64, error) { + projectName = strings.TrimSpace(projectName) + if projectName == "" { + return nil, nil + } + project, err := h.store.GetProject(r.Context(), projectName) + if err != nil { + return nil, err + } + return &project.NumericID, nil +} + +func buildWebhookMetadata(in webhookThoughtRequest, headerKey string, now time.Time) thoughttypes.WebhookMetadata { + sourceMetadata := in.SourceMetadata + if len(sourceMetadata) == 0 { + sourceMetadata = in.Metadata + } + return thoughttypes.WebhookMetadata{ + ReceivedAt: now.Format(time.RFC3339), + IDempotencyKey: firstNonEmpty(in.IDempotencyKey, headerKey), + ExternalID: strings.TrimSpace(in.ExternalID), + SourceMetadata: sanitizeWebhookMetadata(sourceMetadata), + } +} + +func normalizeWebhookThoughtMetadata(in webhookThoughtRequest, webhookMeta thoughttypes.WebhookMetadata, capture config.CaptureConfig) thoughttypes.ThoughtMetadata { + return metadata.Normalize(thoughttypes.ThoughtMetadata{ + People: in.People, + ActionItems: in.ActionItems, + DatesMentioned: in.DatesMentioned, + Topics: in.Topics, + Type: in.Type, + Source: firstNonEmpty(in.Source, "webhook"), + Webhook: &webhookMeta, + }, capture) +} + +func sanitizeWebhookMetadata(in map[string]any) map[string]any { + if len(in) == 0 { + return nil + } + out := make(map[string]any, len(in)) + for key, value := range in { + key = strings.TrimSpace(key) + if key == "" { + continue + } + if sanitized, ok := sanitizeWebhookMetadataValue(value, 0); ok { + out[key] = sanitized + } + } + if len(out) == 0 { + return nil + } + return out +} + +func sanitizeWebhookMetadataValue(value any, depth int) (any, bool) { + if depth > 3 { + return nil, false + } + switch v := value.(type) { + case nil, bool, float64, string: + return v, true + case []any: + if len(v) > 50 { + v = v[:50] + } + out := make([]any, 0, len(v)) + for _, item := range v { + if sanitized, ok := sanitizeWebhookMetadataValue(item, depth+1); ok { + out = append(out, sanitized) + } + } + return out, true + case map[string]any: + if len(v) > 50 { + return nil, false + } + return sanitizeWebhookMetadata(v), true + default: + return nil, false + } +} + +func writeWebhookThoughtResponse(w http.ResponseWriter, status int, out webhookThoughtResponse) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(out) +} diff --git a/internal/app/webhooks_test.go b/internal/app/webhooks_test.go new file mode 100644 index 0000000..961b76a --- /dev/null +++ b/internal/app/webhooks_test.go @@ -0,0 +1,86 @@ +package app + +import ( + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "git.warky.dev/wdevs/amcs/internal/config" +) + +func TestParseWebhookThoughtRequestRequiresJSON(t *testing.T) { + req := httptestRequest("text/plain", `{"content":"hello"}`) + + _, err := parseWebhookThoughtRequest(req) + if err == nil { + t.Fatal("expected error for non-json content type") + } +} + +func TestParseWebhookThoughtRequestRequiresContent(t *testing.T) { + req := httptestRequest("application/json", `{"source":"n8n"}`) + + _, err := parseWebhookThoughtRequest(req) + if err == nil || !strings.Contains(err.Error(), "content is required") { + t.Fatalf("error = %v, want content required", err) + } +} + +func TestBuildWebhookMetadataUsesHeaderIdempotencyAndSanitizesMetadata(t *testing.T) { + now := time.Date(2026, 7, 15, 4, 0, 0, 0, time.UTC) + got := buildWebhookMetadata(webhookThoughtRequest{ + ExternalID: " ext-1 ", + Metadata: map[string]any{ + "service": "n8n", + "unsafe": struct{}{}, + "nested": map[string]any{"ok": true}, + }, + }, " key-1 ", now) + + if got.IDempotencyKey != "key-1" { + t.Fatalf("IDempotencyKey = %q, want key-1", got.IDempotencyKey) + } + if got.ExternalID != "ext-1" { + t.Fatalf("ExternalID = %q, want ext-1", got.ExternalID) + } + if got.ReceivedAt != "2026-07-15T04:00:00Z" { + t.Fatalf("ReceivedAt = %q", got.ReceivedAt) + } + if _, ok := got.SourceMetadata["unsafe"]; ok { + t.Fatal("unsafe metadata value was not removed") + } + if got.SourceMetadata["service"] != "n8n" { + t.Fatalf("service metadata = %#v", got.SourceMetadata["service"]) + } +} + +func TestNormalizeWebhookThoughtMetadata(t *testing.T) { + webhookMeta := buildWebhookMetadata(webhookThoughtRequest{IDempotencyKey: "abc"}, "", time.Date(2026, 7, 15, 4, 0, 0, 0, time.UTC)) + got := normalizeWebhookThoughtMetadata(webhookThoughtRequest{ + Source: "github", + Type: "task", + Topics: []string{"ci", "ci", ""}, + People: []string{" Sam "}, + }, webhookMeta, config.CaptureConfig{}) + + if got.Source != "github" { + t.Fatalf("Source = %q, want github", got.Source) + } + if got.Type != "task" { + t.Fatalf("Type = %q, want task", got.Type) + } + if len(got.Topics) != 1 || got.Topics[0] != "ci" { + t.Fatalf("Topics = %#v, want [ci]", got.Topics) + } + if got.Webhook == nil || got.Webhook.IDempotencyKey != "abc" { + t.Fatalf("Webhook = %#v, want idempotency key abc", got.Webhook) + } +} + +func httptestRequest(contentType, body string) *http.Request { + req := httptest.NewRequest(http.MethodPost, "/webhooks/thoughts", strings.NewReader(body)) + req.Header.Set("Content-Type", contentType) + return req +} diff --git a/internal/metadata/normalize.go b/internal/metadata/normalize.go index 51cd6dd..3c2c4c5 100644 --- a/internal/metadata/normalize.go +++ b/internal/metadata/normalize.go @@ -53,6 +53,7 @@ func Normalize(in thoughttypes.ThoughtMetadata, capture config.CaptureConfig) th Type: normalizeType(in.Type), Source: normalizeSource(in.Source), Attachments: normalizeAttachments(in.Attachments), + Webhook: normalizeWebhook(in.Webhook), MetadataStatus: normalizeMetadataStatus(in.MetadataStatus), MetadataUpdatedAt: strings.TrimSpace(in.MetadataUpdatedAt), MetadataLastAttemptedAt: strings.TrimSpace(in.MetadataLastAttemptedAt), @@ -201,10 +202,31 @@ func Merge(base, patch thoughttypes.ThoughtMetadata, capture config.CaptureConfi if len(patch.Attachments) > 0 { merged.Attachments = append(append([]thoughttypes.ThoughtAttachment{}, merged.Attachments...), patch.Attachments...) } + if patch.Webhook != nil { + merged.Webhook = patch.Webhook + } return Normalize(merged, capture) } +func normalizeWebhook(value *thoughttypes.WebhookMetadata) *thoughttypes.WebhookMetadata { + if value == nil { + return nil + } + out := &thoughttypes.WebhookMetadata{ + ReceivedAt: strings.TrimSpace(value.ReceivedAt), + IDempotencyKey: strings.TrimSpace(value.IDempotencyKey), + ExternalID: strings.TrimSpace(value.ExternalID), + } + if len(value.SourceMetadata) > 0 { + out.SourceMetadata = value.SourceMetadata + } + if out.ReceivedAt == "" && out.IDempotencyKey == "" && out.ExternalID == "" && len(out.SourceMetadata) == 0 { + return nil + } + return out +} + func normalizeAttachments(values []thoughttypes.ThoughtAttachment) []thoughttypes.ThoughtAttachment { seen := make(map[string]struct{}, len(values)) result := make([]thoughttypes.ThoughtAttachment, 0, len(values)) diff --git a/internal/store/thoughts.go b/internal/store/thoughts.go index 942a8dc..0b65e5d 100644 --- a/internal/store/thoughts.go +++ b/internal/store/thoughts.go @@ -68,6 +68,22 @@ func (db *DB) InsertThought(ctx context.Context, thought thoughttypes.Thought, e return created, nil } +func (db *DB) GetThoughtByWebhookIDempotencyKey(ctx context.Context, key string) (thoughttypes.Thought, error) { + row := db.pool.QueryRow(ctx, ` + select id, guid, content, metadata, project_id, archived_at, created_at, updated_at + from thoughts + where metadata->'webhook'->>'idempotency_key' = $1 + order by created_at desc + limit 1 + `, strings.TrimSpace(key)) + + var model generatedmodels.ModelPublicThoughts + if err := row.Scan(&model.ID, &model.GUID, &model.Content, &model.Metadata, &model.ProjectID, &model.ArchivedAt, &model.CreatedAt, &model.UpdatedAt); err != nil { + return thoughttypes.Thought{}, err + } + return thoughtFromModel(model) +} + func (db *DB) SearchThoughts(ctx context.Context, embedding []float32, embeddingModel string, threshold float64, limit int, filter map[string]any) ([]thoughttypes.SearchResult, error) { filterJSON, err := json.Marshal(filter) if err != nil { diff --git a/internal/types/thought.go b/internal/types/thought.go index 4f01af9..0b278af 100644 --- a/internal/types/thought.go +++ b/internal/types/thought.go @@ -14,12 +14,20 @@ type ThoughtMetadata struct { Type string `json:"type"` Source string `json:"source"` Attachments []ThoughtAttachment `json:"attachments,omitempty"` + Webhook *WebhookMetadata `json:"webhook,omitempty"` MetadataStatus string `json:"metadata_status,omitempty"` MetadataUpdatedAt string `json:"metadata_updated_at,omitempty"` MetadataLastAttemptedAt string `json:"metadata_last_attempted_at,omitempty"` MetadataError string `json:"metadata_error,omitempty"` } +type WebhookMetadata struct { + ReceivedAt string `json:"received_at"` + IDempotencyKey string `json:"idempotency_key,omitempty"` + ExternalID string `json:"external_id,omitempty"` + SourceMetadata map[string]any `json:"source_metadata,omitempty"` +} + type ThoughtAttachment struct { FileID uuid.UUID `json:"file_id"` Name string `json:"name"` @@ -30,19 +38,19 @@ type ThoughtAttachment struct { } type StoredFile struct { - ID int64 `json:"id"` - GUID uuid.UUID `json:"guid"` - ThoughtID *int64 `json:"thought_id,omitempty"` - ProjectID *int64 `json:"project_id,omitempty"` - Name string `json:"name"` - MediaType string `json:"media_type"` - Kind string `json:"kind"` - Encoding string `json:"encoding"` - SizeBytes int64 `json:"size_bytes"` - SHA256 string `json:"sha256"` - Content []byte `json:"-"` - CreatedAt time.Time `json:"created_at"` - UpdatedAt time.Time `json:"updated_at"` + ID int64 `json:"id"` + GUID uuid.UUID `json:"guid"` + ThoughtID *int64 `json:"thought_id,omitempty"` + ProjectID *int64 `json:"project_id,omitempty"` + Name string `json:"name"` + MediaType string `json:"media_type"` + Kind string `json:"kind"` + Encoding string `json:"encoding"` + SizeBytes int64 `json:"size_bytes"` + SHA256 string `json:"sha256"` + Content []byte `json:"-"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` } type StoredFileFilter struct {