From d5ca5b915fb3ad728df4a8a243dd472ccb049cd6 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Sun, 15 Mar 2026 07:53:53 +0200 Subject: [PATCH] feat: add webhook, k8s, and attachments gc CLI subcommands Add admin CLI subcommands for managing webhooks, K8s job handlers, and running attachment garbage collection via the Unix admin socket. Server-side: - Add WebhookServiceProvider and K8sServiceProvider interfaces to admin pkg - Add 7 new command handlers: webhook.{register,list,delete}, k8s.{register,list,delete}, attachments.gc - Wire webhook and k8s services into admin.Services struct in main.go CLI-side: - Add `synapbus webhook {register,list,delete}` commands - Add `synapbus k8s {register,list,delete}` commands - Add `synapbus attachments gc` command - All commands follow existing patterns (adminRequest, printTable, printJSON) Tests: - Add cmd/synapbus/admin_test.go with 9 tests covering command registration, required flag validation, and existing command preservation Co-Authored-By: Claude Opus 4.6 --- cmd/synapbus/admin.go | 215 ++++++++++++++++++++- cmd/synapbus/admin_test.go | 210 ++++++++++++++++++++ cmd/synapbus/main.go | 2 + internal/admin/server.go | 19 ++ internal/admin/socket.go | 384 +++++++++++++++++++++++++++++++++++++ 5 files changed, 829 insertions(+), 1 deletion(-) create mode 100644 cmd/synapbus/admin_test.go diff --git a/cmd/synapbus/admin.go b/cmd/synapbus/admin.go index 4fb4bdf..f0698da 100644 --- a/cmd/synapbus/admin.go +++ b/cmd/synapbus/admin.go @@ -705,10 +705,223 @@ func addAdminCommands(rootCmd *cobra.Command) { retentionCmd.AddCommand(retentionStatusCmd) + // ----- webhook commands ----- + webhookCmd := &cobra.Command{ + Use: "webhook", + Short: "Manage webhooks", + } + + var ( + webhookRegisterURL string + webhookRegisterEvents string + webhookRegisterSecret string + webhookRegisterAgent string + ) + webhookRegisterCmd := &cobra.Command{ + Use: "register", + Short: "Register a webhook", + RunE: func(cmd *cobra.Command, args []string) error { + resp, err := adminRequest("webhook.register", map[string]string{ + "url": webhookRegisterURL, + "events": webhookRegisterEvents, + "secret": webhookRegisterSecret, + "agent_name": webhookRegisterAgent, + }) + if err != nil { + return err + } + printJSON(resp["data"]) + return nil + }, + } + webhookRegisterCmd.Flags().StringVar(&webhookRegisterURL, "url", "", "Webhook endpoint URL") + webhookRegisterCmd.Flags().StringVar(&webhookRegisterEvents, "events", "", "Comma-separated event types (e.g. message.received,channel.message)") + webhookRegisterCmd.Flags().StringVar(&webhookRegisterSecret, "secret", "", "HMAC signing secret") + webhookRegisterCmd.Flags().StringVar(&webhookRegisterAgent, "agent", "", "Agent name to hook events for") + webhookRegisterCmd.MarkFlagRequired("url") + webhookRegisterCmd.MarkFlagRequired("events") + webhookRegisterCmd.MarkFlagRequired("secret") + webhookRegisterCmd.MarkFlagRequired("agent") + + var webhookListAgent string + webhookListCmd := &cobra.Command{ + Use: "list", + Short: "List webhooks", + RunE: func(cmd *cobra.Command, args []string) error { + reqArgs := map[string]interface{}{} + if webhookListAgent != "" { + reqArgs["agent_name"] = webhookListAgent + } + resp, err := adminRequest("webhook.list", reqArgs) + if err != nil { + return err + } + rows := toMapSlice(resp["data"]) + if len(rows) == 0 { + fmt.Println("No webhooks found.") + return nil + } + printTable([]string{"ID", "AGENT", "URL", "EVENTS", "STATUS", "FAILURES", "CREATED_AT"}, toTableRows(rows, map[string]string{ + "ID": "id", "AGENT": "agent_name", "URL": "url", "EVENTS": "events", + "STATUS": "status", "FAILURES": "consecutive_failures", "CREATED_AT": "created_at", + })) + return nil + }, + } + webhookListCmd.Flags().StringVar(&webhookListAgent, "agent", "", "Filter by agent name") + + var webhookDeleteID int64 + webhookDeleteCmd := &cobra.Command{ + Use: "delete", + Short: "Delete a webhook", + RunE: func(cmd *cobra.Command, args []string) error { + resp, err := adminRequest("webhook.delete", map[string]interface{}{ + "id": webhookDeleteID, + }) + if err != nil { + return err + } + printJSON(resp["data"]) + return nil + }, + } + webhookDeleteCmd.Flags().Int64Var(&webhookDeleteID, "id", 0, "Webhook ID to delete") + webhookDeleteCmd.MarkFlagRequired("id") + + webhookCmd.AddCommand(webhookRegisterCmd, webhookListCmd, webhookDeleteCmd) + + // ----- k8s commands ----- + k8sCmd := &cobra.Command{ + Use: "k8s", + Short: "Manage Kubernetes job handlers", + } + + var ( + k8sRegisterImage string + k8sRegisterEvents string + k8sRegisterAgent string + k8sRegisterNamespace string + k8sRegisterMemory string + k8sRegisterCPU string + k8sRegisterEnv string + k8sRegisterTimeout int + ) + k8sRegisterCmd := &cobra.Command{ + Use: "register", + Short: "Register a K8s job handler", + RunE: func(cmd *cobra.Command, args []string) error { + reqArgs := map[string]interface{}{ + "image": k8sRegisterImage, + "events": k8sRegisterEvents, + "agent_name": k8sRegisterAgent, + } + if k8sRegisterNamespace != "" { + reqArgs["namespace"] = k8sRegisterNamespace + } + if k8sRegisterMemory != "" { + reqArgs["resources_memory"] = k8sRegisterMemory + } + if k8sRegisterCPU != "" { + reqArgs["resources_cpu"] = k8sRegisterCPU + } + if k8sRegisterEnv != "" { + reqArgs["env"] = k8sRegisterEnv + } + if k8sRegisterTimeout > 0 { + reqArgs["timeout_seconds"] = k8sRegisterTimeout + } + resp, err := adminRequest("k8s.register", reqArgs) + if err != nil { + return err + } + printJSON(resp["data"]) + return nil + }, + } + k8sRegisterCmd.Flags().StringVar(&k8sRegisterImage, "image", "", "Container image") + k8sRegisterCmd.Flags().StringVar(&k8sRegisterEvents, "events", "", "Comma-separated event types") + k8sRegisterCmd.Flags().StringVar(&k8sRegisterAgent, "agent", "", "Agent name") + k8sRegisterCmd.Flags().StringVar(&k8sRegisterNamespace, "namespace", "", "Kubernetes namespace (optional)") + k8sRegisterCmd.Flags().StringVar(&k8sRegisterMemory, "memory", "", "Memory resource limit (e.g. 256Mi)") + k8sRegisterCmd.Flags().StringVar(&k8sRegisterCPU, "cpu", "", "CPU resource limit (e.g. 500m)") + k8sRegisterCmd.Flags().StringVar(&k8sRegisterEnv, "env", "", "Comma-separated KEY=VALUE environment variables") + k8sRegisterCmd.Flags().IntVar(&k8sRegisterTimeout, "timeout", 300, "Job timeout in seconds") + k8sRegisterCmd.MarkFlagRequired("image") + k8sRegisterCmd.MarkFlagRequired("events") + k8sRegisterCmd.MarkFlagRequired("agent") + + var k8sListAgent string + k8sListCmd := &cobra.Command{ + Use: "list", + Short: "List K8s job handlers", + RunE: func(cmd *cobra.Command, args []string) error { + reqArgs := map[string]interface{}{} + if k8sListAgent != "" { + reqArgs["agent_name"] = k8sListAgent + } + resp, err := adminRequest("k8s.list", reqArgs) + if err != nil { + return err + } + rows := toMapSlice(resp["data"]) + if len(rows) == 0 { + fmt.Println("No K8s handlers found.") + return nil + } + printTable([]string{"ID", "AGENT", "IMAGE", "EVENTS", "NAMESPACE", "STATUS", "CREATED_AT"}, toTableRows(rows, map[string]string{ + "ID": "id", "AGENT": "agent_name", "IMAGE": "image", "EVENTS": "events", + "NAMESPACE": "namespace", "STATUS": "status", "CREATED_AT": "created_at", + })) + return nil + }, + } + k8sListCmd.Flags().StringVar(&k8sListAgent, "agent", "", "Filter by agent name") + + var k8sDeleteID int64 + k8sDeleteCmd := &cobra.Command{ + Use: "delete", + Short: "Delete a K8s job handler", + RunE: func(cmd *cobra.Command, args []string) error { + resp, err := adminRequest("k8s.delete", map[string]interface{}{ + "id": k8sDeleteID, + }) + if err != nil { + return err + } + printJSON(resp["data"]) + return nil + }, + } + k8sDeleteCmd.Flags().Int64Var(&k8sDeleteID, "id", 0, "Handler ID to delete") + k8sDeleteCmd.MarkFlagRequired("id") + + k8sCmd.AddCommand(k8sRegisterCmd, k8sListCmd, k8sDeleteCmd) + + // ----- attachments commands ----- + attachmentsCmd := &cobra.Command{ + Use: "attachments", + Short: "Manage attachments", + } + + attachmentsGCCmd := &cobra.Command{ + Use: "gc", + Short: "Run attachment garbage collection to remove orphaned files", + RunE: func(cmd *cobra.Command, args []string) error { + resp, err := adminRequest("attachments.gc", nil) + if err != nil { + return err + } + printJSON(resp["data"]) + return nil + }, + } + + attachmentsCmd.AddCommand(attachmentsGCCmd) + // ----- add persistent flag and commands to root ----- rootCmd.PersistentFlags().StringVar(&adminSocket, "socket", "./data/synapbus.sock", "Path to admin Unix socket") - rootCmd.AddCommand(userCmd, agentCmd, auditCmd, backupCmd, messagesCmd, channelsCmd, conversationsCmd, embeddingsCmd, dbCmd, retentionCmd) + rootCmd.AddCommand(userCmd, agentCmd, auditCmd, backupCmd, messagesCmd, channelsCmd, conversationsCmd, embeddingsCmd, dbCmd, retentionCmd, webhookCmd, k8sCmd, attachmentsCmd) } // toTableRows remaps []map[string]string using a header->key mapping. diff --git a/cmd/synapbus/admin_test.go b/cmd/synapbus/admin_test.go new file mode 100644 index 0000000..09bc3d3 --- /dev/null +++ b/cmd/synapbus/admin_test.go @@ -0,0 +1,210 @@ +package main + +import ( + "testing" + + "github.com/spf13/cobra" +) + +// findSubcommand finds a subcommand by name in a cobra.Command tree. +func findSubcommand(root *cobra.Command, names ...string) *cobra.Command { + cmd := root + for _, name := range names { + found := false + for _, sub := range cmd.Commands() { + if sub.Name() == name { + cmd = sub + found = true + break + } + } + if !found { + return nil + } + } + return cmd +} + +// buildTestRoot creates a root command with all admin commands registered. +func buildTestRoot() *cobra.Command { + root := &cobra.Command{Use: "synapbus"} + addAdminCommands(root) + return root +} + +func TestWebhookCommandsRegistered(t *testing.T) { + root := buildTestRoot() + + tests := []struct { + path []string + }{ + {[]string{"webhook"}}, + {[]string{"webhook", "register"}}, + {[]string{"webhook", "list"}}, + {[]string{"webhook", "delete"}}, + } + + for _, tt := range tests { + cmd := findSubcommand(root, tt.path...) + if cmd == nil { + t.Errorf("command %v not found", tt.path) + } + } +} + +func TestK8sCommandsRegistered(t *testing.T) { + root := buildTestRoot() + + tests := []struct { + path []string + }{ + {[]string{"k8s"}}, + {[]string{"k8s", "register"}}, + {[]string{"k8s", "list"}}, + {[]string{"k8s", "delete"}}, + } + + for _, tt := range tests { + cmd := findSubcommand(root, tt.path...) + if cmd == nil { + t.Errorf("command %v not found", tt.path) + } + } +} + +func TestAttachmentsCommandsRegistered(t *testing.T) { + root := buildTestRoot() + + tests := []struct { + path []string + }{ + {[]string{"attachments"}}, + {[]string{"attachments", "gc"}}, + } + + for _, tt := range tests { + cmd := findSubcommand(root, tt.path...) + if cmd == nil { + t.Errorf("command %v not found", tt.path) + } + } +} + +func TestWebhookRegisterRequiredFlags(t *testing.T) { + root := buildTestRoot() + cmd := findSubcommand(root, "webhook", "register") + if cmd == nil { + t.Fatal("webhook register command not found") + } + + requiredFlags := []string{"url", "events", "secret", "agent"} + for _, flag := range requiredFlags { + f := cmd.Flag(flag) + if f == nil { + t.Errorf("flag --%s not found on webhook register", flag) + continue + } + ann := f.Annotations + if ann == nil { + t.Errorf("flag --%s should be required", flag) + continue + } + if _, ok := ann[cobra.BashCompOneRequiredFlag]; !ok { + t.Errorf("flag --%s should be required", flag) + } + } +} + +func TestWebhookDeleteRequiredFlags(t *testing.T) { + root := buildTestRoot() + cmd := findSubcommand(root, "webhook", "delete") + if cmd == nil { + t.Fatal("webhook delete command not found") + } + + f := cmd.Flag("id") + if f == nil { + t.Fatal("flag --id not found on webhook delete") + } + ann := f.Annotations + if ann == nil { + t.Fatal("flag --id should be required") + } + if _, ok := ann[cobra.BashCompOneRequiredFlag]; !ok { + t.Fatal("flag --id should be required") + } +} + +func TestK8sRegisterRequiredFlags(t *testing.T) { + root := buildTestRoot() + cmd := findSubcommand(root, "k8s", "register") + if cmd == nil { + t.Fatal("k8s register command not found") + } + + requiredFlags := []string{"image", "events", "agent"} + for _, flag := range requiredFlags { + f := cmd.Flag(flag) + if f == nil { + t.Errorf("flag --%s not found on k8s register", flag) + continue + } + ann := f.Annotations + if ann == nil { + t.Errorf("flag --%s should be required", flag) + continue + } + if _, ok := ann[cobra.BashCompOneRequiredFlag]; !ok { + t.Errorf("flag --%s should be required", flag) + } + } +} + +func TestK8sDeleteRequiredFlags(t *testing.T) { + root := buildTestRoot() + cmd := findSubcommand(root, "k8s", "delete") + if cmd == nil { + t.Fatal("k8s delete command not found") + } + + f := cmd.Flag("id") + if f == nil { + t.Fatal("flag --id not found on k8s delete") + } + ann := f.Annotations + if ann == nil { + t.Fatal("flag --id should be required") + } + if _, ok := ann[cobra.BashCompOneRequiredFlag]; !ok { + t.Fatal("flag --id should be required") + } +} + +func TestK8sRegisterOptionalFlags(t *testing.T) { + root := buildTestRoot() + cmd := findSubcommand(root, "k8s", "register") + if cmd == nil { + t.Fatal("k8s register command not found") + } + + optionalFlags := []string{"namespace", "memory", "cpu", "env", "timeout"} + for _, flag := range optionalFlags { + f := cmd.Flag(flag) + if f == nil { + t.Errorf("optional flag --%s not found on k8s register", flag) + } + } +} + +func TestExistingCommandsStillPresent(t *testing.T) { + root := buildTestRoot() + + // Verify existing commands are not broken by our additions. + existingCmds := []string{"user", "agent", "audit", "backup", "messages", "channels", "conversations", "embeddings", "db", "retention"} + for _, name := range existingCmds { + cmd := findSubcommand(root, name) + if cmd == nil { + t.Errorf("existing command %q not found after adding new commands", name) + } + } +} diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index 5f29ad5..a02d538 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -564,6 +564,8 @@ func runServe(cmd *cobra.Command, args []string) error { if retentionWorker != nil { adminSvcs.RetentionWorker = retentionWorker } + adminSvcs.WebhookService = webhookService + adminSvcs.K8sService = k8sService adminServer := admin.NewServer(adminSocketPath, db.DB, adminSvcs, logger) if err := adminServer.Start(); err != nil { return fmt.Errorf("start admin socket: %w", err) diff --git a/internal/admin/server.go b/internal/admin/server.go index 506cb8d..643361f 100644 --- a/internal/admin/server.go +++ b/internal/admin/server.go @@ -2,6 +2,7 @@ package admin import ( + "context" "database/sql" "log/slog" "net" @@ -10,11 +11,27 @@ import ( "github.com/synapbus/synapbus/internal/attachments" "github.com/synapbus/synapbus/internal/auth" "github.com/synapbus/synapbus/internal/channels" + "github.com/synapbus/synapbus/internal/k8s" "github.com/synapbus/synapbus/internal/messaging" "github.com/synapbus/synapbus/internal/search" "github.com/synapbus/synapbus/internal/trace" + "github.com/synapbus/synapbus/internal/webhooks" ) +// WebhookServiceProvider defines the webhook operations needed by the admin socket. +type WebhookServiceProvider interface { + RegisterWebhook(ctx context.Context, agentName, url string, events []string, secret string) (*webhooks.Webhook, error) + ListWebhooks(ctx context.Context, agentName string) ([]*webhooks.Webhook, error) + DeleteWebhook(ctx context.Context, agentName string, webhookID int64) error +} + +// K8sServiceProvider defines the K8s handler operations needed by the admin socket. +type K8sServiceProvider interface { + RegisterHandler(ctx context.Context, agentName string, req k8s.RegisterHandlerRequest) (*k8s.K8sHandler, error) + ListHandlers(ctx context.Context, agentName string) ([]*k8s.K8sHandler, error) + DeleteHandler(ctx context.Context, agentName string, handlerID int64) error +} + // Services holds references to all services the admin socket can control. type Services struct { Users *auth.SQLiteUserStore @@ -27,6 +44,8 @@ type Services struct { VectorIndex *search.VectorIndex SearchService *search.Service AttachmentService *attachments.Service + WebhookService WebhookServiceProvider + K8sService K8sServiceProvider DataDir string RetentionWorker RetentionStatusProvider } diff --git a/internal/admin/socket.go b/internal/admin/socket.go index fe39e73..2b7692c 100644 --- a/internal/admin/socket.go +++ b/internal/admin/socket.go @@ -13,6 +13,7 @@ import ( "strings" "time" + "github.com/synapbus/synapbus/internal/k8s" "github.com/synapbus/synapbus/internal/messaging" "github.com/synapbus/synapbus/internal/trace" ) @@ -188,6 +189,26 @@ func (s *AdminServer) dispatch(req Request) Response { case "retention.status": return s.handleRetentionStatus(ctx) + // --- webhooks --- + case "webhook.register": + return s.handleWebhookRegister(ctx, req.Args) + case "webhook.list": + return s.handleWebhookList(ctx, req.Args) + case "webhook.delete": + return s.handleWebhookDelete(ctx, req.Args) + + // --- k8s --- + case "k8s.register": + return s.handleK8sRegister(ctx, req.Args) + case "k8s.list": + return s.handleK8sList(ctx, req.Args) + case "k8s.delete": + return s.handleK8sDelete(ctx, req.Args) + + // --- attachments --- + case "attachments.gc": + return s.handleAttachmentsGC(ctx) + default: return Response{OK: false, Error: fmt.Sprintf("unknown command: %s", req.Command)} } @@ -1139,5 +1160,368 @@ func (s *AdminServer) handleRetentionStatus(ctx context.Context) Response { }} } +// ---------- webhook handlers ---------- + +func (s *AdminServer) handleWebhookRegister(ctx context.Context, args json.RawMessage) Response { + var p struct { + URL string `json:"url"` + Events string `json:"events"` + Secret string `json:"secret"` + AgentName string `json:"agent_name"` + } + if err := json.Unmarshal(args, &p); err != nil { + return Response{OK: false, Error: "invalid args: " + err.Error()} + } + if p.URL == "" || p.Events == "" || p.Secret == "" || p.AgentName == "" { + return Response{OK: false, Error: "url, events, secret, and agent_name are required"} + } + + if s.services.WebhookService == nil { + return Response{OK: false, Error: "webhook service not configured"} + } + + events := strings.Split(p.Events, ",") + for i := range events { + events[i] = strings.TrimSpace(events[i]) + } + + wh, err := s.services.WebhookService.RegisterWebhook(ctx, p.AgentName, p.URL, events, p.Secret) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + + return Response{OK: true, Data: map[string]interface{}{ + "id": wh.ID, + "url": wh.URL, + "events": wh.Events, + "status": wh.Status, + }} +} + +func (s *AdminServer) handleWebhookList(ctx context.Context, args json.RawMessage) Response { + var p struct { + AgentName string `json:"agent_name"` + } + if args != nil { + json.Unmarshal(args, &p) + } + + if s.services.WebhookService == nil { + return Response{OK: false, Error: "webhook service not configured"} + } + + if p.AgentName == "" { + // Admin: list all webhooks by querying DB directly. + rows, err := s.db.QueryContext(ctx, + `SELECT id, agent_name, url, events, status, consecutive_failures, created_at + FROM webhooks ORDER BY created_at DESC`) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + defer rows.Close() + + type whRow struct { + ID int64 `json:"id"` + AgentName string `json:"agent_name"` + URL string `json:"url"` + Events []string `json:"events"` + Status string `json:"status"` + ConsecutiveFailures int `json:"consecutive_failures"` + CreatedAt string `json:"created_at"` + } + + var result []whRow + for rows.Next() { + var r whRow + var eventsJSON string + var createdAt time.Time + if err := rows.Scan(&r.ID, &r.AgentName, &r.URL, &eventsJSON, &r.Status, &r.ConsecutiveFailures, &createdAt); err != nil { + return Response{OK: false, Error: "scan: " + err.Error()} + } + json.Unmarshal([]byte(eventsJSON), &r.Events) + if r.Events == nil { + r.Events = []string{} + } + r.CreatedAt = createdAt.Format(time.RFC3339) + result = append(result, r) + } + if result == nil { + result = []whRow{} + } + return Response{OK: true, Data: result} + } + + webhookList, err := s.services.WebhookService.ListWebhooks(ctx, p.AgentName) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + + type whRow struct { + ID int64 `json:"id"` + AgentName string `json:"agent_name"` + URL string `json:"url"` + Events []string `json:"events"` + Status string `json:"status"` + ConsecutiveFailures int `json:"consecutive_failures"` + CreatedAt string `json:"created_at"` + } + + result := make([]whRow, len(webhookList)) + for i, wh := range webhookList { + result[i] = whRow{ + ID: wh.ID, + AgentName: wh.AgentName, + URL: wh.URL, + Events: wh.Events, + Status: wh.Status, + ConsecutiveFailures: wh.ConsecutiveFailures, + CreatedAt: wh.CreatedAt.Format(time.RFC3339), + } + } + return Response{OK: true, Data: result} +} + +func (s *AdminServer) handleWebhookDelete(ctx context.Context, args json.RawMessage) Response { + var p struct { + ID int64 `json:"id"` + } + if err := json.Unmarshal(args, &p); err != nil { + return Response{OK: false, Error: "invalid args: " + err.Error()} + } + if p.ID <= 0 { + return Response{OK: false, Error: "id is required"} + } + + if s.services.WebhookService == nil { + return Response{OK: false, Error: "webhook service not configured"} + } + + // Admin bypass: look up the webhook's agent_name first, then delete. + var agentName string + err := s.db.QueryRowContext(ctx, "SELECT agent_name FROM webhooks WHERE id = ?", p.ID).Scan(&agentName) + if err != nil { + return Response{OK: false, Error: "webhook not found"} + } + + if err := s.services.WebhookService.DeleteWebhook(ctx, agentName, p.ID); err != nil { + return Response{OK: false, Error: err.Error()} + } + + return Response{OK: true, Data: map[string]interface{}{ + "deleted": p.ID, + }} +} + +// ---------- k8s handlers ---------- + +func (s *AdminServer) handleK8sRegister(ctx context.Context, args json.RawMessage) Response { + var p struct { + Image string `json:"image"` + Events string `json:"events"` + AgentName string `json:"agent_name"` + Namespace string `json:"namespace"` + ResourcesMemory string `json:"resources_memory"` + ResourcesCPU string `json:"resources_cpu"` + Env string `json:"env"` + TimeoutSeconds int `json:"timeout_seconds"` + } + if err := json.Unmarshal(args, &p); err != nil { + return Response{OK: false, Error: "invalid args: " + err.Error()} + } + if p.Image == "" || p.Events == "" || p.AgentName == "" { + return Response{OK: false, Error: "image, events, and agent_name are required"} + } + + if s.services.K8sService == nil { + return Response{OK: false, Error: "k8s service not configured"} + } + + events := strings.Split(p.Events, ",") + for i := range events { + events[i] = strings.TrimSpace(events[i]) + } + + // Parse env from comma-separated KEY=VALUE pairs. + envMap := map[string]string{} + if p.Env != "" { + for _, pair := range strings.Split(p.Env, ",") { + pair = strings.TrimSpace(pair) + parts := strings.SplitN(pair, "=", 2) + if len(parts) == 2 { + envMap[parts[0]] = parts[1] + } + } + } + + timeout := p.TimeoutSeconds + if timeout <= 0 { + timeout = 300 + } + + req := k8s.RegisterHandlerRequest{ + Image: p.Image, + Events: events, + Namespace: p.Namespace, + ResourcesMemory: p.ResourcesMemory, + ResourcesCPU: p.ResourcesCPU, + Env: envMap, + TimeoutSeconds: timeout, + } + + handler, err := s.services.K8sService.RegisterHandler(ctx, p.AgentName, req) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + + return Response{OK: true, Data: map[string]interface{}{ + "id": handler.ID, + "image": handler.Image, + "events": handler.Events, + "status": handler.Status, + }} +} + +func (s *AdminServer) handleK8sList(ctx context.Context, args json.RawMessage) Response { + var p struct { + AgentName string `json:"agent_name"` + } + if args != nil { + json.Unmarshal(args, &p) + } + + if s.services.K8sService == nil { + return Response{OK: false, Error: "k8s service not configured"} + } + + if p.AgentName == "" { + // Admin: list all K8s handlers by querying DB directly. + rows, err := s.db.QueryContext(ctx, + `SELECT id, agent_name, image, events, namespace, resources_memory, resources_cpu, timeout_seconds, status, created_at + FROM k8s_handlers ORDER BY created_at DESC`) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + defer rows.Close() + + type handlerRow struct { + ID int64 `json:"id"` + AgentName string `json:"agent_name"` + Image string `json:"image"` + Events []string `json:"events"` + Namespace string `json:"namespace"` + ResourcesMemory string `json:"resources_memory"` + ResourcesCPU string `json:"resources_cpu"` + TimeoutSeconds int `json:"timeout_seconds"` + Status string `json:"status"` + CreatedAt string `json:"created_at"` + } + + var result []handlerRow + for rows.Next() { + var r handlerRow + var eventsJSON string + var createdAt time.Time + if err := rows.Scan(&r.ID, &r.AgentName, &r.Image, &eventsJSON, &r.Namespace, + &r.ResourcesMemory, &r.ResourcesCPU, &r.TimeoutSeconds, &r.Status, &createdAt); err != nil { + return Response{OK: false, Error: "scan: " + err.Error()} + } + json.Unmarshal([]byte(eventsJSON), &r.Events) + if r.Events == nil { + r.Events = []string{} + } + r.CreatedAt = createdAt.Format(time.RFC3339) + result = append(result, r) + } + if result == nil { + result = []handlerRow{} + } + return Response{OK: true, Data: result} + } + + handlers, err := s.services.K8sService.ListHandlers(ctx, p.AgentName) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + + type handlerRow struct { + ID int64 `json:"id"` + AgentName string `json:"agent_name"` + Image string `json:"image"` + Events []string `json:"events"` + Namespace string `json:"namespace"` + ResourcesMemory string `json:"resources_memory"` + ResourcesCPU string `json:"resources_cpu"` + TimeoutSeconds int `json:"timeout_seconds"` + Status string `json:"status"` + CreatedAt string `json:"created_at"` + } + + result := make([]handlerRow, len(handlers)) + for i, h := range handlers { + result[i] = handlerRow{ + ID: h.ID, + AgentName: h.AgentName, + Image: h.Image, + Events: h.Events, + Namespace: h.Namespace, + ResourcesMemory: h.ResourcesMemory, + ResourcesCPU: h.ResourcesCPU, + TimeoutSeconds: h.TimeoutSeconds, + Status: h.Status, + CreatedAt: h.CreatedAt.Format(time.RFC3339), + } + } + return Response{OK: true, Data: result} +} + +func (s *AdminServer) handleK8sDelete(ctx context.Context, args json.RawMessage) Response { + var p struct { + ID int64 `json:"id"` + } + if err := json.Unmarshal(args, &p); err != nil { + return Response{OK: false, Error: "invalid args: " + err.Error()} + } + if p.ID <= 0 { + return Response{OK: false, Error: "id is required"} + } + + if s.services.K8sService == nil { + return Response{OK: false, Error: "k8s service not configured"} + } + + // Admin bypass: look up the handler's agent_name first, then delete. + var agentName string + err := s.db.QueryRowContext(ctx, "SELECT agent_name FROM k8s_handlers WHERE id = ?", p.ID).Scan(&agentName) + if err != nil { + return Response{OK: false, Error: "handler not found"} + } + + if err := s.services.K8sService.DeleteHandler(ctx, agentName, p.ID); err != nil { + return Response{OK: false, Error: err.Error()} + } + + return Response{OK: true, Data: map[string]interface{}{ + "deleted": p.ID, + }} +} + +// ---------- attachments handlers ---------- + +func (s *AdminServer) handleAttachmentsGC(ctx context.Context) Response { + if s.services.AttachmentService == nil { + return Response{OK: false, Error: "attachment service not configured"} + } + + result, err := s.services.AttachmentService.GarbageCollect(ctx) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + + return Response{OK: true, Data: map[string]interface{}{ + "files_removed": result.FilesRemoved, + "bytes_reclaimed": result.BytesReclaimed, + }} +} + // Ensure the messaging import is used. var _ = messaging.StatusPending