From 1d7023be1a1400facf002a5751da5615c5164880 Mon Sep 17 00:00:00 2001 From: mudler-agent Date: Wed, 30 Sep 2026 19:23:24 +0200 Subject: [PATCH] test(agentpool): pin the standalone agent contract (#12378) Adds a Ginkgo contract suite for the standalone (LocalAGI-backed) agent service in core/services/agentpool, with a fake OpenAI-compatible LLM and a harness that boots a real AgentPoolService. It pins CRUD, per-user isolation, export and import, pause and resume, the chat and SSE event contract, status and observables, persistence and the raw-key agent lookup. No production code changes. Specs that pin a known gap are named known gap or known defect, so later phases flip them on purpose. Signed-off-by: Ettore Di Giacinto Assisted-by: Claude:claude-opus-5-5 [Claude Code] --- .../agentpool/contract_helpers_test.go | 304 ++++++++++++++++++ .../standalone_chat_contract_test.go | 157 +++++++++ .../agentpool/standalone_contract_test.go | 186 +++++++++++ .../standalone_persistence_contract_test.go | 101 ++++++ .../standalone_status_contract_test.go | 232 +++++++++++++ 5 files changed, 980 insertions(+) create mode 100644 core/services/agentpool/contract_helpers_test.go create mode 100644 core/services/agentpool/standalone_chat_contract_test.go create mode 100644 core/services/agentpool/standalone_contract_test.go create mode 100644 core/services/agentpool/standalone_persistence_contract_test.go create mode 100644 core/services/agentpool/standalone_status_contract_test.go diff --git a/core/services/agentpool/contract_helpers_test.go b/core/services/agentpool/contract_helpers_test.go new file mode 100644 index 000000000..13faff6a2 --- /dev/null +++ b/core/services/agentpool/contract_helpers_test.go @@ -0,0 +1,304 @@ +package agentpool_test + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + + "github.com/mudler/LocalAGI/core/sse" + "github.com/mudler/LocalAGI/core/state" + "github.com/mudler/LocalAGI/core/types" + "github.com/mudler/LocalAI/core/config" + "github.com/mudler/LocalAI/core/services/agentpool" + . "github.com/onsi/gomega" +) + +type fakeLLMRequest struct { + Model string + Stream bool + Messages []map[string]any + Tools []string + ToolChoice any +} + +// fakeLLM is an OpenAI-compatible chat endpoint. The standalone pool reaches its +// LLM over HTTP (apiURL), so a real server is the only seam that exercises the +// whole agent loop without a model. +type fakeLLM struct { + srv *httptest.Server + mu sync.Mutex + reply string + requests []fakeLLMRequest + toolName string + toolArgs string + // failStatus, when non-zero, makes every chat completion fail with that + // HTTP status so a spec can drive the agent's error path. + failStatus int +} + +func newFakeLLM(reply string) *fakeLLM { + f := &fakeLLM{reply: reply} + f.srv = httptest.NewServer(http.HandlerFunc(f.handle)) + return f +} + +func (f *fakeLLM) URL() string { return f.srv.URL } +func (f *fakeLLM) Close() { f.srv.Close() } + +func (f *fakeLLM) SetReply(r string) { + f.mu.Lock() + defer f.mu.Unlock() + f.reply = r +} + +// SetToolCall makes the fake answer with one call to the named function until +// the conversation carries a tool result, then with the plain reply. Keying on +// the tool message rather than a request counter keeps the fake independent of +// how many planning requests the agent makes before it runs the tool. Only +// the LocalAGI counter-action specs use it today; P2 keeps it for the MCP +// tool fixture that replaces them, since the native executor ignores +// Actions. The +// flip side: tool mode stays on until a request carries a role "tool" +// message, so a client that restarts with trimmed history would get the +// tool call again and loop until its iteration cap. +func (f *fakeLLM) SetToolCall(name, argsJSON string) { + f.mu.Lock() + defer f.mu.Unlock() + f.toolName = name + f.toolArgs = argsJSON +} + +// SetFailure makes every chat completion answer with the given HTTP status and +// an OpenAI-style error body, which is what an unreachable or broken backend +// looks like to the agent. Requests are still recorded so a spec can count +// retries. +func (f *fakeLLM) SetFailure(status int) { + f.mu.Lock() + defer f.mu.Unlock() + f.failStatus = status +} + +func (f *fakeLLM) Requests() []fakeLLMRequest { + f.mu.Lock() + defer f.mu.Unlock() + return append([]fakeLLMRequest(nil), f.requests...) +} + +func (f *fakeLLM) handle(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + var req struct { + Model string `json:"model"` + Stream bool `json:"stream"` + Messages []map[string]any `json:"messages"` + Tools []struct { + Function struct { + Name string `json:"name"` + } `json:"function"` + } `json:"tools"` + ToolChoice any `json:"tool_choice"` + } + _ = json.Unmarshal(body, &req) + + var tools []string + for _, t := range req.Tools { + tools = append(tools, t.Function.Name) + } + hasToolResult := false + for _, m := range req.Messages { + if m["role"] == "tool" { + hasToolResult = true + } + } + + f.mu.Lock() + f.requests = append(f.requests, fakeLLMRequest{Model: req.Model, Stream: req.Stream, Messages: req.Messages, Tools: tools, ToolChoice: req.ToolChoice}) + reply := f.reply + failStatus := f.failStatus + var toolCall map[string]any + if f.toolName != "" && !hasToolResult { + toolCall = map[string]any{ + "index": 0, "id": "call_fake", "type": "function", + "function": map[string]any{"name": f.toolName, "arguments": f.toolArgs}, + } + } + f.mu.Unlock() + + if failStatus != 0 { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(failStatus) + _ = json.NewEncoder(w).Encode(map[string]any{ + "error": map[string]any{"message": "fake backend failure", "type": "server_error", "code": failStatus}, + }) + return + } + + message := map[string]any{"role": "assistant", "content": reply} + finish := "stop" + if toolCall != nil { + // "index" belongs only to streaming deltas, so the non-streaming + // message carries a copy without it. + plain := map[string]any{} + for k, v := range toolCall { + if k != "index" { + plain[k] = v + } + } + message = map[string]any{"role": "assistant", "content": "", "tool_calls": []map[string]any{plain}} + finish = "tool_calls" + } + + if !req.Stream { + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{ + "id": "chatcmpl-fake", "object": "chat.completion", "model": req.Model, + "choices": []map[string]any{{ + "index": 0, + "message": message, + "finish_reason": finish, + }}, + "usage": map[string]any{"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}, + }) + return + } + + w.Header().Set("Content-Type", "text/event-stream") + // A write error means the client went away mid-stream; the handler just + // stops writing so a disconnect never blocks or fails the fake. + write := func(s string) bool { + _, err := io.WriteString(w, s) + return err == nil + } + chunk := func(delta map[string]any, finish any) bool { + b, _ := json.Marshal(map[string]any{ + "id": "chatcmpl-fake", "object": "chat.completion.chunk", "model": req.Model, + "choices": []map[string]any{{"index": 0, "delta": delta, "finish_reason": finish}}, + }) + return write(fmt.Sprintf("data: %s\n\n", b)) + } + first := map[string]any{"role": "assistant", "content": reply} + if toolCall != nil { + first = map[string]any{"role": "assistant", "tool_calls": []map[string]any{toolCall}} + } + if !chunk(first, nil) || !chunk(map[string]any{}, finish) || !write("data: [DONE]\n\n") { + return + } + if fl, ok := w.(http.Flusher); ok { + fl.Flush() + } +} + +// startStandalone boots a real standalone AgentPoolService on stateDir with its +// LLM pointed at llmURL. It registers no cleanup itself: callers own Stop(). +func startStandalone(stateDir, llmURL string) *agentpool.AgentPoolService { + cfg := config.NewApplicationConfig() + cfg.AgentPool = config.AgentPoolConfig{ + Enabled: true, + StateDir: stateDir, + APIURL: llmURL, + DefaultModel: "fake-model", + Timeout: "30s", + } + svc, err := agentpool.NewAgentPoolService(cfg) + Expect(err).ToNot(HaveOccurred()) + Expect(svc.Start(context.Background())).To(Succeed()) + return svc +} + +func newAgentConfig(name string) *state.AgentConfig { + return &state.AgentConfig{ + Name: name, + Model: "fake-model", + Description: "contract test agent", + SystemPrompt: "You are a test agent.", + } +} + +// Engine-specific: awaitRunning and collectSSE reach into LocalAGI types +// (agent.Agent via GetAgentForUser, types.NewJob, sse.Manager and +// sse.NewClient). P1 must re-seat them when LocalAGI types leave the service +// signatures; the specs that call them should not need to change. + +// awaitRunning blocks until the agent's Run loop is serving jobs. The pool +// starts Run in a goroutine and LocalAGI's Scheduler.Start and Scheduler.Stop +// are unsynchronized: a Stop (update, delete, svc.Stop) that lands while Start +// is still running can nil the scheduler context under the poll goroutine and +// crash the test binary. Run starts its workers only after Scheduler.Start has +// returned, and jobQueue is unbuffered, so Execute returning proves Start is +// done. The job's context is already cancelled, so the worker finishes it as +// expired without calling the LLM or recording an observable. This depends on +// LocalAGI not short-circuiting a cancelled job before a worker receives it, +// so re-check it when LocalAGI is bumped; the timeout turns a hang into a +// failure if that ever changes. +func awaitRunning(svc *agentpool.AgentPoolService, userID, name string) { + a := svc.GetAgentForUser(userID, name) + Expect(a).ToNot(BeNil()) + ctx, cancel := context.WithCancel(context.Background()) + cancel() + done := make(chan struct{}) + go func() { + defer close(done) + a.Execute(types.NewJob(types.WithContext(ctx))) + }() + Eventually(done, "10s").Should(BeClosed()) +} + +type sseEvent struct { + Name string + Data map[string]any +} + +// collectSSE registers a listener on the agent's SSE manager and records every +// event until stop is called. Subscribe before Chat so nothing is missed. +// The manager replays its last 10 events on Register, so a spec that chats +// twice with a fresh collector sees the first turn's completed status too. +func collectSSE(svc *agentpool.AgentPoolService, userID, name string) (events func() []sseEvent, stop func()) { + mgr := svc.GetSSEManagerForUser(userID, name) + Expect(mgr).ToNot(BeNil()) + // Include the user: the manager keys listeners by ID, so two users' + // same-named agents must never share one if a manager is ever shared. + client := sse.NewClient("contract-" + userID + "-" + name) + mgr.Register(client) + + var mu sync.Mutex + var got []sseEvent + done := make(chan struct{}) + go func() { + for { + select { + case <-done: + return + case env, ok := <-client.Chan(): + if !ok { + return + } + ev := sseEvent{} + for _, line := range strings.Split(env.String(), "\n") { + switch { + case strings.HasPrefix(line, "event:"): + ev.Name = strings.TrimSpace(strings.TrimPrefix(line, "event:")) + case strings.HasPrefix(line, "data:"): + // The parse error is ignored because the hud event's + // data is not JSON; those events keep only their name. + _ = json.Unmarshal([]byte(strings.TrimSpace(strings.TrimPrefix(line, "data:"))), &ev.Data) + } + } + mu.Lock() + got = append(got, ev) + mu.Unlock() + } + } + }() + return func() []sseEvent { + mu.Lock() + defer mu.Unlock() + return append([]sseEvent(nil), got...) + }, func() { + close(done) + mgr.Unregister(client.ID()) + } +} diff --git a/core/services/agentpool/standalone_chat_contract_test.go b/core/services/agentpool/standalone_chat_contract_test.go new file mode 100644 index 000000000..9a2a452c3 --- /dev/null +++ b/core/services/agentpool/standalone_chat_contract_test.go @@ -0,0 +1,157 @@ +package agentpool_test + +import ( + "net/http" + + "github.com/mudler/LocalAI/core/services/agentpool" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var _ = Describe("standalone chat contract", func() { + var ( + llm *fakeLLM + svc *agentpool.AgentPoolService + ) + + BeforeEach(func() { + llm = newFakeLLM("pong") + // Ginkgo runs cleanups LIFO: registering Close first stops the pool + // before its LLM goes away. + DeferCleanup(llm.Close) + svc = startStandalone(GinkgoT().TempDir(), llm.URL()) + DeferCleanup(svc.Stop) + Expect(svc.CreateAgentForUser("alice", newAgentConfig("chatty"))).To(Succeed()) + awaitRunning(svc, "alice", "chatty") + }) + + It("streams the user message, processing, agent reply and completed status over SSE", func() { + events, stop := collectSSE(svc, "alice", "chatty") + defer stop() + + msgID, err := svc.ChatForUser("alice", "chatty", "ping") + Expect(err).ToNot(HaveOccurred()) + Expect(msgID).ToNot(BeEmpty()) + + Eventually(func() []sseEvent { return statusEvents(events(), "completed") }, "30s", "100ms"). + ShouldNot(BeEmpty(), "no completed status event") + + var user, agent, processing bool + for _, e := range events() { + switch { + case e.Name == "json_message" && e.Data["sender"] == "user": + Expect(e.Data["content"]).To(Equal("ping")) + user = true + case e.Name == "json_message" && e.Data["sender"] == "agent": + Expect(e.Data["content"]).To(ContainSubstring("pong")) + // Current standalone shape: the reply id is the id ChatForUser + // returned plus "-agent". The UI correlates on message_id + // (AgentChat.jsx), which the distributed dispatcher sends; a + // native engine may send either, so flip this deliberately. + Expect(e.Data).To(HaveKeyWithValue("id", msgID+"-agent")) + agent = true + case e.Name == "json_message_status" && e.Data["status"] == "processing": + processing = true + } + } + Expect(user).To(BeTrue(), "user json_message") + Expect(processing).To(BeTrue(), "processing status") + Expect(agent).To(BeTrue(), "agent json_message") + }) + + It("sends the user's message to the LLM under the configured model", func() { + _, err := svc.ChatForUser("alice", "chatty", "ping") + Expect(err).ToNot(HaveOccurred()) + + Eventually(llm.Requests, "30s", "100ms").ShouldNot(BeEmpty()) + req := llm.Requests()[0] + Expect(req.Model).To(Equal("fake-model")) + // Observed shape: content is a plain string, not a parts array. A + // switch to parts would change what an OpenAI-compatible backend sees. + Expect(req.Messages).To(ContainElement(And( + HaveKeyWithValue("role", "user"), + HaveKeyWithValue("content", "ping"), + ))) + }) + + // The chat page clears its "processing" state only on an agent + // json_message or a json_error, so a failed turn must end in json_error + // followed by completed, never in silence. The fake fails every request; + // cogito retries the decision 5 times with a linear 1s..5s backoff, so the + // turn settles after about 15s, hence the 60s budget. The error text is + // cogito's wrapped chain and is not pinned. + It("reports a failing LLM as json_error then completed, with no agent reply", func() { + llm.SetFailure(http.StatusInternalServerError) + events, stop := collectSSE(svc, "alice", "chatty") + defer stop() + + _, err := svc.ChatForUser("alice", "chatty", "ping") + Expect(err).ToNot(HaveOccurred()) + + Eventually(func() []sseEvent { return statusEvents(events(), "completed") }, "60s", "100ms"). + ShouldNot(BeEmpty(), "no completed status event") + + errorAt, completedAt := -1, -1 + for i, e := range events() { + switch { + case e.Name == "json_error" && errorAt < 0: + errorAt = i + Expect(e.Data).To(HaveKeyWithValue("error", And(BeAssignableToTypeOf(""), Not(BeEmpty())))) + case e.Name == "json_message_status" && e.Data["status"] == "completed" && completedAt < 0: + completedAt = i + case e.Name == "json_message" && e.Data["sender"] == "agent": + Fail("a failed turn must not produce an agent json_message") + } + } + Expect(errorAt).To(BeNumerically(">=", 0), "no json_error event") + Expect(errorAt).To(BeNumerically("<", completedAt), "json_error must precede completed") + Expect(llm.Requests()).ToNot(BeEmpty()) + }) + + It("reports chat with an unknown agent as ErrAgentNotFound", func() { + _, err := svc.ChatForUser("alice", "ghost", "hi") + Expect(err).To(MatchError(agentpool.ErrAgentNotFound)) + }) + + It("reports chat with a deleted agent as ErrAgentNotFound", func() { + Expect(svc.DeleteAgentForUser("alice", "chatty")).To(Succeed()) + _, err := svc.ChatForUser("alice", "chatty", "hi") + Expect(err).To(MatchError(agentpool.ErrAgentNotFound)) + }) + + It("does not deliver one user's chat events to another user's agent of the same name", func() { + Expect(svc.CreateAgentForUser("bob", newAgentConfig("chatty"))).To(Succeed()) + awaitRunning(svc, "bob", "chatty") + aliceEvents, stopAlice := collectSSE(svc, "alice", "chatty") + defer stopAlice() + bobEvents, stopBob := collectSSE(svc, "bob", "chatty") + defer stopBob() + + _, err := svc.ChatForUser("alice", "chatty", "ping") + Expect(err).ToNot(HaveOccurred()) + Eventually(func() []sseEvent { return statusEvents(aliceEvents(), "completed") }, "30s", "100ms").ShouldNot(BeEmpty()) + + // LocalAGI pushes a "hud" snapshot of each agent's own state every + // second, so bob's stream is not silent; what must never reach it is + // anything produced by alice's chat. + Consistently(func() []string { + var names []string + for _, e := range bobEvents() { + if e.Name != "hud" { + names = append(names, e.Name) + } + } + return names + }, "1s", "100ms").Should(BeEmpty()) + }) +}) + +func statusEvents(events []sseEvent, status string) []sseEvent { + var out []sseEvent + for _, e := range events { + if e.Name == "json_message_status" && e.Data["status"] == status { + out = append(out, e) + } + } + return out +} diff --git a/core/services/agentpool/standalone_contract_test.go b/core/services/agentpool/standalone_contract_test.go new file mode 100644 index 000000000..e88720656 --- /dev/null +++ b/core/services/agentpool/standalone_contract_test.go @@ -0,0 +1,186 @@ +package agentpool_test + +import ( + "encoding/json" + + "github.com/mudler/LocalAI/core/services/agentpool" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var _ = Describe("standalone agent service contract", func() { + var ( + llm *fakeLLM + dir string + ) + + BeforeEach(func() { + llm = newFakeLLM("hello from the fake model") + // DeferCleanup runs after AfterEach and in LIFO order, so registering + // here lets a pool started later stop before its LLM goes away. + DeferCleanup(llm.Close) + dir = GinkgoT().TempDir() + }) + + It("boots against a fake LLM and lists no agents", func() { + svc := startStandalone(dir, llm.URL()) + DeferCleanup(svc.Stop) + + Expect(svc.ListAgentsForUser("")).To(BeEmpty()) + }) + + Context("agent CRUD", func() { + var svc *agentpool.AgentPoolService + BeforeEach(func() { + svc = startStandalone(dir, llm.URL()) + DeferCleanup(svc.Stop) + }) + + It("creates, reads back, updates and deletes an agent", func() { + Expect(svc.CreateAgentForUser("alice", newAgentConfig("helper"))).To(Succeed()) + awaitRunning(svc, "alice", "helper") + + got := svc.GetAgentConfigForUser("alice", "helper") + Expect(got).ToNot(BeNil()) + // The pool stores the key ("alice:helper"); the API must show the bare name. + Expect(got.Name).To(Equal("helper")) + Expect(got.Model).To(Equal("fake-model")) + Expect(svc.ListAgentsForUser("alice")).To(HaveKeyWithValue("helper", true)) + + updated := newAgentConfig("helper") + updated.Description = "changed" + Expect(svc.UpdateAgentForUser("alice", "helper", updated)).To(Succeed()) + // Update restarts the agent, so the new instance needs the same wait. + awaitRunning(svc, "alice", "helper") + Expect(svc.GetAgentConfigForUser("alice", "helper").Description).To(Equal("changed")) + + Expect(svc.DeleteAgentForUser("alice", "helper")).To(Succeed()) + Expect(svc.GetAgentConfigForUser("alice", "helper")).To(BeNil()) + Expect(svc.ListAgentsForUser("alice")).ToNot(HaveKey("helper")) + }) + + It("reports an update of a missing agent as ErrAgentNotFound", func() { + err := svc.UpdateAgentForUser("alice", "ghost", newAgentConfig("ghost")) + Expect(err).To(MatchError(agentpool.ErrAgentNotFound)) + }) + + It("keeps two users' agents with the same name apart", func() { + a := newAgentConfig("shared-name") + a.Description = "alice's" + b := newAgentConfig("shared-name") + b.Description = "bob's" + Expect(svc.CreateAgentForUser("alice", a)).To(Succeed()) + Expect(svc.CreateAgentForUser("bob", b)).To(Succeed()) + awaitRunning(svc, "alice", "shared-name") + awaitRunning(svc, "bob", "shared-name") + + Expect(svc.GetAgentConfigForUser("alice", "shared-name").Description).To(Equal("alice's")) + Expect(svc.GetAgentConfigForUser("bob", "shared-name").Description).To(Equal("bob's")) + + Expect(svc.DeleteAgentForUser("alice", "shared-name")).To(Succeed()) + Expect(svc.GetAgentConfigForUser("alice", "shared-name")).To(BeNil()) + Expect(svc.GetAgentConfigForUser("bob", "shared-name")).ToNot(BeNil()) + + grouped := svc.ListAllAgentsGrouped() + Expect(grouped).To(HaveKey("bob")) + Expect(grouped).ToNot(HaveKey("alice")) + }) + + It("round-trips a config through export and import without a user", func() { + cfg := newAgentConfig("portable") + cfg.Description = "carry me" + Expect(svc.CreateAgentForUser("", cfg)).To(Succeed()) + awaitRunning(svc, "", "portable") + + data, err := svc.ExportAgentForUser("", "portable") + Expect(err).ToNot(HaveOccurred()) + + Expect(svc.DeleteAgentForUser("", "portable")).To(Succeed()) + Expect(svc.ImportAgentForUser("", data)).To(Succeed()) + awaitRunning(svc, "", "portable") + + got := svc.GetAgentConfigForUser("", "portable") + Expect(got).ToNot(BeNil()) + Expect(got.Description).To(Equal("carry me")) + }) + + // Known defect pinned on purpose: export returns the stored config, whose + // name is the pool key, and import refuses ":" in names. A rewrite that + // fixes this must flip this spec rather than silently change behavior. + It("known defect: exports a user's agent under its pool key, which import then rejects", func() { + cfg := newAgentConfig("portable") + cfg.Description = "carry me" + Expect(svc.CreateAgentForUser("alice", cfg)).To(Succeed()) + awaitRunning(svc, "alice", "portable") + + data, err := svc.ExportAgentForUser("alice", "portable") + Expect(err).ToNot(HaveOccurred()) + var out map[string]any + Expect(json.Unmarshal(data, &out)).To(Succeed()) + Expect(out["name"]).To(Equal("alice:portable")) + + Expect(svc.DeleteAgentForUser("alice", "portable")).To(Succeed()) + Expect(svc.ImportAgentForUser("alice", data)).To(MatchError(ContainSubstring("invalid characters"))) + Expect(svc.GetAgentConfigForUser("alice", "portable")).To(BeNil()) + }) + + // Engine-specific: P2/P5 flips this because the native import drops + // unknown fields (connectors, actions) with a warning, so the export + // will no longer carry them. + // P5 strips connectors and actions and the P2 migration reads old configs, + // so record what a config that carries them looks like today. LocalAGI + // logs "Failed to create IRC client" for this fixture because the IRC + // config has no nickname; that is expected and is not a failure. + It("accepts and returns a config that carries connectors and actions", func() { + raw := []byte(`{ + "name": "legacy", + "model": "fake-model", + "description": "old style", + "connectors": [{"type": "irc", "config": "{}"}], + "actions": [{"name": "search", "config": "{}"}] + }`) + Expect(svc.ImportAgentForUser("alice", raw)).To(Succeed()) + awaitRunning(svc, "alice", "legacy") + + data, err := svc.ExportAgentForUser("alice", "legacy") + Expect(err).ToNot(HaveOccurred()) + var out map[string]any + Expect(json.Unmarshal(data, &out)).To(Succeed()) + Expect(out["connectors"]).To(HaveLen(1)) + Expect(out["actions"]).To(HaveLen(1)) + // The P2 migration reads these element shapes. The action name is + // what LocalAGI actually stores, observed as the name sent in, not a + // resolved alias. + Expect(out["connectors"]).To(ConsistOf(And( + HaveKeyWithValue("type", "irc"), + HaveKeyWithValue("config", BeAssignableToTypeOf("")), + ))) + Expect(out["actions"]).To(ConsistOf(And( + HaveKeyWithValue("name", "search"), + HaveKeyWithValue("config", BeAssignableToTypeOf("")), + ))) + }) + }) + + Context("pause and resume", func() { + It("toggles the active flag reported by the list", func() { + svc := startStandalone(dir, llm.URL()) + DeferCleanup(svc.Stop) + Expect(svc.CreateAgentForUser("alice", newAgentConfig("napper"))).To(Succeed()) + awaitRunning(svc, "alice", "napper") + Expect(svc.ListAgentsForUser("alice")).To(HaveKeyWithValue("napper", true)) + + Expect(svc.PauseAgentForUser("alice", "napper")).To(Succeed()) + Expect(svc.ListAgentsForUser("alice")).To(HaveKeyWithValue("napper", false)) + + Expect(svc.ResumeAgentForUser("alice", "napper")).To(Succeed()) + Expect(svc.ListAgentsForUser("alice")).To(HaveKeyWithValue("napper", true)) + }) + + It("reports pausing a missing agent as ErrAgentNotFound", func() { + svc := startStandalone(dir, llm.URL()) + DeferCleanup(svc.Stop) + Expect(svc.PauseAgentForUser("alice", "ghost")).To(MatchError(agentpool.ErrAgentNotFound)) + }) + }) +}) diff --git a/core/services/agentpool/standalone_persistence_contract_test.go b/core/services/agentpool/standalone_persistence_contract_test.go new file mode 100644 index 000000000..0fd782c27 --- /dev/null +++ b/core/services/agentpool/standalone_persistence_contract_test.go @@ -0,0 +1,101 @@ +package agentpool_test + +import ( + "encoding/json" + "os" + "path/filepath" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var _ = Describe("standalone persistence contract", func() { + var ( + llm *fakeLLM + dir string + ) + + BeforeEach(func() { + llm = newFakeLLM("ok") + dir = GinkgoT().TempDir() + DeferCleanup(llm.Close) + }) + + // The P2 migration imports this file, so its layout is an interface. + It("writes pool.json keyed by userID:name with the agent config as value", func() { + svc := startStandalone(dir, llm.URL()) + Expect(svc.CreateAgentForUser("alice", newAgentConfig("keeper"))).To(Succeed()) + awaitRunning(svc, "alice", "keeper") + Expect(svc.CreateAgentForUser("", newAgentConfig("anon"))).To(Succeed()) + awaitRunning(svc, "", "anon") + svc.Stop() + + raw, err := os.ReadFile(filepath.Join(dir, "pool.json")) + Expect(err).ToNot(HaveOccurred()) + var pool map[string]map[string]any + Expect(json.Unmarshal(raw, &pool)).To(Succeed()) + Expect(pool).To(HaveKey("alice:keeper")) + Expect(pool).To(HaveKey("anon")) + Expect(pool["alice:keeper"]["model"]).To(Equal("fake-model")) + // The stored name repeats the key, prefix included: the P2 importer + // strips the prefix, so it depends on this. + Expect(pool["alice:keeper"]["name"]).To(Equal("alice:keeper")) + Expect(pool["anon"]["name"]).To(Equal("anon")) + }) + + It("restores agents after a restart on the same state dir", func() { + svc := startStandalone(dir, llm.URL()) + Expect(svc.CreateAgentForUser("alice", newAgentConfig("survivor"))).To(Succeed()) + awaitRunning(svc, "alice", "survivor") + svc.Stop() + + again := startStandalone(dir, llm.URL()) + DeferCleanup(again.Stop) + awaitRunning(again, "alice", "survivor") + Expect(again.GetAgentConfigForUser("alice", "survivor")).ToNot(BeNil()) + Expect(again.ListAgentsForUser("alice")).To(HaveKey("survivor")) + }) + + // Pause is only an in-memory flag on the running agent: pool.json has no + // status field and no per-agent file records pause, so a restart brings + // the agent back active. Pinned as a known gap for the native-store migration to + // close on purpose rather than by accident. + It("known gap: does not keep a paused agent paused across a restart", func() { + svc := startStandalone(dir, llm.URL()) + Expect(svc.CreateAgentForUser("alice", newAgentConfig("sleeper"))).To(Succeed()) + awaitRunning(svc, "alice", "sleeper") + Expect(svc.PauseAgentForUser("alice", "sleeper")).To(Succeed()) + Expect(svc.ListAgentsForUser("alice")).To(HaveKeyWithValue("sleeper", false)) + svc.Stop() + + again := startStandalone(dir, llm.URL()) + DeferCleanup(again.Stop) + awaitRunning(again, "alice", "sleeper") + Expect(again.ListAgentsForUser("alice")).To(HaveKeyWithValue("sleeper", true)) + }) + + // The /v1/responses interceptor decides "is this model an agent" with + // GetAgent(name) using the raw pool key, with no user prefix. + // Known gap, not a contract: any caller who sends model "alice:mine" runs + // alice's agent (the interceptor has no user check), while alice's own + // request for "mine" falls through; a later fix must not read as a break. + It("known gap: resolves an agent by its raw key for the responses interceptor", func() { + svc := startStandalone(dir, llm.URL()) + DeferCleanup(svc.Stop) + Expect(svc.CreateAgentForUser("", newAgentConfig("global-agent"))).To(Succeed()) + awaitRunning(svc, "", "global-agent") + Expect(svc.CreateAgentForUser("alice", newAgentConfig("mine"))).To(Succeed()) + awaitRunning(svc, "alice", "mine") + + Expect(svc.GetAgent("global-agent")).ToNot(BeNil()) + Expect(svc.GetAgent("alice:mine")).ToNot(BeNil()) + Expect(svc.GetAgent("mine")).To(BeNil(), "current gap: a user's own agent is not found by its bare name, only by its pool key") + Expect(svc.GetAgent("nope")).To(BeNil()) + }) + + It("exposes the state dir it was started with", func() { + svc := startStandalone(dir, llm.URL()) + DeferCleanup(svc.Stop) + Expect(svc.StateDir()).To(Equal(dir)) + }) +}) diff --git a/core/services/agentpool/standalone_status_contract_test.go b/core/services/agentpool/standalone_status_contract_test.go new file mode 100644 index 000000000..f109b32b2 --- /dev/null +++ b/core/services/agentpool/standalone_status_contract_test.go @@ -0,0 +1,232 @@ +package agentpool_test + +import ( + "encoding/json" + "fmt" + + "github.com/mudler/LocalAGI/core/state" + "github.com/mudler/LocalAGI/core/types" + + "github.com/mudler/LocalAI/core/services/agentpool" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var _ = Describe("standalone status and observables contract", func() { + var ( + llm *fakeLLM + svc *agentpool.AgentPoolService + ) + + BeforeEach(func() { + llm = newFakeLLM("done") + // Cleanups run in reverse, so the pool stops before the fake LLM goes away. + DeferCleanup(llm.Close) + svc = startStandalone(GinkgoT().TempDir(), llm.URL()) + DeferCleanup(svc.Stop) + Expect(svc.CreateAgentForUser("alice", newAgentConfig("observed"))).To(Succeed()) + awaitRunning(svc, "alice", "observed") + }) + + // runOnce chats once and waits for the completed status. ChatForUser sends + // that status as soon as Ask returns, before the job finalizers have written + // the finished observable, so callers that read observables use + // settledObservable instead. + runOnce := func() { + events, stop := collectSSE(svc, "alice", "observed") + defer stop() + _, err := svc.ChatForUser("alice", "observed", "go") + Expect(err).ToNot(HaveOccurred()) + Eventually(func() []sseEvent { return statusEvents(events(), "completed") }, "30s", "100ms"). + ShouldNot(BeEmpty(), "no completed status event") + } + + // settleRun chats once with the named agent and returns its observables + // and SSE events after the job finalizers are done. Three observer.Update + // calls follow Finish (the Execute finalizer, the consumeJob finalizer and + // the deferred MakeLastProgressCompletion update, LocalAGI agent.go + // 1182-1187), and Update re-appends an observable whose id is gone, so + // clearing while one is still pending would bring the observable back. Each Update also sends an observable_update event, so the run is + // settled once the root observable carries a completion, the completed + // status went out, and no further observable_update arrives. + settleRun := func(name string) ([]map[string]any, []sseEvent) { + events, stop := collectSSE(svc, "alice", name) + defer stop() + _, err := svc.ChatForUser("alice", name, "go") + Expect(err).ToNot(HaveOccurred()) + + Eventually(func() []sseEvent { return statusEvents(events(), "completed") }, "30s", "100ms"). + ShouldNot(BeEmpty(), "no completed status event") + Eventually(func(g Gomega) { + raw, err := svc.GetAgentObservablesForUser("alice", name) + g.Expect(err).ToNot(HaveOccurred()) + rootDone := false + for _, r := range raw { + var o map[string]any + g.Expect(json.Unmarshal(r, &o)).To(Succeed()) + if _, child := o["parent_id"]; !child { + _, rootDone = o["completion"] + } + } + g.Expect(rootDone).To(BeTrue()) + }, "30s", "100ms").Should(Succeed()) + + updates := func() int { + n := 0 + for _, e := range events() { + if e.Name == "observable_update" { + n++ + } + } + return n + } + var last int + Eventually(func() bool { + n := updates() + stable := n == last + last = n + return stable + }, "10s", "300ms").Should(BeTrue()) + Consistently(updates, "300ms", "50ms").Should(Equal(last)) + + raw, err := svc.GetAgentObservablesForUser("alice", name) + Expect(err).ToNot(HaveOccurred()) + obs := make([]map[string]any, len(raw)) + for i, r := range raw { + Expect(json.Unmarshal(r, &obs[i])).To(Succeed()) + } + return obs, events() + } + + // settledObservable returns the first observable of a settled plain run. + settledObservable := func() map[string]any { + obs, _ := settleRun("observed") + Expect(obs).ToNot(BeEmpty()) + return obs[0] + } + + It("returns an empty observable list before any run", func() { + obs, err := svc.GetAgentObservablesForUser("alice", "observed") + Expect(err).ToNot(HaveOccurred()) + Expect(obs).To(BeEmpty()) + }) + + // Discovery: a plain-content reply (no tool call) is enough for LocalAGI + // to record a "job" observable, so the action fallback was not needed. + // parent_id is not pinned: it is omitempty and a root job has none. + It("records observables after a run with the fields the agent status UI reads", func() { + first := settledObservable() + Expect(first).To(HaveKey("id")) + Expect(first).To(HaveKey("creation")) + Expect(first).To(HaveKey("completion")) + }) + + It("clears observables", func() { + settledObservable() + + Expect(svc.ClearAgentObservablesForUser("alice", "observed")).To(Succeed()) + obs, err := svc.GetAgentObservablesForUser("alice", "observed") + Expect(err).ToNot(HaveOccurred()) + Expect(obs).To(BeEmpty()) + }) + + It("reports observables of a missing agent as ErrAgentNotFound", func() { + _, err := svc.GetAgentObservablesForUser("alice", "ghost") + Expect(err).To(MatchError(agentpool.ErrAgentNotFound)) + Expect(svc.ClearAgentObservablesForUser("alice", "ghost")).To(MatchError(agentpool.ErrAgentNotFound)) + }) + + // LocalAGI only creates a status entry when an action result is recorded, + // so an agent whose runs never called an action looks the same as a + // missing one. Pinned as current behavior, not as a desirable contract. + It("returns a nil status for an agent with no action results, as for a missing one", func() { + Expect(svc.GetAgentStatusForUser("alice", "observed")).To(BeNil()) + runOnce() + Consistently(func() any { return svc.GetAgentStatusForUser("alice", "observed") }, "500ms", "100ms").Should(BeNil()) + Expect(svc.GetAgentStatusForUser("alice", "ghost")).To(BeNil()) + }) + + // Engine-specific: P2 rewrites these three specs because they depend on + // the LocalAGI counter action and the native executor ignores Actions; + // P2 swaps in an MCP tool fixture. The status spec also reads + // types.ActionState, a LocalAGI type that P1 must re-seat when LocalAGI + // types leave the service signatures. + Context("after a run that calls a tool", func() { + BeforeEach(func() { + cfg := newAgentConfig("tooled") + // counter is pure and in-memory, so the action result is + // deterministic without any external service. + cfg.Actions = []state.ActionsConfig{{Name: "counter", Config: "{}"}} + Expect(svc.CreateAgentForUser("alice", cfg)).To(Succeed()) + awaitRunning(svc, "alice", "tooled") + llm.SetToolCall("counter", `{"name":"contract","adjustment":1}`) + }) + + It("records a status entry the status endpoint renders with action, params and result", func() { + settleRun("tooled") + + st := svc.GetAgentStatusForUser("alice", "tooled") + Expect(st).ToNot(BeNil()) + // Select by action name rather than position: the order of + // status entries is a LocalAGI detail, not part of the contract. + var h types.ActionState + Expect(st.Results()).To(ContainElement(Satisfy(func(s types.ActionState) bool { + return s.ActionCurrentState.Action != nil && + s.ActionCurrentState.Action.Definition().Name.String() == "counter" + }), &h)) + Expect(h.ActionCurrentState.Params).To(HaveKeyWithValue("name", "contract")) + Expect(h.Result).To(ContainSubstring("Created counter 'contract'")) + + // Same format string as GetAgentStatusEndpoint: this text is what + // the agent status page shows, so the params must render as JSON + // through ActionParams.String rather than as a Go map. + rendered := fmt.Sprintf("Reasoning: %s\nAction taken: %s\nParameters: %+v\nResult: %s", + h.Reasoning, h.ActionCurrentState.Action.Definition().Name.String(), h.ActionCurrentState.Params, h.Result) + Expect(rendered).To(ContainSubstring("Action taken: counter\n")) + Expect(rendered).To(ContainSubstring(`"name":"contract"`)) + Expect(rendered).To(ContainSubstring("Result: Created counter 'contract'")) + }) + + // The observables tree nests the action under the job that ran it. + // Only the link is pinned: ids, ordering and conversation contents are + // LocalAGI internals. + It("records a child action observable linked to the root job by parent_id", func() { + obs, _ := settleRun("tooled") + Expect(len(obs)).To(BeNumerically(">", 1)) + + var roots, children []map[string]any + for _, o := range obs { + if _, ok := o["parent_id"]; ok { + children = append(children, o) + } else { + roots = append(roots, o) + } + } + Expect(roots).To(HaveLen(1)) + Expect(children).ToNot(BeEmpty()) + for _, c := range children { + // Compared as decoded JSON, whatever type the ids have: the id + // type is a LocalAGI detail, the parent link is the contract. + Expect(c["parent_id"]).To(Equal(roots[0]["id"])) + } + // "action" is the LocalAGI name, visible in AgentStatus.jsx; a + // native engine may name it differently, so flip deliberately. + Expect(children).To(ContainElement(And( + HaveKeyWithValue("name", "action"), + HaveKeyWithValue("completion", HaveKeyWithValue("action_result", ContainSubstring("Created counter"))), + ))) + }) + + It("still delivers the agent's final reply over SSE", func() { + _, events := settleRun("tooled") + var replies []sseEvent + for _, e := range events { + if e.Name == "json_message" && e.Data["sender"] == "agent" { + replies = append(replies, e) + } + } + Expect(replies).ToNot(BeEmpty()) + Expect(replies[0].Data["content"]).To(ContainSubstring("done")) + }) + }) +})