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 <mudler@localai.io>
Assisted-by: Claude:claude-opus-5-5 [Claude Code]
This commit is contained in:
mudler-agent authored and GitHub committed 2026-09-30 19:23:24 +02:00
1 parent f378fe89d0
commit 1d7023be1a
5 files changed
+980

No files matched your search

@@ -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())
}
}
@@ -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
}
@@ -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))
})
})
})
@@ -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))
})
})
@@ -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"))
})
})
})