feat: add webhook thought ingestion
This commit is contained in:
@@ -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)
|
||||
}
|
||||
Reference in New Issue
Block a user