Add TTFT & ITL Metrics for LLM Inference #28
@@ -11,6 +11,7 @@ import (
|
||||
|
||||
"forgejo.riotpiao.com/rock/homelab-frontend/internal/auth"
|
||||
"forgejo.riotpiao.com/rock/homelab-frontend/internal/config"
|
||||
"forgejo.riotpiao.com/rock/homelab-frontend/internal/notification"
|
||||
"forgejo.riotpiao.com/rock/homelab-frontend/internal/proxy"
|
||||
"forgejo.riotpiao.com/rock/homelab-frontend/internal/server"
|
||||
"forgejo.riotpiao.com/rock/homelab-frontend/internal/serviceadapter"
|
||||
@@ -82,6 +83,29 @@ func main() {
|
||||
}
|
||||
_ = registry.Add(workflowAdapter)
|
||||
|
||||
// Add notification service adapter (internal handler, no upstream proxy)
|
||||
notifHandler := notification.NewHandler()
|
||||
notifAdapter := &serviceadapter.ServiceAdapter{
|
||||
Namespace: "notification",
|
||||
ServiceName: "notification",
|
||||
Handler: notifHandler,
|
||||
Spec: serviceadapter.Spec{
|
||||
ServiceName: "notification",
|
||||
Auth: serviceadapter.Auth{Required: true},
|
||||
Resources: []serviceadapter.Resource{
|
||||
{Name: "send-email", Methods: []serviceadapter.Method{{Verb: "POST", UpstreamPath: "/send-email"}}},
|
||||
{Name: "send-message", Methods: []serviceadapter.Method{{Verb: "POST", UpstreamPath: "/send-message"}}},
|
||||
{Name: "list-messages", Methods: []serviceadapter.Method{{Verb: "GET", UpstreamPath: "/list-messages"}}},
|
||||
{Name: "delete-message", Methods: []serviceadapter.Method{{Verb: "DELETE", UpstreamPath: "/delete-message"}}},
|
||||
{Name: "delete-all-messages", Methods: []serviceadapter.Method{{Verb: "DELETE", UpstreamPath: "/delete-all-messages"}}},
|
||||
{Name: "list-applications", Methods: []serviceadapter.Method{{Verb: "GET", UpstreamPath: "/list-applications"}}},
|
||||
{Name: "create-application", Methods: []serviceadapter.Method{{Verb: "POST", UpstreamPath: "/create-application"}}},
|
||||
{Name: "delete-application", Methods: []serviceadapter.Method{{Verb: "DELETE", UpstreamPath: "/delete-application"}}},
|
||||
},
|
||||
},
|
||||
}
|
||||
_ = registry.Add(notifAdapter)
|
||||
|
||||
// Add other adapters from config
|
||||
for _, a := range cfg.Adapters {
|
||||
_ = registry.Add(a)
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
#!/bin/bash
|
||||
# Example: Gotify CRUD operations via notification service (X-Service routing)
|
||||
|
||||
BASE_URL="${1:-https://api.riotpiao.com}"
|
||||
AUTH_TOKEN="${2:-}"
|
||||
AUTH="-H \"Authorization: Bearer $AUTH_TOKEN\""
|
||||
|
||||
echo "=== Send Gotify Message ==="
|
||||
curl -s -X POST "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: send-message" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" \
|
||||
-d '{
|
||||
"title": "Deployment Complete",
|
||||
"message": "homelab-frontend v1.2.0 deployed to production",
|
||||
"priority": 5
|
||||
}' | jq .
|
||||
|
||||
echo ""
|
||||
echo "=== List Messages ==="
|
||||
curl -s -X GET "$BASE_URL?limit=10" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: list-messages" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" | jq .
|
||||
|
||||
echo ""
|
||||
echo "=== List Applications ==="
|
||||
curl -s -X GET "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: list-applications" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" | jq .
|
||||
|
||||
echo ""
|
||||
echo "=== Create Application ==="
|
||||
curl -s -X POST "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: create-application" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" \
|
||||
-d '{
|
||||
"name": "my-monitor",
|
||||
"description": "Monitoring alerts"
|
||||
}' | jq .
|
||||
|
||||
echo ""
|
||||
echo "=== Delete Message (by ID) ==="
|
||||
curl -s -X DELETE "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: delete-message" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" \
|
||||
-d '{"id": 1}' | jq .
|
||||
|
||||
echo ""
|
||||
echo "=== Delete Application (by ID) ==="
|
||||
curl -s -X DELETE "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: delete-application" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" \
|
||||
-d '{"id": 1}' | jq .
|
||||
+14
-30
@@ -1,36 +1,20 @@
|
||||
#!/bin/bash
|
||||
# Example: Send email via notification/sendMsg endpoint
|
||||
# Example: Send email via notification service (X-Service routing)
|
||||
|
||||
BASE_URL="${1:-https://api.riotpiao.com}"
|
||||
AUTH_TOKEN="${2:-}" # Optional JWT token if auth required
|
||||
AUTH_TOKEN="${2:-}"
|
||||
|
||||
PAYLOAD=$(cat <<'EOF'
|
||||
{
|
||||
"format": "smtp",
|
||||
"title": "System Alert",
|
||||
"message": "CPU usage exceeded 90% threshold",
|
||||
"priority": 7,
|
||||
"extras": {
|
||||
"to_email": "[email protected]",
|
||||
"cc": "[email protected]"
|
||||
}
|
||||
}
|
||||
EOF
|
||||
)
|
||||
|
||||
if [ -n "$AUTH_TOKEN" ]; then
|
||||
curl -X POST "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: sendMsg" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" \
|
||||
-d "$PAYLOAD"
|
||||
else
|
||||
curl -X POST "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: sendMsg" \
|
||||
-H "Content-Type: application/json" \
|
||||
-d "$PAYLOAD"
|
||||
fi
|
||||
# Send email
|
||||
curl -s -X POST "$BASE_URL" \
|
||||
-H "X-Service: notification" \
|
||||
-H "X-Resource: send-email" \
|
||||
-H "Content-Type: application/json" \
|
||||
-H "Authorization: Bearer $AUTH_TOKEN" \
|
||||
-d '{
|
||||
"to": "[email protected]",
|
||||
"cc": "[email protected]",
|
||||
"subject": "System Alert",
|
||||
"body": "CPU usage exceeded 90% threshold"
|
||||
}' | jq .
|
||||
|
||||
echo ""
|
||||
|
||||
@@ -0,0 +1,257 @@
|
||||
package notification
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"time"
|
||||
)
|
||||
|
||||
// GotifyClient is a CRUD client for the Gotify API.
|
||||
type GotifyClient struct {
|
||||
baseURL string
|
||||
appToken string // token for sending messages (application token)
|
||||
clientToken string // token for reading/managing (client token)
|
||||
httpClient *http.Client
|
||||
}
|
||||
|
||||
// NewGotifyClient creates a Gotify API client.
|
||||
// appToken is used for sending messages.
|
||||
// clientToken is used for listing/deleting messages and managing applications.
|
||||
func NewGotifyClient(baseURL, appToken, clientToken string) *GotifyClient {
|
||||
return &GotifyClient{
|
||||
baseURL: baseURL,
|
||||
appToken: appToken,
|
||||
clientToken: clientToken,
|
||||
httpClient: &http.Client{
|
||||
Timeout: 10 * time.Second,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// --- Message Types ---
|
||||
|
||||
// GotifyMessage represents a Gotify message.
|
||||
type GotifyMessage struct {
|
||||
ID int `json:"id,omitempty"`
|
||||
AppID int `json:"appid,omitempty"`
|
||||
Title string `json:"title"`
|
||||
Message string `json:"message"`
|
||||
Priority int `json:"priority,omitempty"`
|
||||
Date string `json:"date,omitempty"`
|
||||
Extras map[string]interface{} `json:"extras,omitempty"`
|
||||
}
|
||||
|
||||
// GotifyMessageList is a paginated list of messages.
|
||||
type GotifyMessageList struct {
|
||||
Messages []GotifyMessage `json:"messages"`
|
||||
Paging GotifyPaging `json:"paging"`
|
||||
}
|
||||
|
||||
// GotifyPaging represents pagination info.
|
||||
type GotifyPaging struct {
|
||||
Size int `json:"size"`
|
||||
Since int `json:"since"`
|
||||
Limit int `json:"limit"`
|
||||
Next string `json:"next,omitempty"`
|
||||
}
|
||||
|
||||
// --- Application Types ---
|
||||
|
||||
// GotifyApplication represents a Gotify application.
|
||||
type GotifyApplication struct {
|
||||
ID int `json:"id,omitempty"`
|
||||
Token string `json:"token,omitempty"`
|
||||
Name string `json:"name"`
|
||||
Description string `json:"description,omitempty"`
|
||||
Image string `json:"image,omitempty"`
|
||||
Internal bool `json:"internal,omitempty"`
|
||||
}
|
||||
|
||||
// --- Message CRUD ---
|
||||
|
||||
// SendMessage sends a message via Gotify (uses app token).
|
||||
func (c *GotifyClient) SendMessage(msg GotifyMessage) (*GotifyMessage, error) {
|
||||
body, err := json.Marshal(msg)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("marshal message: %w", err)
|
||||
}
|
||||
|
||||
req, err := http.NewRequest(http.MethodPost, c.baseURL+"/message", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("X-Gotify-Key", c.appToken)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("send message: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
|
||||
return nil, c.readError(resp)
|
||||
}
|
||||
|
||||
var result GotifyMessage
|
||||
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
|
||||
return nil, fmt.Errorf("decode response: %w", err)
|
||||
}
|
||||
return &result, nil
|
||||
}
|
||||
|
||||
// ListMessages lists messages (uses client token).
|
||||
func (c *GotifyClient) ListMessages(limit int) (*GotifyMessageList, error) {
|
||||
url := fmt.Sprintf("%s/message?limit=%d", c.baseURL, limit)
|
||||
req, err := http.NewRequest(http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("X-Gotify-Key", c.clientToken)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list messages: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, c.readError(resp)
|
||||
}
|
||||
|
||||
var result GotifyMessageList
|
||||
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
|
||||
return nil, fmt.Errorf("decode response: %w", err)
|
||||
}
|
||||
return &result, nil
|
||||
}
|
||||
|
||||
// DeleteMessage deletes a message by ID (uses client token).
|
||||
func (c *GotifyClient) DeleteMessage(id int) error {
|
||||
url := fmt.Sprintf("%s/message/%d", c.baseURL, id)
|
||||
req, err := http.NewRequest(http.MethodDelete, url, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("X-Gotify-Key", c.clientToken)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("delete message: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusNoContent {
|
||||
return c.readError(resp)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteAllMessages deletes all messages (uses client token).
|
||||
func (c *GotifyClient) DeleteAllMessages() error {
|
||||
req, err := http.NewRequest(http.MethodDelete, c.baseURL+"/message", nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("X-Gotify-Key", c.clientToken)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("delete all messages: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusNoContent {
|
||||
return c.readError(resp)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// --- Application CRUD ---
|
||||
|
||||
// ListApplications lists all applications (uses client token).
|
||||
func (c *GotifyClient) ListApplications() ([]GotifyApplication, error) {
|
||||
req, err := http.NewRequest(http.MethodGet, c.baseURL+"/application", nil)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("X-Gotify-Key", c.clientToken)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list applications: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, c.readError(resp)
|
||||
}
|
||||
|
||||
var result []GotifyApplication
|
||||
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
|
||||
return nil, fmt.Errorf("decode response: %w", err)
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// CreateApplication creates a new application (uses client token).
|
||||
func (c *GotifyClient) CreateApplication(app GotifyApplication) (*GotifyApplication, error) {
|
||||
body, err := json.Marshal(app)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("marshal application: %w", err)
|
||||
}
|
||||
|
||||
req, err := http.NewRequest(http.MethodPost, c.baseURL+"/application", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("X-Gotify-Key", c.clientToken)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create application: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
|
||||
return nil, c.readError(resp)
|
||||
}
|
||||
|
||||
var result GotifyApplication
|
||||
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
|
||||
return nil, fmt.Errorf("decode response: %w", err)
|
||||
}
|
||||
return &result, nil
|
||||
}
|
||||
|
||||
// DeleteApplication deletes an application by ID (uses client token).
|
||||
func (c *GotifyClient) DeleteApplication(id int) error {
|
||||
url := fmt.Sprintf("%s/application/%d", c.baseURL, id)
|
||||
req, err := http.NewRequest(http.MethodDelete, url, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create request: %w", err)
|
||||
}
|
||||
req.Header.Set("X-Gotify-Key", c.clientToken)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("delete application: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusNoContent {
|
||||
return c.readError(resp)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// --- Helpers ---
|
||||
|
||||
func (c *GotifyClient) readError(resp *http.Response) error {
|
||||
body, _ := io.ReadAll(resp.Body)
|
||||
return fmt.Errorf("gotify API error (HTTP %d): %s", resp.StatusCode, string(body))
|
||||
}
|
||||
+230
-171
@@ -7,226 +7,285 @@ import (
|
||||
"net/http"
|
||||
"net/smtp"
|
||||
"os"
|
||||
"strings"
|
||||
"strconv"
|
||||
)
|
||||
|
||||
// SendMsgRequest represents a sendMsg API request.
|
||||
type SendMsgRequest struct {
|
||||
Format string `json:"format"` // "smtp" or "sms"
|
||||
Title string `json:"title"`
|
||||
Message string `json:"message"`
|
||||
Priority int `json:"priority,omitempty"`
|
||||
Extras map[string]string `json:"extras,omitempty"` // e.g., {"to_email": "[email protected]", "phone": "+1234567890"}
|
||||
}
|
||||
|
||||
// SendMsgResponse represents a sendMsg API response.
|
||||
type SendMsgResponse struct {
|
||||
Status string `json:"status"`
|
||||
MessageID string `json:"messageId,omitempty"`
|
||||
Error string `json:"error,omitempty"`
|
||||
}
|
||||
|
||||
// Handler handles sendMsg requests and forwards to appropriate channel (email, SMS, or Gotify push).
|
||||
// Handler handles notification requests routed via X-Resource header.
|
||||
// Supports: send-email, send-gotify, list-messages, delete-message,
|
||||
// delete-all-messages, list-applications, create-application, delete-application.
|
||||
type Handler struct {
|
||||
smtpHost string
|
||||
smtpPort string
|
||||
smtpFrom string
|
||||
smtpUser string
|
||||
smtpPass string
|
||||
smsAPIURL string
|
||||
smsAPIKey string
|
||||
gotifyURL string
|
||||
gotifyToken string
|
||||
smtpHost string
|
||||
smtpPort string
|
||||
smtpFrom string
|
||||
smtpUser string
|
||||
smtpPass string
|
||||
gotify *GotifyClient
|
||||
}
|
||||
|
||||
// NewHandler creates a new notification handler from environment variables.
|
||||
// NewHandler creates a notification handler from environment variables.
|
||||
func NewHandler() *Handler {
|
||||
var gotify *GotifyClient
|
||||
gotifyURL := os.Getenv("GOTIFY_URL")
|
||||
if gotifyURL != "" {
|
||||
gotify = NewGotifyClient(
|
||||
gotifyURL,
|
||||
os.Getenv("GOTIFY_APP_TOKEN"),
|
||||
os.Getenv("GOTIFY_CLIENT_TOKEN"),
|
||||
)
|
||||
}
|
||||
|
||||
return &Handler{
|
||||
smtpHost: os.Getenv("SMTP_HOST"),
|
||||
smtpPort: os.Getenv("SMTP_PORT"),
|
||||
smtpFrom: os.Getenv("SMTP_FROM"),
|
||||
smtpUser: os.Getenv("SMTP_USER"),
|
||||
smtpPass: os.Getenv("SMTP_PASS"),
|
||||
smsAPIURL: os.Getenv("SMS_API_URL"),
|
||||
smsAPIKey: os.Getenv("SMS_API_KEY"),
|
||||
gotifyURL: os.Getenv("GOTIFY_URL"),
|
||||
gotifyToken: os.Getenv("GOTIFY_TOKEN"),
|
||||
smtpHost: os.Getenv("SMTP_HOST"),
|
||||
smtpPort: os.Getenv("SMTP_PORT"),
|
||||
smtpFrom: os.Getenv("SMTP_FROM"),
|
||||
smtpUser: os.Getenv("SMTP_USER"),
|
||||
smtpPass: os.Getenv("SMTP_PASS"),
|
||||
gotify: gotify,
|
||||
}
|
||||
}
|
||||
|
||||
// ServeHTTP handles sendMsg requests.
|
||||
// ServeHTTP routes requests by X-Upstream-Path (set by dispatcher after resource matching).
|
||||
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
resource := r.Header.Get("X-Resource")
|
||||
|
||||
var req SendMsgRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
json.NewEncoder(w).Encode(SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "invalid request: " + err.Error(),
|
||||
})
|
||||
return
|
||||
}
|
||||
switch resource {
|
||||
// --- Email ---
|
||||
case "send-email":
|
||||
h.handleSendEmail(w, r)
|
||||
|
||||
// --- Gotify Messages ---
|
||||
case "send-message":
|
||||
h.handleSendGotify(w, r)
|
||||
case "list-messages":
|
||||
h.handleListMessages(w, r)
|
||||
case "delete-message":
|
||||
h.handleDeleteMessage(w, r)
|
||||
case "delete-all-messages":
|
||||
h.handleDeleteAllMessages(w, r)
|
||||
|
||||
// --- Gotify Applications ---
|
||||
case "list-applications":
|
||||
h.handleListApplications(w, r)
|
||||
case "create-application":
|
||||
h.handleCreateApplication(w, r)
|
||||
case "delete-application":
|
||||
h.handleDeleteApplication(w, r)
|
||||
|
||||
// Route based on format
|
||||
var resp SendMsgResponse
|
||||
switch req.Format {
|
||||
case "smtp":
|
||||
resp = h.sendEmail(req)
|
||||
case "sms":
|
||||
resp = h.sendSMS(req)
|
||||
case "gotify":
|
||||
resp = h.sendGotify(req)
|
||||
default:
|
||||
resp = SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "unsupported format: " + req.Format,
|
||||
}
|
||||
h.writeJSON(w, http.StatusNotFound, map[string]string{
|
||||
"error": fmt.Sprintf("unknown resource: %s", resource),
|
||||
})
|
||||
}
|
||||
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if resp.Error != "" {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
} else {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}
|
||||
json.NewEncoder(w).Encode(resp)
|
||||
}
|
||||
|
||||
// sendEmail sends an email via SMTP.
|
||||
func (h *Handler) sendEmail(req SendMsgRequest) SendMsgResponse {
|
||||
toEmail := req.Extras["to_email"]
|
||||
if toEmail == "" {
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "missing to_email in extras",
|
||||
}
|
||||
// --- Email ---
|
||||
|
||||
type SendEmailRequest struct {
|
||||
To string `json:"to"`
|
||||
CC string `json:"cc,omitempty"`
|
||||
Subject string `json:"subject"`
|
||||
Body string `json:"body"`
|
||||
}
|
||||
|
||||
func (h *Handler) handleSendEmail(w http.ResponseWriter, r *http.Request) {
|
||||
var req SendEmailRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request: " + err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
subject := req.Title
|
||||
if req.To == "" {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "missing 'to' field"})
|
||||
return
|
||||
}
|
||||
|
||||
subject := req.Subject
|
||||
if subject == "" {
|
||||
subject = "Notification"
|
||||
}
|
||||
|
||||
// Construct email body
|
||||
body := req.Message
|
||||
if req.Extras != nil {
|
||||
if cc := req.Extras["cc"]; cc != "" {
|
||||
body = fmt.Sprintf("CC: %s\n\n%s", cc, body)
|
||||
}
|
||||
}
|
||||
|
||||
msg := fmt.Sprintf(
|
||||
"From: %s\r\nTo: %s\r\nSubject: %s\r\nContent-Type: text/plain; charset=UTF-8\r\n\r\n%s",
|
||||
h.smtpFrom, toEmail, subject, body,
|
||||
h.smtpFrom, req.To, subject, req.Body,
|
||||
)
|
||||
|
||||
// Send via SMTP
|
||||
smtpAddr := fmt.Sprintf("%s:%s", h.smtpHost, h.smtpPort)
|
||||
auth := smtp.PlainAuth("", h.smtpUser, h.smtpPass, h.smtpHost)
|
||||
|
||||
if err := smtp.SendMail(smtpAddr, auth, h.smtpFrom, []string{toEmail}, []byte(msg)); err != nil {
|
||||
log.Printf("error sending email to %s: %v", toEmail, err)
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "failed to send email: " + err.Error(),
|
||||
}
|
||||
if err := smtp.SendMail(smtpAddr, auth, h.smtpFrom, []string{req.To}, []byte(msg)); err != nil {
|
||||
log.Printf("error sending email to %s: %v", req.To, err)
|
||||
h.writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "failed to send email: " + err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
return SendMsgResponse{
|
||||
Status: "success",
|
||||
MessageID: fmt.Sprintf("email-%s", toEmail),
|
||||
}
|
||||
h.writeJSON(w, http.StatusOK, map[string]string{
|
||||
"status": "success",
|
||||
"messageId": fmt.Sprintf("email-%s", req.To),
|
||||
})
|
||||
}
|
||||
|
||||
// sendSMS sends an SMS via configured provider.
|
||||
// Placeholder: integrate with Twilio, AWS SNS, or similar.
|
||||
func (h *Handler) sendSMS(req SendMsgRequest) SendMsgResponse {
|
||||
phone := req.Extras["phone"]
|
||||
if phone == "" {
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "missing phone in extras",
|
||||
}
|
||||
// --- Gotify Messages ---
|
||||
|
||||
func (h *Handler) handleSendGotify(w http.ResponseWriter, r *http.Request) {
|
||||
if h.gotify == nil {
|
||||
h.writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "Gotify not configured"})
|
||||
return
|
||||
}
|
||||
|
||||
// TODO: Implement SMS provider integration (Twilio, AWS SNS, etc.)
|
||||
// For now, return error
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "SMS not implemented yet",
|
||||
}
|
||||
}
|
||||
|
||||
// sendGotify sends a notification via Gotify server.
|
||||
func (h *Handler) sendGotify(req SendMsgRequest) SendMsgResponse {
|
||||
if h.gotifyURL == "" || h.gotifyToken == "" {
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "Gotify not configured (missing GOTIFY_URL or GOTIFY_TOKEN)",
|
||||
}
|
||||
var msg GotifyMessage
|
||||
if err := json.NewDecoder(r.Body).Decode(&msg); err != nil {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request: " + err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Construct Gotify message
|
||||
gotifyReq := map[string]interface{}{
|
||||
"title": req.Title,
|
||||
"message": req.Message,
|
||||
"priority": req.Priority,
|
||||
}
|
||||
|
||||
// Marshal to JSON
|
||||
body, err := json.Marshal(gotifyReq)
|
||||
result, err := h.gotify.SendMessage(msg)
|
||||
if err != nil {
|
||||
log.Printf("error marshaling Gotify request: %v", err)
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "failed to marshal request: " + err.Error(),
|
||||
}
|
||||
log.Printf("error sending gotify message: %v", err)
|
||||
h.writeJSON(w, http.StatusBadGateway, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// POST to Gotify
|
||||
gotifyEndpoint := fmt.Sprintf("%s/message?token=%s", h.gotifyURL, h.gotifyToken)
|
||||
resp, err := http.Post(gotifyEndpoint, "application/json", strings.NewReader(string(body)))
|
||||
h.writeJSON(w, http.StatusOK, result)
|
||||
}
|
||||
|
||||
func (h *Handler) handleListMessages(w http.ResponseWriter, r *http.Request) {
|
||||
if h.gotify == nil {
|
||||
h.writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "Gotify not configured"})
|
||||
return
|
||||
}
|
||||
|
||||
limit := 50
|
||||
if l := r.URL.Query().Get("limit"); l != "" {
|
||||
if parsed, err := strconv.Atoi(l); err == nil && parsed > 0 {
|
||||
limit = parsed
|
||||
}
|
||||
}
|
||||
|
||||
result, err := h.gotify.ListMessages(limit)
|
||||
if err != nil {
|
||||
log.Printf("error listing gotify messages: %v", err)
|
||||
h.writeJSON(w, http.StatusBadGateway, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
h.writeJSON(w, http.StatusOK, result)
|
||||
}
|
||||
|
||||
func (h *Handler) handleDeleteMessage(w http.ResponseWriter, r *http.Request) {
|
||||
if h.gotify == nil {
|
||||
h.writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "Gotify not configured"})
|
||||
return
|
||||
}
|
||||
|
||||
var req struct {
|
||||
ID int `json:"id"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request: " + err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
if req.ID == 0 {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "missing 'id' field"})
|
||||
return
|
||||
}
|
||||
|
||||
if err := h.gotify.DeleteMessage(req.ID); err != nil {
|
||||
log.Printf("error deleting gotify message %d: %v", req.ID, err)
|
||||
h.writeJSON(w, http.StatusBadGateway, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
h.writeJSON(w, http.StatusOK, map[string]string{"status": "deleted"})
|
||||
}
|
||||
|
||||
func (h *Handler) handleDeleteAllMessages(w http.ResponseWriter, r *http.Request) {
|
||||
if h.gotify == nil {
|
||||
h.writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "Gotify not configured"})
|
||||
return
|
||||
}
|
||||
|
||||
if err := h.gotify.DeleteAllMessages(); err != nil {
|
||||
log.Printf("error deleting all gotify messages: %v", err)
|
||||
h.writeJSON(w, http.StatusBadGateway, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
h.writeJSON(w, http.StatusOK, map[string]string{"status": "all messages deleted"})
|
||||
}
|
||||
|
||||
// --- Gotify Applications ---
|
||||
|
||||
func (h *Handler) handleListApplications(w http.ResponseWriter, r *http.Request) {
|
||||
if h.gotify == nil {
|
||||
h.writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "Gotify not configured"})
|
||||
return
|
||||
}
|
||||
|
||||
result, err := h.gotify.ListApplications()
|
||||
if err != nil {
|
||||
log.Printf("error listing gotify applications: %v", err)
|
||||
h.writeJSON(w, http.StatusBadGateway, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
h.writeJSON(w, http.StatusOK, result)
|
||||
}
|
||||
|
||||
func (h *Handler) handleCreateApplication(w http.ResponseWriter, r *http.Request) {
|
||||
if h.gotify == nil {
|
||||
h.writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "Gotify not configured"})
|
||||
return
|
||||
}
|
||||
|
||||
var app GotifyApplication
|
||||
if err := json.NewDecoder(r.Body).Decode(&app); err != nil {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request: " + err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
result, err := h.gotify.CreateApplication(app)
|
||||
if err != nil {
|
||||
log.Printf("error sending to Gotify: %v", err)
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: "failed to send to Gotify: " + err.Error(),
|
||||
}
|
||||
log.Printf("error creating gotify application: %v", err)
|
||||
h.writeJSON(w, http.StatusBadGateway, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
h.writeJSON(w, http.StatusCreated, result)
|
||||
}
|
||||
|
||||
// Check response status
|
||||
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
|
||||
log.Printf("Gotify returned status %d", resp.StatusCode)
|
||||
return SendMsgResponse{
|
||||
Status: "error",
|
||||
Error: fmt.Sprintf("Gotify returned status %d", resp.StatusCode),
|
||||
}
|
||||
func (h *Handler) handleDeleteApplication(w http.ResponseWriter, r *http.Request) {
|
||||
if h.gotify == nil {
|
||||
h.writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "Gotify not configured"})
|
||||
return
|
||||
}
|
||||
|
||||
// Parse response
|
||||
var gotifyResp map[string]interface{}
|
||||
if err := json.NewDecoder(resp.Body).Decode(&gotifyResp); err != nil {
|
||||
log.Printf("error decoding Gotify response: %v", err)
|
||||
return SendMsgResponse{
|
||||
Status: "success",
|
||||
MessageID: "gotify-sent",
|
||||
}
|
||||
var req struct {
|
||||
ID int `json:"id"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid request: " + err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Extract message ID if available
|
||||
msgID := "gotify-sent"
|
||||
if id, ok := gotifyResp["id"]; ok {
|
||||
msgID = fmt.Sprintf("gotify-%v", id)
|
||||
if req.ID == 0 {
|
||||
h.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "missing 'id' field"})
|
||||
return
|
||||
}
|
||||
|
||||
return SendMsgResponse{
|
||||
Status: "success",
|
||||
MessageID: msgID,
|
||||
if err := h.gotify.DeleteApplication(req.ID); err != nil {
|
||||
log.Printf("error deleting gotify application %d: %v", req.ID, err)
|
||||
h.writeJSON(w, http.StatusBadGateway, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
h.writeJSON(w, http.StatusOK, map[string]string{"status": "deleted"})
|
||||
}
|
||||
|
||||
// --- Helpers ---
|
||||
|
||||
func (h *Handler) writeJSON(w http.ResponseWriter, status int, data interface{}) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(status)
|
||||
json.NewEncoder(w).Encode(data)
|
||||
}
|
||||
|
||||
@@ -80,6 +80,14 @@ func (d *Dispatcher) Dispatch(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
}
|
||||
|
||||
// Internal handler: dispatch directly without reverse proxy
|
||||
if adapter.Handler != nil {
|
||||
// Set X-Upstream-Path so the handler knows which method was matched
|
||||
r.Header.Set("X-Upstream-Path", method.UpstreamPath)
|
||||
adapter.Handler.ServeHTTP(w, r)
|
||||
return
|
||||
}
|
||||
|
||||
upstreamURL := adapter.Spec.Upstream.URL
|
||||
if strings.HasPrefix(upstreamURL, "grpc://") {
|
||||
d.dispatchGRPC(w, r, upstreamURL, method, adapter)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package serviceadapter
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -49,11 +50,14 @@ type Status struct {
|
||||
}
|
||||
|
||||
// ServiceAdapter is a gateway service adapter.
|
||||
// When Handler is set, the dispatcher routes directly to the internal handler
|
||||
// instead of reverse-proxying to Spec.Upstream.URL.
|
||||
type ServiceAdapter struct {
|
||||
Name string // namespace/name
|
||||
Namespace string
|
||||
Name string // namespace/name
|
||||
Namespace string
|
||||
ServiceName string
|
||||
Spec Spec
|
||||
Status Status
|
||||
CreatedAt time.Time
|
||||
Spec Spec
|
||||
Status Status
|
||||
CreatedAt time.Time
|
||||
Handler http.Handler `json:"-" yaml:"-"` // internal handler (skip serialization)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user