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) }