mirror of
https://github.com/langgenius/dify-plugin-daemon.git
synced 2026-07-22 01:35:24 -04:00
888ad788bc
* refactor: introduce local plugin control panel and cleanup environment setup process * fix: args * refactor: new local runtime * temp: stash work for refactor on RemotePluginServer * refactor: unify local runtime lifetime and sperate init environment process * chore: add missing files * stash * refactor: local plugin lifetime control * refactor: complete installation process of control panel * refactor: adapt service layer to new controlpanel * refactor: pluginManager.Install * fix: add routine wrap to InstallServerless, avoid blocking main thread * feat: reinstall serverless runtime * chore: add comments to Reinstall and update confusing naming * refactor: unify install plugin service * refactor: add labels to debugging runtime * refactor: add getters to plugin manager * refactor: split install service to decode/install_task/install service * ??? * refactor: adapt controllers * refactor: session write * refactor: session runtime * Refine install task orchestration (#501) * refactor: installing task * refactor cluster management, decouple lifetime management and cluster * fix cli test command * fix: cleanup TODO comments and implement GracefulStop for instance * feat: add logger to control panel * fix: multiple nil references * refactor: better lifetime control * refactor: better cycle interval * fix(LocalPluginRuntime): prevent returning err when it's not error * fix: avoid adding empty PipExtraArgs * fix: missing errors in Environment init * fix: add truncateMessage to avoid db explosion * cleanup: better lifecycle management * fix: init status at the beginning of installation * optimize: GracefulStop for pluginInstance * refactor: tests * refactor: centralize routine labels (#504) * cleanup: RoutineKey * fix: init routine pool * fix: correctly handle cluster register error * fix: memory leak * fix: add \n to instance write * fix(installer.go): set success to true after succeed for defer func * refactor * fix: missing cwd in testutils * fix: scaleup default runtime nums to 1 when testing * fix: localruntime appconfig in testing module * Update internal/core/local_runtime/load_balancing.go Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> * fix: more efficiency implement in installer_local.go * fix: returns after failing in onDebuggingRuntimeDisconnected * fix: returns after failing in onDebuggingRuntimeDisconnected * fix: splits tests * refactor: naming * refactor: manifest.VersionX * fix: adapt SetDefault to tests * fix: enforce use constants in DBType * fix: generate * fix: linter * cleanup tests * refactor: change package to * cleanup: useless codes * Update internal/cluster/plugin.go Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> * cleanup * refactor: decouple connection_key management from debugging_time * refactor: confused naming * feat: recycle resources to adapt to https://github.com/langgenius/dify-plugin-daemon/pull/500 * refactor: confusing redirecting * fix: support get serverless runtime * fix: race condition in Launching * fix: avoid ManifestValidate in first step of debugging handshake * fix: adding ReleaseAllLocks to finalizers * wtf: what a beautiful code * refactor: rename Stream.Async to Stream.Process * fix: kill process if daed instance was detected * fix: correctly handle failures * fix: consistence of difference interfaces * fix: add stacktrace to panic * fix: only trigger once event * fix: ensure plugin runtime was shutdown * feat: cleanup install tasks * fix: add scale logs --------- Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
152 lines
3.9 KiB
Go
152 lines
3.9 KiB
Go
package testutils
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/langgenius/dify-plugin-daemon/pkg/utils/network"
|
|
"golang.org/x/exp/rand"
|
|
)
|
|
|
|
// FakeOpenAIResponse represents the structure of an OpenAI chat completion response
|
|
type FakeOpenAIResponse struct {
|
|
ID string `json:"id"`
|
|
Object string `json:"object"`
|
|
Created int64 `json:"created"`
|
|
Model string `json:"model"`
|
|
Choices []struct {
|
|
Index int `json:"index"`
|
|
Delta Delta `json:"delta"`
|
|
FinishReason *string `json:"finish_reason"`
|
|
} `json:"choices"`
|
|
}
|
|
|
|
type Delta struct {
|
|
Content string `json:"content"`
|
|
}
|
|
|
|
// StartFakeOpenAIServer starts a fake OpenAI server that streams responses
|
|
// Returns the port number and a cancel function to stop the server
|
|
func StartFakeOpenAIServer() (int, func()) {
|
|
port, err := network.GetRandomPort()
|
|
if err != nil {
|
|
panic(fmt.Sprintf("Failed to get a random port: %v", err))
|
|
}
|
|
|
|
// Find an available port
|
|
listener, err := net.Listen("tcp", fmt.Sprintf(":%d", port))
|
|
if err != nil {
|
|
panic(fmt.Sprintf("Failed to find an available port: %v", err))
|
|
}
|
|
|
|
listener.Close()
|
|
|
|
// Create a new server
|
|
mux := http.NewServeMux()
|
|
server := &http.Server{
|
|
Addr: fmt.Sprintf(":%d", port),
|
|
Handler: mux,
|
|
}
|
|
|
|
// Define the chat completions endpoint
|
|
mux.HandleFunc("/v1/chat/completions", func(w http.ResponseWriter, r *http.Request) {
|
|
// Set headers for streaming response
|
|
w.Header().Set("Content-Type", "text/event-stream")
|
|
w.Header().Set("Cache-Control", "no-cache")
|
|
w.Header().Set("Connection", "keep-alive")
|
|
|
|
// Generate a random ID
|
|
id := fmt.Sprintf("chatcmpl-%d", rand.Intn(1000000))
|
|
|
|
// Create a list of random words to return
|
|
words := []string{
|
|
"hello", "world", "this", "is", "a", "fake", "openai", "server",
|
|
"that", "streams", "responses", "for", "testing", "purposes",
|
|
"only", "it", "will", "return", "words", "every", "hundred",
|
|
"milliseconds", "until", "it", "reaches", "the", "limit",
|
|
"of", "four", "hundred", "tokens", "please", "use", "this",
|
|
"for", "benchmarking", "your", "plugin", "system", "thank",
|
|
"you", "for", "using", "our", "service", "have", "a", "nice",
|
|
"day", "goodbye", "see", "you", "soon", "take", "care",
|
|
}
|
|
|
|
flusher, ok := w.(http.Flusher)
|
|
if !ok {
|
|
http.Error(w, "Streaming not supported", http.StatusInternalServerError)
|
|
return
|
|
}
|
|
|
|
// Stream 400 words, one every 100ms
|
|
for i := 0; i < 100; i++ {
|
|
word := words[i%len(words)]
|
|
|
|
// Add space before words (except the first one)
|
|
if i > 0 {
|
|
word = " " + word
|
|
}
|
|
|
|
response := FakeOpenAIResponse{
|
|
ID: id,
|
|
Object: "chat.completion.chunk",
|
|
Created: time.Now().Unix(),
|
|
Model: "gpt-3.5-turbo",
|
|
Choices: []struct {
|
|
Index int `json:"index"`
|
|
Delta Delta `json:"delta"`
|
|
FinishReason *string `json:"finish_reason"`
|
|
}{
|
|
{
|
|
Index: i,
|
|
Delta: Delta{
|
|
Content: word,
|
|
},
|
|
FinishReason: nil,
|
|
},
|
|
},
|
|
}
|
|
|
|
// For the last message, set finish_reason to "stop"
|
|
if i == 99 {
|
|
response.Choices[0].FinishReason = &[]string{"stop"}[0]
|
|
}
|
|
|
|
data, _ := json.Marshal(response)
|
|
fmt.Fprintf(w, "data: %s\n\n", string(data))
|
|
flusher.Flush()
|
|
|
|
// Sleep for 100ms
|
|
time.Sleep(10 * time.Millisecond)
|
|
|
|
// Check if the client has disconnected
|
|
select {
|
|
case <-r.Context().Done():
|
|
return
|
|
default:
|
|
}
|
|
}
|
|
|
|
// Send the [DONE] message
|
|
fmt.Fprintf(w, "data: [DONE]\n\n")
|
|
flusher.Flush()
|
|
})
|
|
|
|
// Start the server in a goroutine
|
|
go func() {
|
|
if err := server.ListenAndServe(); err != nil && !strings.Contains(err.Error(), "server closed") {
|
|
fmt.Printf("Fake OpenAI server error: %v\n", err)
|
|
}
|
|
}()
|
|
|
|
// Return the port and a cancel function
|
|
return int(port), func() {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
server.Shutdown(ctx)
|
|
}
|
|
}
|