mirror of
https://github.com/mudler/LocalAI.git
synced 2026-10-10 15:52:29 -04:00
* feat(messaging): add shared subject rules Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(messaging): cover BroadcastRoots, ControlRoots and SubjectRoot Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(messaging): add Broadcaster and enforce subject rules in every carrier Broadcaster is the fan-out half of MessagingClient. The NATS client and the in-memory FakeBus now refuse a subject outside the served roots and any wildcard other than a whole single token, and FakeBus shares MatchSubject instead of its own copy. FakeBus Unsubscribe now removes its own subscription instead of the first one with the same subject. A shared conformance suite in messagingtest runs against both carriers. The distributed e2e specs that used invented test.* subjects, and the one that subscribed with a > filter, now use subjects from subjects.go. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: depend on Broadcaster where only publish and subscribe are used Narrowed to messaging.Broadcaster: nodes/staging_progress.go, nodes/install_progress_publisher.go, galleryop/operation.go, galleryop/service.go, agentpool/user_services.go, agentpool/agent_jobs.go, openresponses/store.go, openresponses/sync.go, syncstate/syncstate.go, finetune/service.go, quantization/service.go and failover/distsync/distsync.go. SubscribeJSON now takes a Broadcaster because it only calls Subscribe, which lets the narrowed consumers use it. Stayed wide: worker/supervisor.go, because its client field also serves the SubscribeReply handlers in worker/lifecycle.go. The request/reply, queue and wiring files (nodes/unloader.go, nodes/file_stager_s3.go, jobs/dispatcher.go, agents/dispatcher.go, agents/events.go, worker/file_staging.go, cli/agent_worker.go, http/app.go) are unchanged by design. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): name the no-route condition and confine the carrier error Consumers matched nats.ErrNoResponders, which names an absence, to demote a node. They now match ErrNoRoute, the control path maps the carrier's failure onto it, and timeouts and worker refusals are pinned as not being no-route. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs(nodes): state which FileStager implementations return ErrNoRoute Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): build backend clients through one node-aware seam Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: describe the distributed transport seams Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: correct comments that overclaim after the seams refactor Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(agent-worker): refuse an unserved LOCALAI_AGENT_SUBJECT at startup The messaging client now refuses a subject whose root no carrier serves. An agent worker started with a custom LOCALAI_AGENT_SUBJECT such as tenant-a.agent.execute used to start and then wait on a subject the frontend never publishes to. After the subject rules landed it exited at subscribe time with an error that did not name the setting. Behaviour change: the worker now checks LOCALAI_AGENT_SUBJECT before it registers or connects, and exits with an error that names the variable and says to use a served subject under the agent root, for example agent.execute. The served roots are not widened: a custom root was never delivered by the frontend, and a wider set would reopen the drift the subject rules exist to close. The flag help and the agent worker docs state the constraint. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(nodes): pin the reactions to ErrNoRoute Three callers react to ErrNoRoute and had no spec: the reconciler's upgrade drain falls back to the legacy forced install, the reconciler marks the node unhealthy when a pending op has no route, and the backend-op fan-out marks the node unhealthy. Each spec drives the real caller with a scripted no-responders reply and reads the result from the registry or the recorded requests. A fourth spec pins the other side: a pending op that times out leaves the node healthy and only counts the attempt, so mapping timeouts onto ErrNoRoute would fail here. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(messaging): pin client subject checks and fail the carrier suite in CI Add specs that call Publish, Request, Subscribe, QueueSubscribe, SubscribeReply and QueueSubscribeReply on a client with no connection. Each call must return ErrUnservedSubject for bogus.thing and ErrUnsupportedWildcard for jobs.>. This proves that the subject check runs before the connection is used, and needs no server. The NATS conformance suite is the only check that runs the subject rules against a real carrier. Before this change it skipped without output when Docker was missing. Now it fails when CI is set, so a Linux runner without Docker cannot hide it. It still skips on local runs and on macOS CI, which has no Docker. Add SubjectNodeBackendInstallProgress to the list of constructors that must build served subjects, and ask contributors to extend the list. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: state what ErrNoRoute may change, and group the distributed guides The seams note said MarkUnhealthy was the only state change allowed on ErrNoRoute. A pending backend op still records the failed attempt, counts toward the reconciler's retry limit and is dead-lettered after the maximum attempts. The note now says that MarkUnhealthy is the only change to the node's own state, and that the per-op accounting is not a verdict about the node. The note also documents that the NATS conformance run fails under CI when Docker is missing. The distributed-seams row moves next to the distributed-state row in the topics table. The liveness ping spec header now says no route is a reason to skip the worker, not proof that the worker is gone. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): give the backend client factory the node id Mechanical: the method gains a nodeID parameter and the eight test fakes are updated. No behaviour change. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): drop the optional node-aware factory The node id is now in the main method, so the optional interface and its helper had no behaviour of their own. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): dial backend probes through the client factory Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): dial workers' file servers through a per-node dialer Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(http): proxy backend logs through the per-node worker dialer The admin backend-logs proxy (list, lines and the WebSocket stream) now reaches a worker through the same per-node dialer as the HTTP file stager, so every frontend-to-worker dial goes through one seam. The shared direct dialer keeps alive for 15s where the proxy used 30s. Harmless for requests bounded at 15s. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(http): keep the backend-logs proxy independent of the admin connection The proxy request had no context before the dialer change and is bounded only by its 15s timeout. Keep it that way. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: move the worker control payloads to workerctl Mechanical move of the request and reply structs, the install progress event and the file payloads out of messaging. The verbs no longer belong to one carrier. No alias is left behind. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(worker): serve the lifecycle verbs through a controlServer The worker registers one handler per verb and a NATS server maps each verb to its subject. Registration errors now name the verb. node.stop is served with SubscribeReply, which is identical on the wire because the handler never replies. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(worker): report install progress through the control sink Install and upgrade now emit download progress through the sink the control server hands them. The debounce and the terminal flush stay in the handler path, built over that sink by the new nodes.NewDebouncedInstallProgressSink, which replaces NewDebouncedInstallProgressPublisher. The subject and payload on the wire are unchanged. The supervisor no longer holds the bus, and installFn and upgradeFn let specs drive both verbs without a gallery. The malformed-request log lines are restored for install, upgrade, backend.delete, model.unload, model.stop and model.delete, with the reply bytes unchanged. The signal adapter is renamed noReply, which also lets worker.go import os/signal without an alias again. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(worker): serve the file-staging verbs through a controlServer An empty list-dir answer is now {} rather than {"files":null}, because the typed reply omits an empty Files slice. The frontend decodes both to a nil slice in nodes/file_stager_s3.go. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(messaging): add WorkQueue and the NATS producer This is the producer side of the competing-consumer seam. The work kinds map one to one to today's subjects and queue groups: task to jobs.new and mcp-ci to jobs.mcp-ci.new (both in group workers), agent-run to agent.execute (group agent-workers). Enqueue publishes the payload as Publish does today, with one JSON marshal. FakeBus now records queue groups and keeps reply handlers so later specs can pin and drive them. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(messaging): add the NATS WorkConsumer An in-flight limit of one runs the handler inline on the delivery goroutine, as the MCP CI consumer does today. Any other limit spawns per delivery, as the agent consumer does. Queue groups are unchanged. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: publish queued work through WorkQueue The job dispatcher, the agent pool and the agent scheduler enqueue through messaging.WorkQueue; the NATS implementation publishes to the same subjects as before. DistributedServices builds the queue next to the NATS client and hands it to the dispatcher and the agent pool, whose distributed mode switch now reads a non-nil WorkQueue. The unused AgentPoolService.SetNATSClient is removed. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: consume queued work through WorkConsumer The agent dispatcher and the MCP CI consumer register through messaging.WorkConsumer. The NATS implementation keeps the inline one-at-a-time model for MCP CI and the per-delivery model for agent runs. handleMCPCIJob reports on the events publisher the carrier hands it instead of a captured client. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: delete the consumers nothing in production reached jobs.new has a producer and no production consumer, and the agent dispatcher's Dispatch was only called from tests. Publishing jobs.new is unchanged. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(mcp): send MCP requests to agent workers through AgentControl Timeouts still honour only the deadline, not cancellation, exactly as today. The NATS no-responders error maps to ErrNoRoute and a timeout does not. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(agent-worker): serve MCP requests and backend.stop through agentRPCServer The agent worker's MCP tool and discovery reply subscriptions and its backend stop listener move behind an unexported agentRPCServer interface, served on NATS by nodes.NATSAgentRPCServer. The handlers become typed mcp.ToolHandler and mcp.DiscoveryHandler values that answer every failure with a reply carrying Error. Queue group (agent-workers), inline execution on the delivery goroutine, the background handler context, the unmarshal error reply texts and the reply-less backend stop subscription are unchanged. The backend stop handler takes the decoded backend name, so it can still close that backend's MCP sessions. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(messaging): remove helpers that only tests used BroadcastRoots, ControlRoots and SubjectRoot had no production caller. The roots spec now asserts every served root through ValidateSubject instead. MatchSubject moves back into the test support package, the only place that used it, with its table. NATSAgentRPCServer drops the subscription list it stored and never read, and NewNATSAgentRPCServer gets a doc comment. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(mcp): round trip the agent RPC server over a real NATS server One spec sends a tool request and a discovery request through NATSAgentControl to NATSAgentRPCServer and checks that the handlers see the decoded requests and the replies come back. It also puts an undecodable body on the tool subject and checks the server answers with an unmarshal error instead of leaving the requester to time out. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: describe the distributed transport seams The developer note now lists the final seams: fan-out, queues, both halves of the control verbs and of agent RPC, and the dial. It records the open items a second carrier has to handle. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test: pin the in-flight limit each queue consumer asks for The work queue specs pin what Consume does for a given limit, but nothing pinned which limit each production consumer passes. Changing the agent worker's MCP CI limit from 1 to 0 would have let MCP CI jobs run concurrently on each worker with every test green. Move the MCP CI Consume call into startMCPCIConsumer with the same wiring and pin that it asks for (WorkMCPCI, 1). Pin that NATSDispatcher.Start asks for (WorkAgentRun, maxConcurrent) for several limits. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: remove helpers the branch left without a caller SubjectJobCancelWildcard lost its last subscriber when the frontend stopped listening on jobs.*.cancel; the NATS permissions and conformance suite spell the subject out, so nothing reads the constant. decodeBackendStopRequest returned a stopAll flag that production dropped and only a test read. decodeBackendStop is now the single decoder with the same semantics: an empty body is stop-all, an empty Backend is stop-all, malformed JSON is an error. stopBackends still derives stop-all from Backend, so no reply changes. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(messaging): keep an explicitly empty agent queue a plain subscription Before the work queue seam the agent worker passed LOCALAI_AGENT_QUEUE straight to QueueSubscribe, so an explicitly empty value made a plain subscription and every agent worker ran every agent run. WithAgentRunRoute replaced an empty queue with agent-workers, which silently changed that. Keep the queue as given once the option is applied. An empty subject still falls back to agent.execute, since it never had a meaning of its own. The flag default stays agent-workers, so only an explicitly empty value reaches this. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: correct comments and record the PR B notes Fix the recordingFactory comment (it also records the parallel flag), document that a negative maxInFlight is unbounded and that Unsubscribe from a handler deadlocks, and say a permanently undecodable payload returns nil. Record controlHandler's undecodable return as a kept exception, and add the second carrier notes to the developer note: the reconciler has no ClientFactory option, the logs proxy honours HTTP_PROXY, verbs one carrier serves need an opt-out, terminal replies come from the result event, and agent runs publish through the NATS-bound EventBridge, which is not an additive change. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> --------- Signed-off-by: Ettore Di Giacinto <mudler@localai.io> Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
1230 lines
44 KiB
Go
1230 lines
44 KiB
Go
package openai
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/labstack/echo/v4"
|
|
"github.com/mudler/LocalAI/core/backend"
|
|
"github.com/mudler/LocalAI/core/config"
|
|
mcpTools "github.com/mudler/LocalAI/core/http/endpoints/mcp"
|
|
"github.com/mudler/LocalAI/core/http/middleware"
|
|
"github.com/mudler/LocalAI/core/schema"
|
|
"github.com/mudler/LocalAI/core/services/cloudproxy"
|
|
"github.com/mudler/LocalAI/pkg/functions"
|
|
reason "github.com/mudler/LocalAI/pkg/reasoning"
|
|
|
|
"github.com/mudler/LocalAI/core/templates"
|
|
pb "github.com/mudler/LocalAI/pkg/grpc/proto"
|
|
"github.com/mudler/LocalAI/pkg/model"
|
|
|
|
"github.com/mudler/xlog"
|
|
)
|
|
|
|
// messageText returns the textual content of a message, preferring the
|
|
// middleware-populated StringContent and falling back to a string Content.
|
|
func messageText(m schema.Message) string {
|
|
if m.StringContent != "" {
|
|
return m.StringContent
|
|
}
|
|
if s, ok := m.Content.(string); ok {
|
|
return s
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// hasSystemMessage reports whether the message slice already contains a
|
|
// non-empty system-role message — used to avoid clobbering a caller-supplied
|
|
// system prompt when the LocalAI Assistant modality is on. Empty / whitespace
|
|
// system turns (historically sent by the web Chat UI) are ignored so they do
|
|
// not suppress the model config system_prompt.
|
|
func hasSystemMessage(messages []schema.Message) bool {
|
|
for _, m := range messages {
|
|
if m.Role == "system" && strings.TrimSpace(messageText(m)) != "" {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// stripEmptySystemMessages drops system-role messages whose content is empty
|
|
// or whitespace-only. An explicit blank system turn would otherwise satisfy
|
|
// tokenizer chat templates' `messages[0].role == "system"` check and suppress
|
|
// both the model's configured system_prompt and any template default.
|
|
func stripEmptySystemMessages(messages []schema.Message) []schema.Message {
|
|
out := messages[:0:0]
|
|
for _, m := range messages {
|
|
if m.Role == "system" && strings.TrimSpace(messageText(m)) == "" {
|
|
continue
|
|
}
|
|
out = append(out, m)
|
|
}
|
|
return out
|
|
}
|
|
|
|
// normalizeLateSystemMessages handles system-role messages that appear after the
|
|
// leading system block, according to template.system_messages_after_first:
|
|
// "merge" folds them into the first system message (created if absent), "user"
|
|
// forwards them as user-role turns at their original position. Any other value
|
|
// returns the messages unchanged. Needed for tokenizer templates that reject
|
|
// late system turns (Qwen3.8: "System message must be at the beginning") while
|
|
// agent frameworks append instructions mid-conversation.
|
|
func normalizeLateSystemMessages(messages []schema.Message, mode string) []schema.Message {
|
|
if mode != "merge" && mode != "user" {
|
|
return messages
|
|
}
|
|
lead := 0
|
|
for lead < len(messages) && messages[lead].Role == "system" {
|
|
lead++
|
|
}
|
|
late := false
|
|
for _, m := range messages[lead:] {
|
|
if m.Role == "system" {
|
|
late = true
|
|
break
|
|
}
|
|
}
|
|
if !late {
|
|
return messages
|
|
}
|
|
out := make([]schema.Message, 0, len(messages)+1)
|
|
out = append(out, messages[:lead]...)
|
|
if mode == "merge" && lead == 0 {
|
|
out = append(out, schema.Message{Role: "system"})
|
|
}
|
|
for _, m := range messages[lead:] {
|
|
if m.Role != "system" {
|
|
out = append(out, m)
|
|
continue
|
|
}
|
|
text := strings.TrimSpace(messageText(m))
|
|
if text == "" {
|
|
continue
|
|
}
|
|
switch mode {
|
|
case "merge":
|
|
first := &out[0]
|
|
joined := strings.TrimSpace(messageText(*first))
|
|
if joined != "" {
|
|
joined += "\n\n"
|
|
}
|
|
joined += text
|
|
first.Content = joined
|
|
first.StringContent = joined
|
|
case "user":
|
|
m.Role = "user"
|
|
m.Content = text
|
|
m.StringContent = text
|
|
out = append(out, m)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// mergeToolCallDeltas merges streaming tool call deltas into complete tool calls.
|
|
// In SSE streaming, a single tool call arrives as multiple chunks sharing the same Index:
|
|
// the first chunk carries the ID, Type, and Name; subsequent chunks append to Arguments.
|
|
func mergeToolCallDeltas(existing []schema.ToolCall, deltas []schema.ToolCall) []schema.ToolCall {
|
|
byIndex := make(map[int]int, len(existing)) // tool call Index -> position in slice
|
|
for i, tc := range existing {
|
|
byIndex[tc.Index] = i
|
|
}
|
|
for _, d := range deltas {
|
|
pos, found := byIndex[d.Index]
|
|
if !found {
|
|
byIndex[d.Index] = len(existing)
|
|
existing = append(existing, d)
|
|
continue
|
|
}
|
|
// Merge into existing entry
|
|
tc := &existing[pos]
|
|
if d.ID != "" {
|
|
tc.ID = d.ID
|
|
}
|
|
if d.Type != "" {
|
|
tc.Type = d.Type
|
|
}
|
|
if d.FunctionCall.Name != "" {
|
|
tc.FunctionCall.Name = d.FunctionCall.Name
|
|
}
|
|
tc.FunctionCall.Arguments += d.FunctionCall.Arguments
|
|
}
|
|
return existing
|
|
}
|
|
|
|
// applyAutoparserOverride replaces the Go-side reasoning-extraction result with
|
|
// the C++ autoparser's classified ChatDeltas when those deltas contain
|
|
// actionable content or reasoning. It preserves the original logprobs.
|
|
//
|
|
// When the autoparser did not classify any reasoning (deltaReasoning == "") but
|
|
// deltaContent still carries an unparsed reasoning tag pair (e.g. the
|
|
// non-jinja "pure content" fallback path on a <think> model — issue #9985),
|
|
// the Go-side reasoning extractor is run on deltaContent as a defensive
|
|
// fallback so <think>…</think> blocks do not leak into the OpenAI `content`
|
|
// field.
|
|
func applyAutoparserOverride(
|
|
chatDeltas []*pb.ChatDelta,
|
|
thinkingStartToken string,
|
|
reasoningConfig reason.Config,
|
|
existing []schema.Choice,
|
|
) []schema.Choice {
|
|
if len(chatDeltas) == 0 {
|
|
return existing
|
|
}
|
|
deltaContent := functions.ContentFromChatDeltas(chatDeltas)
|
|
deltaReasoning := functions.ReasoningFromChatDeltas(chatDeltas)
|
|
if deltaContent == "" && deltaReasoning == "" {
|
|
return existing
|
|
}
|
|
// Fallback for non-jinja models (issue #9985): when the C++ autoparser
|
|
// did not classify reasoning but the raw content still contains a known
|
|
// reasoning tag pair, run Go-side extraction on the content so that the
|
|
// <think>…</think> block does not leak into the OpenAI `content` field.
|
|
// When the autoparser DID populate ReasoningContent, leave its
|
|
// content/reasoning split alone — trust the parser. We replace
|
|
// deltaContent unconditionally because ExtractReasoningWithConfig is a
|
|
// no-op when no tag pair matches; this also strips empty thinking
|
|
// blocks like "<think></think>" that some models emit when reasoning
|
|
// is disabled.
|
|
if deltaReasoning == "" && deltaContent != "" {
|
|
// Complete-response extraction: only honor a prefilled <think> start
|
|
// token when deltaContent actually closes the reasoning block. Without
|
|
// it the model answered directly and the whole answer must stay in
|
|
// content rather than be swallowed as unclosed reasoning. See
|
|
// reason.ExtractReasoningComplete.
|
|
deltaReasoning, deltaContent = reason.ExtractReasoningComplete(deltaContent, thinkingStartToken, reasoningConfig)
|
|
}
|
|
xlog.Debug("[ChatDeltas] non-SSE no-tools: overriding result with C++ autoparser deltas",
|
|
"content_len", len(deltaContent), "reasoning_len", len(deltaReasoning))
|
|
stopReason := FinishReasonStop
|
|
message := &schema.Message{Role: "assistant", Content: &deltaContent}
|
|
if deltaReasoning != "" {
|
|
message.Reasoning = &deltaReasoning
|
|
}
|
|
newChoice := schema.Choice{FinishReason: &stopReason, Index: 0, Message: message}
|
|
if len(existing) > 0 && existing[0].Logprobs != nil {
|
|
newChoice.Logprobs = existing[0].Logprobs
|
|
}
|
|
return []schema.Choice{newChoice}
|
|
}
|
|
|
|
// ChatEndpoint is the OpenAI Completion API endpoint https://platform.openai.com/docs/api-reference/chat/create
|
|
// @Summary Generate a chat completions for a given prompt and model.
|
|
// @Tags inference
|
|
// @Param request body schema.OpenAIRequest true "query params"
|
|
// @Success 200 {object} schema.OpenAIResponse "Response"
|
|
// @Router /v1/chat/completions [post]
|
|
func ChatEndpoint(cl *config.ModelConfigLoader, ml *model.ModelLoader, evaluator *templates.Evaluator, startupOptions *config.ApplicationConfig, agentControl mcpTools.AgentControl, assistantHolder *mcpTools.LocalAIAssistantHolder, compressor middleware.ChatCompressor) echo.HandlerFunc {
|
|
return func(c echo.Context) error {
|
|
var textContentToReturn string
|
|
id := uuid.New().String()
|
|
created := int(time.Now().Unix())
|
|
|
|
input, ok := c.Get(middleware.CONTEXT_LOCALS_KEY_LOCALAI_REQUEST).(*schema.OpenAIRequest)
|
|
if !ok || input.Model == "" {
|
|
return echo.ErrBadRequest
|
|
}
|
|
|
|
extraUsage := c.Request().Header.Get("Extra-Usage") != ""
|
|
|
|
config, ok := c.Get(middleware.CONTEXT_LOCALS_KEY_MODEL_CONFIG).(*config.ModelConfig)
|
|
if !ok || config == nil {
|
|
return echo.ErrBadRequest
|
|
}
|
|
|
|
xlog.Debug("Chat endpoint configuration read", "config", config)
|
|
|
|
// Drop blank system turns from the web UI (and similar clients) so they
|
|
// cannot suppress the model YAML system_prompt / tokenizer defaults.
|
|
input.Messages = stripEmptySystemMessages(input.Messages)
|
|
input.Messages = normalizeLateSystemMessages(input.Messages, config.TemplateConfig.SystemMessagesAfterFirst)
|
|
|
|
// Tokenizer-template models pass messages through to the backend as-is,
|
|
// so apply the configured system_prompt when the request did not supply
|
|
// one. Go-template models already receive SystemPrompt via PromptTemplateData.
|
|
if config.TemplateConfig.UseTokenizerTemplate && config.SystemPrompt != "" && !hasSystemMessage(input.Messages) {
|
|
prompt := config.SystemPrompt
|
|
input.Messages = append([]schema.Message{{Role: "system", Content: prompt, StringContent: prompt}}, input.Messages...)
|
|
}
|
|
|
|
// Cloud-proxy bail. Bypasses the local pipeline (templating,
|
|
// MCP injection, gRPC backend) and forwards via the cloud-
|
|
// proxy backend, which does the outbound HTTP. Request-side PII
|
|
// redaction already ran in the middleware; the response is
|
|
// forwarded unmodified.
|
|
if config.IsCloudProxyBackendPassthrough() {
|
|
if err := middleware.CompressChatRequest(c, compressor); err != nil {
|
|
return err
|
|
}
|
|
return forwardCloudProxyOpenAIViaBackend(c, config, input, ml, startupOptions)
|
|
}
|
|
|
|
funcs := input.Functions
|
|
shouldUseFn := len(input.Functions) > 0 && config.ShouldUseFunctions()
|
|
strictMode := false
|
|
|
|
// MCP tool injection: when mcp_servers is set in metadata and model has MCP config
|
|
var mcpExecutor mcpTools.ToolExecutor
|
|
mcpServers := mcpTools.MCPServersFromMetadata(input.Metadata)
|
|
|
|
// LocalAI Assistant modality: an admin opted into the in-process MCP
|
|
// admin tool surface. Runs *before* the regular MCP block — when both
|
|
// are set, the assistant tools win (the admin cannot mix them with
|
|
// per-model MCP servers in the same chat session by design).
|
|
assistantMode := mcpTools.LocalAIAssistantFromMetadata(input.Metadata)
|
|
if assistantMode {
|
|
if err := requireAssistantAccess(c, startupOptions.Auth.Enabled); err != nil {
|
|
return err
|
|
}
|
|
// Read the disable flag live: an admin can flip it via /api/settings
|
|
// and the next request must see the change without a restart.
|
|
if startupOptions.DisableLocalAIAssistant {
|
|
return echo.NewHTTPError(http.StatusServiceUnavailable, "LocalAI Assistant is disabled on this server")
|
|
}
|
|
if assistantHolder == nil || !assistantHolder.HasTools() {
|
|
return echo.NewHTTPError(http.StatusServiceUnavailable, "LocalAI Assistant is not available on this server")
|
|
}
|
|
mcpExecutor = assistantHolder.Executor()
|
|
mcpFuncs, discErr := mcpExecutor.DiscoverTools(c.Request().Context())
|
|
if discErr != nil {
|
|
xlog.Error("Failed to discover LocalAI Assistant tools", "error", discErr)
|
|
return echo.NewHTTPError(http.StatusInternalServerError, "discover assistant tools: "+discErr.Error())
|
|
}
|
|
for _, fn := range mcpFuncs {
|
|
funcs = append(funcs, fn)
|
|
input.Tools = append(input.Tools, functions.Tool{Type: "function", Function: fn})
|
|
}
|
|
shouldUseFn = len(funcs) > 0 && config.ShouldUseFunctions()
|
|
|
|
// Prepend the embedded system prompt unless the caller supplied
|
|
// their own system message. Why: the prompt is what teaches the
|
|
// model the safety rules and recipes. If a caller already has a
|
|
// system message they're responsible for keeping the assistant
|
|
// safe, so we leave it alone.
|
|
if !hasSystemMessage(input.Messages) {
|
|
prompt := assistantHolder.SystemPrompt()
|
|
input.Messages = append([]schema.Message{{Role: "system", Content: prompt, StringContent: prompt}}, input.Messages...)
|
|
}
|
|
|
|
xlog.Debug("LocalAI Assistant tools injected", "count", len(mcpFuncs))
|
|
}
|
|
|
|
// MCP prompt and resource injection (extracted before tool injection)
|
|
mcpPromptName, mcpPromptArgs := mcpTools.MCPPromptFromMetadata(input.Metadata)
|
|
mcpResourceURIs := mcpTools.MCPResourcesFromMetadata(input.Metadata)
|
|
|
|
if (len(mcpServers) > 0 || mcpPromptName != "" || len(mcpResourceURIs) > 0) && (config.MCP.Servers != "" || config.MCP.Stdio != "") {
|
|
remote, stdio, mcpErr := config.MCP.MCPConfigFromYAML()
|
|
if mcpErr == nil {
|
|
mcpExecutor = mcpTools.NewToolExecutor(c.Request().Context(), agentControl, config.Name, remote, stdio, mcpServers)
|
|
|
|
// Prompt and resource injection (pre-processing step — resolves locally regardless of distributed mode)
|
|
namedSessions, sessErr := mcpTools.NamedSessionsFromMCPConfig(config.Name, remote, stdio, mcpServers)
|
|
if sessErr == nil && len(namedSessions) > 0 {
|
|
mcpCtx, _ := mcpTools.InjectMCPContext(c.Request().Context(), namedSessions, mcpPromptName, mcpPromptArgs, mcpResourceURIs)
|
|
if mcpCtx != nil {
|
|
input.Messages = append(mcpCtx.PromptMessages, input.Messages...)
|
|
mcpTools.AppendResourceSuffix(input.Messages, mcpCtx.ResourceSuffix)
|
|
}
|
|
}
|
|
|
|
// Tool injection via executor
|
|
if mcpExecutor.HasTools() {
|
|
mcpFuncs, discErr := mcpExecutor.DiscoverTools(c.Request().Context())
|
|
if discErr == nil {
|
|
for _, fn := range mcpFuncs {
|
|
funcs = append(funcs, fn)
|
|
input.Tools = append(input.Tools, functions.Tool{Type: "function", Function: fn})
|
|
}
|
|
shouldUseFn = len(funcs) > 0 && config.ShouldUseFunctions()
|
|
xlog.Debug("MCP tools injected", "count", len(mcpFuncs), "total_funcs", len(funcs))
|
|
} else {
|
|
xlog.Error("Failed to discover MCP tools", "error", discErr)
|
|
}
|
|
}
|
|
} else {
|
|
xlog.Error("Failed to parse MCP config", "error", mcpErr)
|
|
}
|
|
}
|
|
if err := middleware.CompressChatRequest(c, compressor); err != nil {
|
|
return err
|
|
}
|
|
|
|
xlog.Debug("Tool call routing decision",
|
|
"shouldUseFn", shouldUseFn,
|
|
"len(input.Functions)", len(input.Functions),
|
|
"len(input.Tools)", len(input.Tools),
|
|
"config.ShouldUseFunctions()", config.ShouldUseFunctions(),
|
|
"config.FunctionToCall()", config.FunctionToCall(),
|
|
)
|
|
|
|
for _, f := range input.Functions {
|
|
if f.Strict {
|
|
strictMode = true
|
|
break
|
|
}
|
|
}
|
|
|
|
// Allow the user to set custom actions via config file
|
|
// to be "embedded" in each model
|
|
noActionName := "answer"
|
|
noActionDescription := "use this action to answer without performing any action"
|
|
|
|
if config.FunctionsConfig.NoActionFunctionName != "" {
|
|
noActionName = config.FunctionsConfig.NoActionFunctionName
|
|
}
|
|
if config.FunctionsConfig.NoActionDescriptionName != "" {
|
|
noActionDescription = config.FunctionsConfig.NoActionDescriptionName
|
|
}
|
|
|
|
// If we are using a response format, we need to generate a grammar for it
|
|
if config.ResponseFormatMap != nil {
|
|
d := schema.ChatCompletionResponseFormat{}
|
|
dat, err := json.Marshal(config.ResponseFormatMap)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
err = json.Unmarshal(dat, &d)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
switch d.Type {
|
|
case "json_object":
|
|
input.Grammar = functions.JSONBNF
|
|
case "json_schema":
|
|
d := schema.JsonSchemaRequest{}
|
|
dat, err := json.Marshal(config.ResponseFormatMap)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
err = json.Unmarshal(dat, &d)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fs := &functions.JSONFunctionStructure{
|
|
AnyOf: []functions.Item{d.JsonSchema.Schema},
|
|
}
|
|
g, err := fs.Grammar(config.FunctionsConfig.GrammarOptions()...)
|
|
if err == nil {
|
|
input.Grammar = g
|
|
} else {
|
|
xlog.Error("Failed generating grammar", "error", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
config.Grammar = input.Grammar
|
|
|
|
if shouldUseFn {
|
|
xlog.Debug("Response needs to process functions")
|
|
}
|
|
|
|
switch {
|
|
// Generates grammar with internal's LocalAI engine
|
|
case (!config.FunctionsConfig.GrammarConfig.NoGrammar || strictMode) && shouldUseFn:
|
|
noActionGrammar := functions.Function{
|
|
Name: noActionName,
|
|
Description: noActionDescription,
|
|
Parameters: map[string]any{
|
|
"properties": map[string]any{
|
|
"message": map[string]any{
|
|
"type": "string",
|
|
"description": "The message to reply the user with",
|
|
},
|
|
},
|
|
},
|
|
}
|
|
|
|
// Append the no action function
|
|
if !config.FunctionsConfig.DisableNoAction && !strictMode {
|
|
funcs = append(funcs, noActionGrammar)
|
|
}
|
|
|
|
// Force picking one of the functions by the request
|
|
if config.FunctionToCall() != "" {
|
|
funcs = funcs.Select(config.FunctionToCall())
|
|
}
|
|
|
|
// Update input grammar or json_schema based on use_llama_grammar option
|
|
jsStruct := config.FunctionsConfig.ToJSONStructure(funcs)
|
|
g, err := jsStruct.Grammar(config.FunctionsConfig.GrammarOptions()...)
|
|
if err == nil {
|
|
config.Grammar = g
|
|
} else {
|
|
xlog.Error("Failed generating grammar", "error", err)
|
|
}
|
|
case input.JSONFunctionGrammarObject != nil:
|
|
g, err := input.JSONFunctionGrammarObject.Grammar(config.FunctionsConfig.GrammarOptions()...)
|
|
if err == nil {
|
|
config.Grammar = g
|
|
} else {
|
|
xlog.Error("Failed generating grammar", "error", err)
|
|
}
|
|
|
|
default:
|
|
// Force picking one of the functions by the request
|
|
if config.FunctionToCall() != "" {
|
|
funcs = funcs.Select(config.FunctionToCall())
|
|
}
|
|
}
|
|
|
|
// process functions if we have any defined or if we have a function call string
|
|
|
|
// functions are not supported in stream mode (yet?)
|
|
toStream := input.Stream
|
|
|
|
xlog.Debug("Parameters", "config", config)
|
|
|
|
var predInput string
|
|
|
|
// If we are using the tokenizer template, we don't need to process the messages
|
|
// unless we are processing functions
|
|
if !config.TemplateConfig.UseTokenizerTemplate {
|
|
predInput = evaluator.TemplateMessages(*input, input.Messages, config, funcs, shouldUseFn)
|
|
|
|
xlog.Debug("Prompt (after templating)", "prompt", predInput)
|
|
if config.Grammar != "" {
|
|
xlog.Debug("Grammar", "grammar", config.Grammar)
|
|
}
|
|
}
|
|
|
|
switch {
|
|
case toStream:
|
|
|
|
xlog.Debug("Stream request received")
|
|
c.Response().Header().Set("Content-Type", "text/event-stream")
|
|
c.Response().Header().Set("Cache-Control", "no-cache")
|
|
c.Response().Header().Set("Connection", "keep-alive")
|
|
c.Response().Header().Set("X-Correlation-ID", id)
|
|
|
|
mcpStreamMaxIterations := 10
|
|
if config.Agent.MaxIterations > 0 {
|
|
mcpStreamMaxIterations = config.Agent.MaxIterations
|
|
}
|
|
hasMCPToolsStream := mcpExecutor != nil && mcpExecutor.HasTools()
|
|
|
|
for mcpStreamIter := 0; mcpStreamIter <= mcpStreamMaxIterations; mcpStreamIter++ {
|
|
// Re-template on MCP iterations
|
|
if mcpStreamIter > 0 {
|
|
if err := middleware.CompressChatRequest(c, compressor); err != nil {
|
|
fmt.Fprintf(c.Response().Writer, "data: {\"error\":{\"message\":%q,\"type\":\"context_compression_error\"}}\n\n", err.Error())
|
|
fmt.Fprintf(c.Response().Writer, "data: [DONE]\n\n")
|
|
c.Response().Flush()
|
|
return nil
|
|
}
|
|
if !config.TemplateConfig.UseTokenizerTemplate {
|
|
predInput = evaluator.TemplateMessages(*input, input.Messages, config, funcs, shouldUseFn)
|
|
xlog.Debug("MCP stream re-templating", "iteration", mcpStreamIter)
|
|
}
|
|
}
|
|
|
|
responses := make(chan schema.OpenAIResponse)
|
|
ended := make(chan streamWorkerResult, 1)
|
|
|
|
go func() {
|
|
if !shouldUseFn {
|
|
u, err := processStream(predInput, input, config, cl, startupOptions, ml, responses, id, created)
|
|
ended <- streamWorkerResult{usage: u, err: err}
|
|
} else {
|
|
u, err := processStreamWithTools(noActionName, predInput, input, config, cl, startupOptions, ml, responses, id, created, &textContentToReturn)
|
|
ended <- streamWorkerResult{usage: u, err: err}
|
|
}
|
|
}()
|
|
|
|
var finalUsage backend.TokenUsage
|
|
// wroteChunk records whether this iteration sent anything. The
|
|
// chunks go to the raw writer, so c.Response().Committed stays
|
|
// false and cannot tell.
|
|
wroteChunk := false
|
|
toolsCalled := false
|
|
var collectedToolCalls []schema.ToolCall
|
|
var collectedContent string
|
|
|
|
LOOP:
|
|
for {
|
|
select {
|
|
case <-input.Context.Done():
|
|
// Context was cancelled (client disconnected or request cancelled)
|
|
xlog.Debug("Request context cancelled, stopping stream")
|
|
input.Cancel()
|
|
break LOOP
|
|
case ev := <-responses:
|
|
if len(ev.Choices) == 0 {
|
|
xlog.Debug("No choices in the response, skipping")
|
|
continue
|
|
}
|
|
if len(ev.Choices[0].Delta.ToolCalls) > 0 {
|
|
toolsCalled = true
|
|
// Collect and merge tool call deltas for MCP execution
|
|
if hasMCPToolsStream {
|
|
collectedToolCalls = mergeToolCallDeltas(collectedToolCalls, ev.Choices[0].Delta.ToolCalls)
|
|
}
|
|
}
|
|
// Extract the raw content delta string once per chunk;
|
|
// both the MCP collector and the PII filter need it
|
|
// and the type-switch is otherwise duplicated.
|
|
var rawContent string
|
|
haveContent := false
|
|
if ev.Choices[0].Delta != nil && ev.Choices[0].Delta.Content != nil {
|
|
switch v := ev.Choices[0].Delta.Content.(type) {
|
|
case string:
|
|
rawContent = v
|
|
haveContent = true
|
|
case *string:
|
|
if v != nil {
|
|
rawContent = *v
|
|
haveContent = true
|
|
}
|
|
}
|
|
}
|
|
// Collect content for MCP conversation history and automatic tool parsing fallback.
|
|
// We collect the RAW (unfiltered) content so the model's tool-call
|
|
// markup keeps parsing correctly even when PII redaction would mask
|
|
// substrings.
|
|
if (hasMCPToolsStream || config.FunctionsConfig.AutomaticToolParsingFallback) && haveContent {
|
|
collectedContent += rawContent
|
|
}
|
|
respData, err := json.Marshal(ev)
|
|
if err != nil {
|
|
xlog.Debug("Failed to marshal response", "error", err)
|
|
input.Cancel()
|
|
continue
|
|
}
|
|
xlog.Debug("Sending chunk", "chunk", string(respData))
|
|
_, err = fmt.Fprintf(c.Response().Writer, "data: %s\n\n", string(respData))
|
|
if err != nil {
|
|
xlog.Debug("Sending chunk failed", "error", err)
|
|
input.Cancel()
|
|
return err
|
|
}
|
|
wroteChunk = true
|
|
c.Response().Flush()
|
|
case res := <-ended:
|
|
if res.err == nil {
|
|
finalUsage = res.usage
|
|
break LOOP
|
|
}
|
|
xlog.Error("Stream ended with error", "error", res.err)
|
|
|
|
// Nothing was sent yet, so the status line is still
|
|
// open: answer with a real HTTP error instead of an
|
|
// error chunk on a 200 stream. Clients then handle it
|
|
// like the same failure on a non-streaming request.
|
|
if mcpStreamIter == 0 && !wroteChunk {
|
|
h := c.Response().Header()
|
|
h.Del("Content-Type")
|
|
h.Del("Cache-Control")
|
|
h.Del("Connection")
|
|
return backendRequestError(res.err)
|
|
}
|
|
|
|
errorResp := schema.ErrorResponse{
|
|
Error: &schema.APIError{
|
|
Message: res.err.Error(),
|
|
Type: "server_error",
|
|
Code: "server_error",
|
|
},
|
|
}
|
|
respData, marshalErr := json.Marshal(errorResp)
|
|
if marshalErr != nil {
|
|
xlog.Error("Failed to marshal error response", "error", marshalErr)
|
|
fmt.Fprintf(c.Response().Writer, "data: {\"error\":{\"message\":\"Internal error\",\"type\":\"server_error\"}}\n\n")
|
|
} else {
|
|
fmt.Fprintf(c.Response().Writer, "data: %s\n\n", respData)
|
|
}
|
|
fmt.Fprintf(c.Response().Writer, "data: [DONE]\n\n")
|
|
c.Response().Flush()
|
|
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// Drain responses channel to unblock the background goroutine if it's
|
|
// still trying to send (e.g., after client disconnect). The goroutine
|
|
// calls close(responses) when done, which terminates the drain.
|
|
if input.Context.Err() != nil {
|
|
go func() {
|
|
for range responses {
|
|
}
|
|
}()
|
|
<-ended
|
|
}
|
|
|
|
// MCP streaming tool execution: if we collected MCP tool calls, execute and loop
|
|
if hasMCPToolsStream && toolsCalled && len(collectedToolCalls) > 0 {
|
|
var hasMCPCalls bool
|
|
for _, tc := range collectedToolCalls {
|
|
if mcpExecutor != nil && mcpExecutor.IsTool(tc.FunctionCall.Name) {
|
|
hasMCPCalls = true
|
|
break
|
|
}
|
|
}
|
|
if hasMCPCalls {
|
|
// Append assistant message with tool_calls
|
|
assistantMsg := schema.Message{
|
|
Role: "assistant",
|
|
Content: collectedContent,
|
|
ToolCalls: collectedToolCalls,
|
|
}
|
|
input.Messages = append(input.Messages, assistantMsg)
|
|
|
|
// Execute MCP tool calls and stream results as tool_result events
|
|
for _, tc := range collectedToolCalls {
|
|
if mcpExecutor == nil || !mcpExecutor.IsTool(tc.FunctionCall.Name) {
|
|
continue
|
|
}
|
|
xlog.Debug("Executing MCP tool (stream)", "tool", tc.FunctionCall.Name, "iteration", mcpStreamIter)
|
|
toolResult, toolErr := mcpExecutor.ExecuteTool(c.Request().Context(), tc.FunctionCall.Name, tc.FunctionCall.Arguments)
|
|
if toolErr != nil {
|
|
xlog.Error("MCP tool execution failed", "tool", tc.FunctionCall.Name, "error", toolErr)
|
|
toolResult = fmt.Sprintf("Error: %v", toolErr)
|
|
}
|
|
input.Messages = append(input.Messages, schema.Message{
|
|
Role: "tool",
|
|
Content: toolResult,
|
|
StringContent: toolResult,
|
|
ToolCallID: tc.ID,
|
|
Name: tc.FunctionCall.Name,
|
|
})
|
|
|
|
// Stream tool result event to client
|
|
mcpEvent := map[string]any{
|
|
"type": "mcp_tool_result",
|
|
"name": tc.FunctionCall.Name,
|
|
"result": toolResult,
|
|
}
|
|
if mcpEventData, err := json.Marshal(mcpEvent); err == nil {
|
|
fmt.Fprintf(c.Response().Writer, "data: %s\n\n", mcpEventData)
|
|
c.Response().Flush()
|
|
}
|
|
}
|
|
|
|
xlog.Debug("MCP streaming tools executed, re-running inference", "iteration", mcpStreamIter)
|
|
continue // next MCP stream iteration
|
|
}
|
|
}
|
|
|
|
// Automatic tool parsing fallback for streaming: when no tools were
|
|
// requested but the model emitted tool call markup, parse and emit them.
|
|
if !shouldUseFn && config.FunctionsConfig.AutomaticToolParsingFallback && collectedContent != "" && !toolsCalled {
|
|
parsed := functions.ParseFunctionCall(collectedContent, config.FunctionsConfig)
|
|
for i, fc := range parsed {
|
|
toolCallID := fc.ID
|
|
if toolCallID == "" {
|
|
toolCallID = id
|
|
}
|
|
toolCallMsg := schema.OpenAIResponse{
|
|
ID: id,
|
|
Created: created,
|
|
Model: input.Model,
|
|
Choices: []schema.Choice{{
|
|
Delta: &schema.Message{
|
|
Role: "assistant",
|
|
ToolCalls: []schema.ToolCall{{
|
|
Index: i,
|
|
ID: toolCallID,
|
|
Type: "function",
|
|
FunctionCall: schema.FunctionCall{
|
|
Name: fc.Name,
|
|
Arguments: fc.Arguments,
|
|
},
|
|
}},
|
|
},
|
|
Index: 0,
|
|
}},
|
|
Object: "chat.completion.chunk",
|
|
}
|
|
respData, _ := json.Marshal(toolCallMsg)
|
|
fmt.Fprintf(c.Response().Writer, "data: %s\n\n", respData)
|
|
c.Response().Flush()
|
|
toolsCalled = true
|
|
}
|
|
}
|
|
|
|
// No MCP tools to execute, send final stop message
|
|
finishReason := FinishReasonStop
|
|
if toolsCalled && len(input.Tools) > 0 {
|
|
finishReason = FinishReasonToolCalls
|
|
} else if toolsCalled {
|
|
finishReason = FinishReasonFunctionCall
|
|
} else if reachedTokenBudget(finalUsage.Completion, config.Maxtokens) {
|
|
// Generation stopped because it hit the max_tokens ceiling
|
|
// rather than a natural stop — report "length" (issue #9716).
|
|
finishReason = FinishReasonLength
|
|
}
|
|
|
|
// Final delta chunk: empty delta with finish_reason set. Per
|
|
// OpenAI streaming spec this chunk does NOT carry usage —
|
|
// the optional trailer (below) does, gated on include_usage.
|
|
resp := &schema.OpenAIResponse{
|
|
ID: id,
|
|
Created: created,
|
|
Model: input.Model, // we have to return what the user sent here, due to OpenAI spec.
|
|
Choices: []schema.Choice{
|
|
{
|
|
FinishReason: &finishReason,
|
|
Index: 0,
|
|
Delta: &schema.Message{},
|
|
},
|
|
},
|
|
Object: "chat.completion.chunk",
|
|
}
|
|
respData, _ := json.Marshal(resp)
|
|
|
|
middleware.StampUsage(c, input.Model, finalUsage.Prompt, finalUsage.Completion)
|
|
|
|
fmt.Fprintf(c.Response().Writer, "data: %s\n\n", respData)
|
|
|
|
// Trailing usage chunk per OpenAI spec: emit only when the
|
|
// caller opted in via stream_options.include_usage. Shape:
|
|
// {"choices":[],"usage":{...},"object":"chat.completion.chunk",...}
|
|
//
|
|
// finalUsage is the authoritative TokenUsage returned by the
|
|
// worker function (process / processTools) via the `ended`
|
|
// channel. The worker reads it from ComputeChoices' return
|
|
// value, which is the cumulative count produced by the backend
|
|
// over the whole prediction. Issue #9927 was caused by the
|
|
// tools-path worker not surfacing this value at all.
|
|
if input.StreamOptions != nil && input.StreamOptions.IncludeUsage {
|
|
trailerUsage := streamUsageFromTokenUsage(finalUsage, extraUsage)
|
|
trailerUsage.CompressionMeta = middleware.CompressionMetadata(c)
|
|
trailer := streamUsageTrailerJSON(id, input.Model, created, trailerUsage)
|
|
_, _ = fmt.Fprintf(c.Response().Writer, "data: %s\n\n", trailer)
|
|
}
|
|
|
|
fmt.Fprintf(c.Response().Writer, "data: [DONE]\n\n")
|
|
c.Response().Flush()
|
|
xlog.Debug("Stream ended")
|
|
return nil
|
|
} // end MCP stream iteration loop
|
|
|
|
// Safety fallback
|
|
fmt.Fprintf(c.Response().Writer, "data: [DONE]\n\n")
|
|
c.Response().Flush()
|
|
return nil
|
|
|
|
// no streaming mode
|
|
default:
|
|
mcpMaxIterations := 10
|
|
if config.Agent.MaxIterations > 0 {
|
|
mcpMaxIterations = config.Agent.MaxIterations
|
|
}
|
|
hasMCPTools := mcpExecutor != nil && mcpExecutor.HasTools()
|
|
|
|
for mcpIteration := 0; mcpIteration <= mcpMaxIterations; mcpIteration++ {
|
|
// Re-template on each MCP iteration since messages may have changed
|
|
if mcpIteration > 0 {
|
|
if err := middleware.CompressChatRequest(c, compressor); err != nil {
|
|
return err
|
|
}
|
|
if !config.TemplateConfig.UseTokenizerTemplate {
|
|
predInput = evaluator.TemplateMessages(*input, input.Messages, config, funcs, shouldUseFn)
|
|
xlog.Debug("MCP re-templating", "iteration", mcpIteration, "prompt_len", len(predInput))
|
|
}
|
|
}
|
|
|
|
// Detect if thinking token is already in prompt or template
|
|
var template string
|
|
if config.TemplateConfig.UseTokenizerTemplate {
|
|
template = config.GetModelTemplate() // TODO: this should be the parsed jinja template. But for now this is the best we can do.
|
|
} else {
|
|
template = predInput
|
|
}
|
|
thinkingStartToken := reason.DetectThinkingStartToken(template, &config.ReasoningConfig)
|
|
if config.TemplateConfig.UseTokenizerTemplate {
|
|
thinkingStartToken = reason.DetectThinkingStartTokenInTemplate(template, &config.ReasoningConfig)
|
|
}
|
|
|
|
xlog.Debug("Thinking start token", "thinkingStartToken", thinkingStartToken, "template", template)
|
|
|
|
// When shouldUseFn, the callback just stores the raw text — tool parsing
|
|
// is deferred to after ComputeChoices so we can check chat deltas first
|
|
// and avoid redundant Go-side parsing.
|
|
var cbRawResult, cbReasoning string
|
|
|
|
tokenCallback := func(s string, c *[]schema.Choice) {
|
|
reasoning, s := reason.ExtractReasoningWithConfig(s, thinkingStartToken, config.ReasoningConfig)
|
|
|
|
if !shouldUseFn {
|
|
stopReason := FinishReasonStop
|
|
message := &schema.Message{Role: "assistant", Content: &s}
|
|
if reasoning != "" {
|
|
message.Reasoning = &reasoning
|
|
}
|
|
*c = append(*c, schema.Choice{FinishReason: &stopReason, Index: 0, Message: message})
|
|
return
|
|
}
|
|
|
|
// Store raw text for deferred tool parsing
|
|
cbRawResult = s
|
|
cbReasoning = reasoning
|
|
}
|
|
|
|
var result []schema.Choice
|
|
var tokenUsage backend.TokenUsage
|
|
var err error
|
|
|
|
var chatDeltas []*pb.ChatDelta
|
|
result, tokenUsage, chatDeltas, err = ComputeChoices(
|
|
input,
|
|
predInput,
|
|
config,
|
|
cl,
|
|
startupOptions,
|
|
ml,
|
|
tokenCallback,
|
|
nil,
|
|
func(attempt int) bool {
|
|
if !shouldUseFn {
|
|
return false
|
|
}
|
|
// Retry when backend produced only reasoning and no content/tool calls.
|
|
// Full tool parsing is deferred until after ComputeChoices returns
|
|
// (when chat deltas are available), but we can detect the empty case here.
|
|
if cbRawResult == "" && textContentToReturn == "" {
|
|
xlog.Warn("Backend produced reasoning without actionable content, retrying",
|
|
"reasoning_len", len(cbReasoning), "attempt", attempt+1)
|
|
cbRawResult = ""
|
|
cbReasoning = ""
|
|
textContentToReturn = ""
|
|
return true
|
|
}
|
|
return false
|
|
},
|
|
)
|
|
if err != nil {
|
|
return backendRequestError(err)
|
|
}
|
|
|
|
// For non-tool requests: prefer C++ autoparser chat deltas over
|
|
// Go-side tag extraction (which can mangle output when thinkingStartToken
|
|
// differs from the model's actual reasoning tags, e.g. Gemma 4).
|
|
if !shouldUseFn {
|
|
result = applyAutoparserOverride(chatDeltas, thinkingStartToken, config.ReasoningConfig, result)
|
|
}
|
|
|
|
// Tool parsing is deferred here (only when shouldUseFn) so chat deltas are available
|
|
if shouldUseFn {
|
|
var funcResults []functions.FuncCallResults
|
|
|
|
// Try pre-parsed tool calls from C++ autoparser first
|
|
if deltaToolCalls := functions.ToolCallsFromChatDeltas(chatDeltas); len(deltaToolCalls) > 0 {
|
|
xlog.Debug("[ChatDeltas] non-SSE: using C++ autoparser tool calls, skipping Go-side parsing", "count", len(deltaToolCalls))
|
|
funcResults = deltaToolCalls
|
|
textContentToReturn = functions.ContentFromChatDeltas(chatDeltas)
|
|
cbReasoning = functions.ReasoningFromChatDeltas(chatDeltas)
|
|
} else if deltaContent := functions.ContentFromChatDeltas(chatDeltas); len(chatDeltas) > 0 && deltaContent != "" {
|
|
// ChatDeltas have content but no tool calls — model answered without using tools.
|
|
// This happens with thinking models (e.g. Gemma 4) where the Go-side reasoning
|
|
// extraction misclassifies clean content as reasoning, leaving cbRawResult empty.
|
|
xlog.Debug("[ChatDeltas] non-SSE: using C++ autoparser content (no tool calls)", "content_len", len(deltaContent))
|
|
textContentToReturn = deltaContent
|
|
cbReasoning = functions.ReasoningFromChatDeltas(chatDeltas)
|
|
} else {
|
|
// Fallback: parse tool calls from raw text
|
|
xlog.Debug("[ChatDeltas] non-SSE: no chat deltas, falling back to Go-side text parsing")
|
|
textContentToReturn = functions.ParseTextContent(cbRawResult, config.FunctionsConfig)
|
|
cbRawResult = functions.CleanupLLMResult(cbRawResult, config.FunctionsConfig)
|
|
funcResults = functions.ParseFunctionCall(cbRawResult, config.FunctionsConfig)
|
|
}
|
|
|
|
// Content-based tool call fallback: if no tool calls were found,
|
|
// try parsing the raw result — ParseFunctionCall handles detection internally.
|
|
if len(funcResults) == 0 {
|
|
contentFuncResults := functions.ParseFunctionCall(cbRawResult, config.FunctionsConfig)
|
|
if len(contentFuncResults) > 0 {
|
|
funcResults = contentFuncResults
|
|
textContentToReturn = functions.StripToolCallMarkup(cbRawResult)
|
|
}
|
|
}
|
|
|
|
noActionsToRun := len(funcResults) > 0 && funcResults[0].Name == noActionName || len(funcResults) == 0
|
|
|
|
switch {
|
|
case noActionsToRun:
|
|
// Use textContentToReturn if available (e.g. from ChatDeltas),
|
|
// otherwise fall back to cbRawResult for legacy Go-side parsing.
|
|
questionInput := cbRawResult
|
|
if textContentToReturn != "" {
|
|
questionInput = textContentToReturn
|
|
}
|
|
qResult, qErr := handleQuestion(config, funcResults, questionInput, predInput)
|
|
if qErr != nil {
|
|
xlog.Error("error handling question", "error", qErr)
|
|
}
|
|
|
|
stopReason := FinishReasonStop
|
|
message := &schema.Message{Role: "assistant", Content: &qResult}
|
|
if cbReasoning != "" {
|
|
message.Reasoning = &cbReasoning
|
|
}
|
|
result = append(result, schema.Choice{
|
|
FinishReason: &stopReason,
|
|
Message: message,
|
|
})
|
|
default:
|
|
toolCallsReason := FinishReasonToolCalls
|
|
toolChoice := schema.Choice{
|
|
FinishReason: &toolCallsReason,
|
|
Message: &schema.Message{
|
|
Role: "assistant",
|
|
},
|
|
}
|
|
if cbReasoning != "" {
|
|
toolChoice.Message.Reasoning = &cbReasoning
|
|
}
|
|
|
|
for _, ss := range funcResults {
|
|
name, args := ss.Name, ss.Arguments
|
|
toolCallID := ss.ID
|
|
if toolCallID == "" {
|
|
toolCallID = id
|
|
}
|
|
if len(input.Tools) > 0 {
|
|
toolChoice.Message.Content = textContentToReturn
|
|
toolChoice.Message.ToolCalls = append(toolChoice.Message.ToolCalls,
|
|
schema.ToolCall{
|
|
ID: toolCallID,
|
|
Type: "function",
|
|
FunctionCall: schema.FunctionCall{
|
|
Name: name,
|
|
Arguments: args,
|
|
},
|
|
},
|
|
)
|
|
} else {
|
|
// Deprecated function_call format
|
|
functionCallReason := FinishReasonFunctionCall
|
|
message := &schema.Message{
|
|
Role: "assistant",
|
|
Content: &textContentToReturn,
|
|
FunctionCall: map[string]any{
|
|
"name": name,
|
|
"arguments": args,
|
|
},
|
|
}
|
|
if cbReasoning != "" {
|
|
message.Reasoning = &cbReasoning
|
|
}
|
|
result = append(result, schema.Choice{
|
|
FinishReason: &functionCallReason,
|
|
Message: message,
|
|
})
|
|
}
|
|
}
|
|
|
|
if len(input.Tools) > 0 {
|
|
result = append(result, toolChoice)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Automatic tool parsing fallback: when no tools/functions were in the
|
|
// request but the model emitted tool call markup, parse and surface them.
|
|
if !shouldUseFn && config.FunctionsConfig.AutomaticToolParsingFallback && len(result) > 0 {
|
|
for i, choice := range result {
|
|
if choice.Message == nil || choice.Message.Content == nil {
|
|
continue
|
|
}
|
|
contentStr, ok := choice.Message.Content.(string)
|
|
if !ok || contentStr == "" {
|
|
continue
|
|
}
|
|
parsed := functions.ParseFunctionCall(contentStr, config.FunctionsConfig)
|
|
if len(parsed) == 0 {
|
|
continue
|
|
}
|
|
stripped := functions.StripToolCallMarkup(contentStr)
|
|
toolCallsReason := FinishReasonToolCalls
|
|
result[i].FinishReason = &toolCallsReason
|
|
if stripped != "" {
|
|
result[i].Message.Content = &stripped
|
|
} else {
|
|
result[i].Message.Content = nil
|
|
}
|
|
for _, fc := range parsed {
|
|
toolCallID := fc.ID
|
|
if toolCallID == "" {
|
|
toolCallID = id
|
|
}
|
|
result[i].Message.ToolCalls = append(result[i].Message.ToolCalls,
|
|
schema.ToolCall{
|
|
ID: toolCallID,
|
|
Type: "function",
|
|
FunctionCall: schema.FunctionCall{
|
|
Name: fc.Name,
|
|
Arguments: fc.Arguments,
|
|
},
|
|
},
|
|
)
|
|
}
|
|
}
|
|
}
|
|
|
|
// MCP server-side tool execution loop:
|
|
// If we have MCP tools and the model returned tool_calls, execute MCP tools
|
|
// and re-run inference with the results appended to the conversation.
|
|
if hasMCPTools && len(result) > 0 {
|
|
var mcpCallsExecuted bool
|
|
for _, choice := range result {
|
|
if choice.Message == nil || len(choice.Message.ToolCalls) == 0 {
|
|
continue
|
|
}
|
|
// Check if any tool calls are MCP tools
|
|
var hasMCPCalls bool
|
|
for _, tc := range choice.Message.ToolCalls {
|
|
if mcpExecutor != nil && mcpExecutor.IsTool(tc.FunctionCall.Name) {
|
|
hasMCPCalls = true
|
|
break
|
|
}
|
|
}
|
|
if !hasMCPCalls {
|
|
continue
|
|
}
|
|
|
|
// Append assistant message with tool_calls to conversation
|
|
assistantContent := ""
|
|
if choice.Message.Content != nil {
|
|
if s, ok := choice.Message.Content.(string); ok {
|
|
assistantContent = s
|
|
} else if sp, ok := choice.Message.Content.(*string); ok && sp != nil {
|
|
assistantContent = *sp
|
|
}
|
|
}
|
|
assistantMsg := schema.Message{
|
|
Role: "assistant",
|
|
Content: assistantContent,
|
|
ToolCalls: choice.Message.ToolCalls,
|
|
}
|
|
input.Messages = append(input.Messages, assistantMsg)
|
|
|
|
// Execute each MCP tool call and append results
|
|
for _, tc := range choice.Message.ToolCalls {
|
|
if mcpExecutor == nil || !mcpExecutor.IsTool(tc.FunctionCall.Name) {
|
|
continue
|
|
}
|
|
xlog.Debug("Executing MCP tool", "tool", tc.FunctionCall.Name, "arguments", tc.FunctionCall.Arguments, "iteration", mcpIteration)
|
|
toolResult, toolErr := mcpExecutor.ExecuteTool(c.Request().Context(), tc.FunctionCall.Name, tc.FunctionCall.Arguments)
|
|
if toolErr != nil {
|
|
xlog.Error("MCP tool execution failed", "tool", tc.FunctionCall.Name, "error", toolErr)
|
|
toolResult = fmt.Sprintf("Error: %v", toolErr)
|
|
}
|
|
input.Messages = append(input.Messages, schema.Message{
|
|
Role: "tool",
|
|
Content: toolResult,
|
|
StringContent: toolResult,
|
|
ToolCallID: tc.ID,
|
|
Name: tc.FunctionCall.Name,
|
|
})
|
|
mcpCallsExecuted = true
|
|
}
|
|
}
|
|
|
|
if mcpCallsExecuted {
|
|
xlog.Debug("MCP tools executed, re-running inference", "iteration", mcpIteration, "messages", len(input.Messages))
|
|
continue // next MCP iteration
|
|
}
|
|
}
|
|
|
|
// If generation hit the max_tokens ceiling, report "length"
|
|
// instead of a natural "stop" (issue #9716). Mirrors the
|
|
// streaming path; tool/function finish reasons are untouched.
|
|
if reachedTokenBudget(tokenUsage.Completion, config.Maxtokens) {
|
|
for i := range result {
|
|
if result[i].FinishReason != nil && *result[i].FinishReason == FinishReasonStop {
|
|
lengthReason := FinishReasonLength
|
|
result[i].FinishReason = &lengthReason
|
|
}
|
|
}
|
|
}
|
|
|
|
// No MCP tools to execute (or no MCP tools configured), return response
|
|
usage := schema.OpenAIUsage{
|
|
PromptTokens: tokenUsage.Prompt,
|
|
CompletionTokens: tokenUsage.Completion,
|
|
TotalTokens: tokenUsage.Prompt + tokenUsage.Completion,
|
|
}
|
|
if extraUsage {
|
|
usage.TimingTokenGeneration = tokenUsage.TimingTokenGeneration
|
|
usage.TimingPromptProcessing = tokenUsage.TimingPromptProcessing
|
|
}
|
|
usage.CompressionMeta = middleware.CompressionMetadata(c)
|
|
|
|
resp := &schema.OpenAIResponse{
|
|
ID: id,
|
|
Created: created,
|
|
Model: input.Model, // we have to return what the user sent here, due to OpenAI spec.
|
|
Choices: result,
|
|
Object: "chat.completion",
|
|
Usage: &usage,
|
|
}
|
|
respData, _ := json.Marshal(resp)
|
|
xlog.Debug("Response", "response", string(respData))
|
|
|
|
middleware.StampUsage(c, input.Model, usage.PromptTokens, usage.CompletionTokens)
|
|
|
|
// Return the prediction in the response body
|
|
return c.JSON(200, resp)
|
|
} // end MCP iteration loop
|
|
|
|
// Should not reach here, but safety fallback
|
|
return fmt.Errorf("MCP iteration limit reached")
|
|
}
|
|
}
|
|
}
|
|
|
|
func handleQuestion(config *config.ModelConfig, funcResults []functions.FuncCallResults, result, prompt string) (string, error) {
|
|
if len(funcResults) == 0 && result != "" {
|
|
xlog.Debug("nothing function results but we had a message from the LLM")
|
|
|
|
return result, nil
|
|
}
|
|
|
|
xlog.Debug("nothing to do, computing a reply")
|
|
arg := ""
|
|
if len(funcResults) > 0 {
|
|
arg = funcResults[0].Arguments
|
|
}
|
|
// If there is a message that the LLM already sends as part of the JSON reply, use it
|
|
arguments := map[string]any{}
|
|
if err := json.Unmarshal([]byte(arg), &arguments); err != nil {
|
|
xlog.Debug("handleQuestion: function result did not contain a valid JSON object")
|
|
}
|
|
m, exists := arguments["message"]
|
|
if exists {
|
|
switch message := m.(type) {
|
|
case string:
|
|
if message != "" {
|
|
xlog.Debug("Reply received from LLM", "message", message)
|
|
message = backend.Finetune(*config, prompt, message)
|
|
xlog.Debug("Reply received from LLM(finetuned)", "message", message)
|
|
|
|
return message, nil
|
|
}
|
|
}
|
|
}
|
|
|
|
xlog.Debug("No action received from LLM, without a message, computing a reply")
|
|
|
|
return "", nil
|
|
}
|
|
|
|
// forwardCloudProxyOpenAIViaBackend marshals the OpenAI request and
|
|
// hands off to the cloud-proxy gRPC backend which does the outbound
|
|
// HTTP. The chat endpoint owns the body construction because it's the
|
|
// only place the request lands as a parsed *schema.OpenAIRequest.
|
|
// Request-side PII redaction already ran in the middleware; the
|
|
// response is forwarded unmodified.
|
|
func forwardCloudProxyOpenAIViaBackend(c echo.Context, cfg *config.ModelConfig, input *schema.OpenAIRequest, ml *model.ModelLoader, appConfig *config.ApplicationConfig) error {
|
|
body, err := json.Marshal(input)
|
|
if err != nil {
|
|
return echo.NewHTTPError(http.StatusBadRequest, "cloudproxy: marshal request: "+err.Error())
|
|
}
|
|
return cloudproxy.ForwardViaBackend(c, cfg, body, ml, appConfig)
|
|
}
|