mirror of
https://github.com/langgenius/dify-plugin-daemon.git
synced 2026-07-23 02:05:27 -04:00
ca3d00229e
* use slog instead of log package and format to new log schema * update the environment name to LOG_OUTPUT_FORMAT * add the env to .env.example * fix log reference error * change the order of milldlewares * delete unused code * fix the concurrently session potential race condition * fix the log format in tests * update the duplicate code * refactor: convert log functions to slog structured format - Change log.Error/Info/Warn/Debug/Panic to accept msg + key-value pairs - Remove printf-style formatting from log functions - Update log calls in internal/cluster, internal/db, internal/core/session_manager - Remove unused 'initialized' variable from log package - Remaining files will be updated in follow-up commits 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com> * refactor: update all log call sites to use slog structured format Convert all log.Error, log.Info, log.Warn, log.Debug, and log.Panic calls from printf-style formatting to slog key-value pairs. Before: log.Error("failed to do something: %s", err.Error()) After: log.Error("failed to do something", "error", err) 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com> * refactor: update cmd/ log calls to use slog structured format 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com> * feat: implement GnetLogger for structured logging in gnet * refactor: remove deprecated log visibility functions and related calls * feat: enhance session management with trace and identity context propagation * feat: implement serverless transaction handler and writer for plugin runtime * refactor: rename context field to traceCtx in RealBackwardsInvocation --------- Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com> Co-authored-by: Yeuoly <admin@srmxy.cn>
181 lines
4.9 KiB
Go
181 lines
4.9 KiB
Go
package http_requests
|
|
|
|
import (
|
|
"bytes"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"time"
|
|
|
|
routinepkg "github.com/langgenius/dify-plugin-daemon/pkg/routine"
|
|
"github.com/langgenius/dify-plugin-daemon/pkg/utils/log"
|
|
"github.com/langgenius/dify-plugin-daemon/pkg/utils/parser"
|
|
"github.com/langgenius/dify-plugin-daemon/pkg/utils/routine"
|
|
"github.com/langgenius/dify-plugin-daemon/pkg/utils/stream"
|
|
)
|
|
|
|
func parseJsonBody(resp *http.Response, ret interface{}) error {
|
|
defer resp.Body.Close()
|
|
jsonDecoder := json.NewDecoder(resp.Body)
|
|
return jsonDecoder.Decode(ret)
|
|
}
|
|
|
|
func RequestAndParse[T any](client *http.Client, url string, method string, options ...HttpOptions) (*T, error) {
|
|
var ret T
|
|
|
|
// check if ret is a map, if so, create a new map
|
|
if _, ok := any(ret).(map[string]any); ok {
|
|
ret = *new(T)
|
|
}
|
|
|
|
resp, err := Request(client, url, method, options...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// get read timeout
|
|
readTimeout := int64(60000)
|
|
for _, option := range options {
|
|
if option.Type == HttpOptionTypeReadTimeout {
|
|
readTimeout = option.Value.(int64)
|
|
break
|
|
}
|
|
}
|
|
time.AfterFunc(time.Millisecond*time.Duration(readTimeout), func() {
|
|
// close the response body if timeout
|
|
resp.Body.Close()
|
|
})
|
|
|
|
err = parseJsonBody(resp, &ret)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &ret, nil
|
|
}
|
|
|
|
func GetAndParse[T any](client *http.Client, url string, options ...HttpOptions) (*T, error) {
|
|
return RequestAndParse[T](client, url, "GET", options...)
|
|
}
|
|
|
|
func PostAndParse[T any](client *http.Client, url string, options ...HttpOptions) (*T, error) {
|
|
return RequestAndParse[T](client, url, "POST", options...)
|
|
}
|
|
|
|
func PutAndParse[T any](client *http.Client, url string, options ...HttpOptions) (*T, error) {
|
|
return RequestAndParse[T](client, url, "PUT", options...)
|
|
}
|
|
|
|
func DeleteAndParse[T any](client *http.Client, url string, options ...HttpOptions) (*T, error) {
|
|
return RequestAndParse[T](client, url, "DELETE", options...)
|
|
}
|
|
|
|
func PatchAndParse[T any](client *http.Client, url string, options ...HttpOptions) (*T, error) {
|
|
return RequestAndParse[T](client, url, "PATCH", options...)
|
|
}
|
|
|
|
func RequestAndParseStream[T any](client *http.Client, url string, method string, options ...HttpOptions) (*stream.Stream[T], error) {
|
|
resp, err := Request(client, url, method, options...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if resp.StatusCode != http.StatusOK {
|
|
defer resp.Body.Close()
|
|
errorText, _ := io.ReadAll(resp.Body)
|
|
return nil, fmt.Errorf("request failed with status code: %d and respond with: %s", resp.StatusCode, errorText)
|
|
}
|
|
|
|
ch := stream.NewStream[T](1024)
|
|
|
|
// get read timeout
|
|
readTimeout := int64(60000)
|
|
raiseErrorWhenStreamDataNotMatch := false
|
|
usingLengthPrefixed := false
|
|
for _, option := range options {
|
|
if option.Type == HttpOptionTypeReadTimeout {
|
|
readTimeout = option.Value.(int64)
|
|
} else if option.Type == HttpOptionTypeRaiseErrorWhenStreamDataNotMatch {
|
|
raiseErrorWhenStreamDataNotMatch = option.Value.(bool)
|
|
} else if option.Type == HttpOptionTypeUsingLengthPrefixed {
|
|
usingLengthPrefixed = option.Value.(bool)
|
|
}
|
|
}
|
|
time.AfterFunc(time.Millisecond*time.Duration(readTimeout), func() {
|
|
// close the response body if timeout
|
|
resp.Body.Close()
|
|
})
|
|
|
|
// Common data processor function to reduce code duplication
|
|
processData := func(data []byte) error {
|
|
// unmarshal
|
|
t, err := parser.UnmarshalJsonBytes[T](data)
|
|
if err != nil {
|
|
if raiseErrorWhenStreamDataNotMatch {
|
|
return err
|
|
} else {
|
|
log.Warn("stream data not match", "url", url, "data", string(data))
|
|
return nil
|
|
}
|
|
}
|
|
|
|
ch.Write(t)
|
|
return nil
|
|
}
|
|
|
|
routine.Submit(routinepkg.Labels{
|
|
routinepkg.RoutineLabelKeyModule: "http_requests",
|
|
routinepkg.RoutineLabelKeyMethod: "RequestAndParseStream",
|
|
}, func() {
|
|
defer resp.Body.Close()
|
|
|
|
var err error
|
|
if usingLengthPrefixed {
|
|
// at most 30MB a single chunk
|
|
err = parser.LengthPrefixedChunking(resp.Body, 0x0f, 1024*1024*30, processData)
|
|
} else {
|
|
err = parser.LineBasedChunking(resp.Body, 1024*1024*30, func(data []byte) error {
|
|
if len(data) == 0 {
|
|
return nil
|
|
}
|
|
|
|
if bytes.HasPrefix(data, []byte("data:")) {
|
|
// split
|
|
data = data[5:]
|
|
}
|
|
|
|
if bytes.HasPrefix(data, []byte("event:")) {
|
|
// TODO: handle event
|
|
return nil
|
|
}
|
|
|
|
// trim space
|
|
data = bytes.TrimSpace(data)
|
|
|
|
return processData(data)
|
|
})
|
|
}
|
|
|
|
if err != nil {
|
|
ch.WriteError(err)
|
|
}
|
|
|
|
ch.Close()
|
|
})
|
|
|
|
return ch, nil
|
|
}
|
|
|
|
func GetAndParseStream[T any](client *http.Client, url string, options ...HttpOptions) (*stream.Stream[T], error) {
|
|
return RequestAndParseStream[T](client, url, "GET", options...)
|
|
}
|
|
|
|
func PostAndParseStream[T any](client *http.Client, url string, options ...HttpOptions) (*stream.Stream[T], error) {
|
|
return RequestAndParseStream[T](client, url, "POST", options...)
|
|
}
|
|
|
|
func PutAndParseStream[T any](client *http.Client, url string, options ...HttpOptions) (*stream.Stream[T], error) {
|
|
return RequestAndParseStream[T](client, url, "PUT", options...)
|
|
}
|