mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-08-02 05:47:31 +08:00
## Summary
Adds a side-effect-free DataFlow canvas **debug (dry-run) mode** plus a
**debug run log with a "View result" panel**, so a canvas can be
executed synchronously and inspected end-to-end (per-component progress
and parsed chunks) without persisting anything.
### Dry-run execution (inline parsed chunks)
- `task/debug.go`: `NewDebugTaskContext` builds an in-memory
`TaskContext` with **`KB.ID == ""` — the single debug signal used across
the ingestion pipeline**. A canvas debug run has no knowledgebase, and
production ingestion always supplies one, so `kb_id == ""` occurs ONLY
in debug mode. Components gate their own side effects on this signal
without any dedicated debug vocabulary (the former `CANVAS_DEBUG_DOC_ID`
marker constant is removed).
- `task/pipeline_executor.go`: `validateTaskContext` no longer requires
a KB when `KB.ID == ""` (debug); debug runs return `collectDebugOutput`
(chunks) instead of a no-op; uploaded bytes are delivered as
`inputs['binary']` for doc-less runs; `injectDebugPageCap` caps the
parser to the first pages for a fast preview via the production
`override_params` channel (Parser cpnID + family).
- `component/tokenizer.go`: `shouldHaveEmbedding` skips embedding when
`kb_id == ""` — the embedder is configured on the knowledgebase, so a
debug run has nothing to resolve against and stays side-effect free.
- `chunker/register.go`: chunk images are uploaded to MinIO only when a
KB is present (persist run). **In debug mode the raw image bytes are
intentionally dropped (`delete(ck, "image")`)** — the debug preview does
not render chunk images, and dropping the bytes keeps them out of memory
and out of the Redis-stored debug log. This is a deliberate trade-off,
not an oversight.
- `component/file.go`: pass through in-memory binary bytes, skipping
`doc_id` -> storage resolution.
- `handler/agent.go` + `agent_webhook.go`: detect `dataflow_canvas` and
run a sync debug returning chunks inline on the existing
chat/completions endpoint; reject DataFlow canvases from webhooks (fixes
the previously dead `== "DataFlow"` check; mirrors Python
`agent_api.py`).
- `parser_dispatch.go`: export `ParserFileFamily` for the executor's
page-cap injection.
### Debug run log + "View result"
Mirrors Python's debug-log contract so the front-end can replay each
component's progress and parsed output:
- `task/debug_log_sink.go`: a `DebugLogSink` records every component's
lifecycle into a `[{component_id, trace}]` array (each trace entry
carries `message`, `progress`, `timestamp`, `elapsed_time`). `Flush`
appends a terminal `END` marker whose first trace message is non-empty
so the front-end detects completion. On failure the END marker is
prefixed `[ERROR]` yet still carries the run, so the failure timeline
renders instead of being stuck empty. Timestamps and `elapsed_time` are
in seconds (matching the rest of the app).
- `task/debug_result_dsl.go`: `BuildDebugResultDSL` builds the `dsl` the
END marker carries — the Go analogue of Python's `Graph.__str__` +
END-marker `dsl` in `rag/flow/pipeline.py`. It combines the static DSL
structure (component_name / downstream / params / graph.nodes) with the
run output map (`output["state"][<id>]`) to emit, per component,
`obj.params.outputs[<format>].value` (chunks / text / json / html /
markdown) — the exact keys the front-end `dataflow-result` page reads to
render each step's parsed chunks. Raw embedding vectors (including the
dimension-scoped `q_<dim>_vec` keys) are stripped so the stored log
stays Python-scale.
- `task/pipeline_executor.go`: after the run, attach the built `dsl` to
the END marker via the `ResultSink` capability.
- `handler/agent.go`: `runCanvasPipelineDebug` generates a stable
`message_id` up-front and always flushes the log (success or failure);
`respondWithDebugResult` returns `message_id` in **both** the success
and the error envelope so the front-end can poll the log. The debug-log
endpoint `GET /agents/:id/logs/:message_id` serves the array.
- `web/src/pages/agent/hooks/use-run-dataflow.ts`: on a run failure,
also surface `message_id` via `setMessageId` so the log sheet renders
the failure timeline (the `[ERROR]` END marker is already written).
Guarded by `if (msgId)`, so it is a safe no-op when the back-end does
not return an id.
## Behavioral notes
- Debug parses only the first pages (`debugPageCapPages`) for a fast
preview; an explicit `pages` cap already present in the ParserConfig is
respected.
- Debug mode does not keep chunk images (see above) and does not compute
embeddings — it exercises parse + chunk only.
## Test plan
- Go: `debug_test.go`, `debug_log_sink_test.go` (trace pairing, END
marker, `[ERROR]` prefix, fractional-second timestamp/elapsed_time, size
caps, and a real-pipeline test asserting the END-marker `dsl` carries
non-empty per-component `params.outputs` with chunks),
`debug_result_dsl_test.go` (flat and real nested `output["state"]`
shapes, vector stripping, format priority),
`debug_pages_integration_test.go`, `pipeline_executor_persist_test.go`,
`handler/agent_pipeline_debug_test.go`, `handler/agent_logs_test.go`
(incl. `TestRunCanvasPipelineDebug_ErrorStillExposesMessageID` /
`TestRespondWithDebugResult_ErrorCarriesMessageID` locking `message_id`
on failure), plus updates to `agent_test.go` / `agent_webhook_test.go` /
`chunker/image_upload_test.go` / `tokenizer*.go`.
- `go build ./...` and `./build.sh --test` for affected packages.
🤖 Generated with [CodeBuddy Code](https://cnb.cool/codebuddy)
---------
Co-authored-by: CodeBuddy Code <noreply@tencent.com>
1683 lines
60 KiB
Go
1683 lines
60 KiB
Go
//
|
|
// Copyright 2026 The InfiniFlow Authors. All Rights Reserved.
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
//
|
|
|
|
package handler
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"mime/multipart"
|
|
"ragflow/internal/engine/redis"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/gin-gonic/gin"
|
|
"github.com/google/uuid"
|
|
"go.uber.org/zap"
|
|
|
|
"ragflow/internal/agent/canvas"
|
|
"ragflow/internal/common"
|
|
"ragflow/internal/dao"
|
|
"ragflow/internal/entity"
|
|
pipelinepkg "ragflow/internal/ingestion/pipeline"
|
|
"ragflow/internal/ingestion/task"
|
|
"ragflow/internal/service"
|
|
"ragflow/internal/service/file"
|
|
"ragflow/internal/utility"
|
|
|
|
dslpkg "ragflow/internal/agent/dsl"
|
|
)
|
|
|
|
// AgentHandler agent handler
|
|
// fileUploader is the subset of FileService used by agent handlers.
|
|
//
|
|
// The full FileService also has UploadFile, but it is consumed by
|
|
// the FileHandler (handler/file.go), not by any agent handler, so
|
|
// the interface deliberately does NOT list it. (Code review CR1.)
|
|
type agentFileService interface {
|
|
DownloadAgentFile(ctx context.Context, tenantID, location string) ([]byte, error)
|
|
// UploadInfos stores raw bytes in the per-user downloads bucket and
|
|
// returns lightweight descriptors.
|
|
UploadInfos(ctx context.Context, userID string, files []*multipart.FileHeader) ([]map[string]interface{}, error)
|
|
// UploadFromURL downloads a remote file (with SSRF protection) and
|
|
// stores it as an info blob.
|
|
UploadFromURL(ctx context.Context, tenantID, rawURL string) (map[string]interface{}, error)
|
|
}
|
|
|
|
// chatAgentService is the subset of AgentService used by the chat-completion
|
|
// endpoints (AgentChatCompletions, RunAgent). Kept as a separate interface so
|
|
// handler tests can inject a fake RunAgent without standing up the full
|
|
// AgentService (DB DAOs, eino runner, etc.). The production wiring in
|
|
// NewAgentHandler assigns the concrete *service.AgentService — which
|
|
// satisfies this interface because its RunAgent signature matches.
|
|
type chatAgentService interface {
|
|
RunAgent(ctx context.Context, userID, canvasID, sessionID, version string, userInput any, files []map[string]interface{}) (<-chan canvas.RunEvent, error)
|
|
}
|
|
|
|
// documentAccessChecker is the minimal surface RerunAgent needs
|
|
// from DocumentService. Defined as an interface (instead of taking
|
|
// the concrete *service.DocumentService) so handler tests can
|
|
// inject a deny-all stub without spinning up the full service
|
|
// (DB DAOs, storage clients, …). The production *service.DocumentService
|
|
// satisfies this interface because its Accessible signature
|
|
// matches.
|
|
type documentAccessChecker interface {
|
|
Accessible(docID, userID string) bool
|
|
}
|
|
|
|
// AgentHandler agent handler
|
|
type AgentHandler struct {
|
|
agentService *service.AgentService
|
|
chatRunner chatAgentService
|
|
fileService agentFileService
|
|
loader canvasLoader
|
|
// documentService is optional. Wired in cmd/server_main.go after
|
|
// NewAgentHandler (which doesn't take it to preserve the existing
|
|
// test-friendly signature). When nil, RerunAgent falls back to
|
|
// tenant-only authorization (i.e. cannot verify the doc, so the
|
|
// check is skipped — same shape as the pre-port behaviour).
|
|
documentService documentAccessChecker
|
|
// redisGet fetches a raw string from Redis. Defaults to the global
|
|
// client (redis.Get) so production behaviour is unchanged; tests inject
|
|
// a miniredis-backed getter to exercise GetAgentLogs without a live
|
|
// Redis (mirrors the newExecutor injection pattern).
|
|
redisGet func(key string) (string, error)
|
|
// redisStore writes the debug-run log array. Defaults to the global
|
|
// client (redis.Get, which satisfies task.DebugLogStore); tests inject a
|
|
// miniredis-backed writer so runCanvasPipelineDebug can be exercised without a
|
|
// live Redis.
|
|
redisStore task.DebugLogStore
|
|
// newExecutor builds the pipeline executor for a debug run. Defaults to
|
|
// task.NewPipelineExecutor; tests inject a fake so the debug path can be
|
|
// driven without a real canvas/DSL.
|
|
newExecutor func(taskCtx *task.TaskContext, canvasID string, docBulkSize int) (debugExecutor, error)
|
|
}
|
|
|
|
// debugExecutor is the minimal executor surface runCanvasPipelineDebug needs. It is an
|
|
// interface (not *task.PipelineExecutor) so tests can substitute a fake that
|
|
// emits progress through the attached sink. *task.PipelineExecutor satisfies it.
|
|
type debugExecutor interface {
|
|
WithProgressSink(sink pipelinepkg.ProgressSink) *task.PipelineExecutor
|
|
Execute(ctx context.Context) (*task.PipelineResult, error)
|
|
}
|
|
|
|
// WithDocumentService injects the document service used by
|
|
// RerunAgent to enforce DocumentService.accessible(docID, tenantID)
|
|
// before re-running. Returns the receiver for chaining in
|
|
// server_main wiring.
|
|
func (h *AgentHandler) WithDocumentService(s documentAccessChecker) *AgentHandler {
|
|
h.documentService = s
|
|
return h
|
|
}
|
|
|
|
// NewAgentHandler create agent handler
|
|
|
|
func NewAgentHandler(agentService *service.AgentService, fileService *file.FileService) *AgentHandler {
|
|
return &AgentHandler{
|
|
agentService: agentService,
|
|
chatRunner: agentService,
|
|
fileService: fileService,
|
|
loader: agentService,
|
|
redisGet: func(key string) (string, error) { return redis.Get().Get(key) },
|
|
redisStore: redis.Get(),
|
|
newExecutor: func(taskCtx *task.TaskContext, canvasID string, docBulkSize int) (debugExecutor, error) {
|
|
return task.NewPipelineExecutor(taskCtx, canvasID, docBulkSize)
|
|
},
|
|
}
|
|
}
|
|
|
|
// WithRedisGetter overrides the Redis string getter used by GetAgentLogs.
|
|
// Tests pass a miniredis-backed getter so the debug-log polling endpoint can
|
|
// be exercised end-to-end without a live Redis.
|
|
func (h *AgentHandler) WithRedisGetter(f func(key string) (string, error)) *AgentHandler {
|
|
h.redisGet = f
|
|
return h
|
|
}
|
|
|
|
// WithRedisStore overrides the Redis writer used by runCanvasPipelineDebug to persist
|
|
// the debug-run log. Tests pass a miniredis-backed writer.
|
|
func (h *AgentHandler) WithRedisStore(s task.DebugLogStore) *AgentHandler {
|
|
h.redisStore = s
|
|
return h
|
|
}
|
|
|
|
// WithNewExecutor overrides the pipeline executor factory used by
|
|
// runCanvasPipelineDebug. Tests inject a fake that emits progress through the
|
|
// attached sink without a real canvas/DSL.
|
|
func (h *AgentHandler) WithNewExecutor(f func(taskCtx *task.TaskContext, canvasID string, docBulkSize int) (debugExecutor, error)) *AgentHandler {
|
|
h.newExecutor = f
|
|
return h
|
|
}
|
|
|
|
// ListAgents lists agent canvases for the current user.
|
|
// @Summary List Agents
|
|
// @Description List agent canvases accessible to the current user (Home dashboard tile)
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param keywords query string false "Filter by title keyword"
|
|
// @Param page query int false "Page number (0 = no pagination)"
|
|
// @Param page_size query int false "Items per page (0 = no pagination)"
|
|
// @Param orderby query string false "Order-by field (default: create_time)"
|
|
// @Param desc query bool false "Descending order (default: true)"
|
|
// @Param owner_ids query string false "Comma-separated owner IDs to filter (default: all authorised tenants)"
|
|
// @Param canvas_category query string false "Canvas category (default: agent_canvas)"
|
|
// @Success 200 {object} service.ListAgentsResponse
|
|
// @Router /api/v1/agents [get]
|
|
func (h *AgentHandler) ListAgents(c *gin.Context) {
|
|
user, errorCode, errorMessage := GetUser(c)
|
|
if errorCode != common.CodeSuccess {
|
|
common.ErrorWithCode(c, errorCode, errorMessage)
|
|
return
|
|
}
|
|
|
|
keywords := c.Query("keywords")
|
|
canvasCategory := c.Query("canvas_category")
|
|
canvasType := c.Query("canvas_type")
|
|
|
|
page := 0
|
|
if v := c.Query("page"); v != "" {
|
|
if p, err := strconv.Atoi(v); err == nil && p > 0 {
|
|
page = p
|
|
}
|
|
}
|
|
|
|
pageSize := 0
|
|
if v := c.Query("page_size"); v != "" {
|
|
if ps, err := strconv.Atoi(v); err == nil && ps > 0 {
|
|
pageSize = ps
|
|
}
|
|
}
|
|
|
|
orderby := c.DefaultQuery("orderby", "create_time")
|
|
|
|
desc := true
|
|
if v := c.Query("desc"); v != "" {
|
|
desc = strings.ToLower(v) != "false"
|
|
}
|
|
|
|
var ownerIDs []string
|
|
if raw := c.Query("owner_ids"); raw != "" {
|
|
for _, id := range strings.Split(raw, ",") {
|
|
id = strings.TrimSpace(id)
|
|
if id != "" {
|
|
ownerIDs = append(ownerIDs, id)
|
|
}
|
|
}
|
|
}
|
|
var tags []string
|
|
if raw := c.Query("tags"); raw != "" {
|
|
for _, tag := range strings.Split(raw, ",") {
|
|
tag = strings.TrimSpace(tag)
|
|
if tag != "" {
|
|
tags = append(tags, tag)
|
|
}
|
|
}
|
|
}
|
|
ctx := c.Request.Context()
|
|
|
|
result, code, err := h.agentService.ListAgents(
|
|
ctx,
|
|
user.ID,
|
|
keywords,
|
|
page,
|
|
pageSize,
|
|
orderby,
|
|
desc,
|
|
ownerIDs,
|
|
canvasCategory,
|
|
canvasType,
|
|
tags,
|
|
)
|
|
if err != nil {
|
|
common.ResponseWithCodeData(c, code, false, err.Error())
|
|
return
|
|
}
|
|
|
|
common.SuccessWithData(c, result, "success")
|
|
}
|
|
|
|
// mapAgentError normalises service-layer errors onto the existing
|
|
// {code, data, message} response envelope used by every other handler.
|
|
//
|
|
// Four classes:
|
|
// - service.ErrAgentNotOwner -> "Only the owner..." (DELETE only, 103)
|
|
// - service.ErrAgentSessionBusy -> "session already running" (103)
|
|
// - dao.ErrUserCanvasNotFound -> "Make sure you have permission..." (103)
|
|
// - service.ErrAgentStorageError -> "Internal storage error" (500)
|
|
//
|
|
// The first two surface as OPERATING_ERROR(103) so the front-end
|
|
// cannot use the response code to enumerate other users' canvases
|
|
// (IDOR mitigation). ErrAgentStorageError maps to SERVER_ERROR(500)
|
|
// with a sanitized message — the raw DAO / DB error string MUST
|
|
// NOT reach the client (DSNs, table names, gorm stack frames would
|
|
// otherwise leak). Everything else becomes CodeDataError so the
|
|
// front-end can surface the message verbatim.
|
|
func mapAgentError(err error) (common.ErrorCode, string) {
|
|
if err == nil {
|
|
return common.CodeSuccess, ""
|
|
}
|
|
if errors.Is(err, service.ErrAgentNotOwner) {
|
|
return common.CodeOperatingError, "Only the owner of the agent is authorized for this operation."
|
|
}
|
|
if errors.Is(err, service.ErrAgentSessionBusy) {
|
|
return common.CodeOperatingError, "This agent session is already running."
|
|
}
|
|
if errors.Is(err, dao.ErrUserCanvasNotFound) ||
|
|
errors.Is(err, dao.ErrUserCanvasVersionNotFound) {
|
|
return common.CodeOperatingError, "Make sure you have permission to access the agent."
|
|
}
|
|
if errors.Is(err, service.ErrAgentStorageError) {
|
|
return common.CodeServerError, "Internal storage error while accessing the agent."
|
|
}
|
|
return common.CodeDataError, err.Error()
|
|
}
|
|
|
|
// CreateAgent creates a new agent canvas.
|
|
// @Summary Create Agent
|
|
// @Tags agents
|
|
// @Accept json
|
|
// @Produce json
|
|
// @Param request body service.CreateAgentRequest true "agent create request"
|
|
// @Success 200 {object} entity.UserCanvas
|
|
// @Router /api/v1/agents [post]
|
|
func (h *AgentHandler) CreateAgent(c *gin.Context) {
|
|
user, errorCode, errorMessage := GetUser(c)
|
|
if errorCode != common.CodeSuccess {
|
|
common.ErrorWithCode(c, errorCode, errorMessage)
|
|
return
|
|
}
|
|
var req service.CreateAgentRequest
|
|
if err := c.ShouldBindJSON(&req); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
req.UserID = user.ID
|
|
row, code, err := h.agentService.CreateAgent(c.Request.Context(), &req)
|
|
if err != nil {
|
|
common.ErrorWithCode(c, code, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, row, "success")
|
|
}
|
|
|
|
// GetAgent returns one canvas by ID.
|
|
// @Summary Get Agent
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Success 200 {object} entity.UserCanvas
|
|
// @Router /api/v1/agents/{canvas_id} [get]
|
|
func (h *AgentHandler) GetAgent(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
row, err := h.agentService.GetAgent(c.Request.Context(), user.ID, canvasID)
|
|
if err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
// Defensive: any historical v1 / Go-v2-only row in user_canvas.dsl
|
|
// is rendered as a missing graph by the front-end. Normalize in
|
|
// place (NormalizeForCanvas is a no-op when graph.nodes is set) so
|
|
// the response is always renderable without a migration.
|
|
if row != nil {
|
|
row.DSL = dslpkg.NormalizeForCanvas(row.DSL)
|
|
}
|
|
// Attach last_publish_time so the front-end can render when the agent
|
|
// was last published (Python get_agent parity).
|
|
lastPublishTime, err := h.agentService.GetLastPublishTime(c.Request.Context(), canvasID)
|
|
if err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
common.SuccessWithData(c, &agentDetailResponse{UserCanvas: row, LastPublishTime: lastPublishTime}, "success")
|
|
}
|
|
|
|
// agentDetailResponse wraps the canvas row with derived fields that Python's
|
|
// get_agent handler adds to the response.
|
|
type agentDetailResponse struct {
|
|
*entity.UserCanvas
|
|
LastPublishTime *int64 `json:"last_publish_time,omitempty"`
|
|
}
|
|
|
|
// updateAgentRequest is the wire shape for PUT /api/v1/agents/:canvas_id.
|
|
type updateAgentRequest map[string]interface{}
|
|
|
|
// UpdateAgent applies a partial update to the canvas draft.
|
|
// @Summary Update Agent
|
|
// @Tags agents
|
|
// @Accept json
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Param request body updateAgentRequest true "agent update payload"
|
|
// @Success 200 {object} map[string]interface{}
|
|
// @Router /api/v1/agents/{canvas_id} [put]
|
|
func (h *AgentHandler) UpdateAgent(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
var req updateAgentRequest
|
|
if err := c.ShouldBindJSON(&req); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
if req == nil {
|
|
req = updateAgentRequest{}
|
|
}
|
|
if err := h.agentService.UpdateAgent(c.Request.Context(), user.ID, canvasID, map[string]interface{}(req)); err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
canvas, err := h.agentService.GetAgent(c.Request.Context(), user.ID, canvasID)
|
|
if err != nil || canvas == nil {
|
|
common.SuccessWithData(c, map[string]interface{}{}, "success")
|
|
return
|
|
}
|
|
common.SuccessWithData(c, map[string]interface{}{
|
|
"update_time": canvas.UpdateTime,
|
|
}, "success")
|
|
}
|
|
|
|
// DeleteAgent removes the canvas and cascades to its versions.
|
|
// @Summary Delete Agent
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Success 200 {object} map[string]interface{}
|
|
// @Router /api/v1/agents/{canvas_id} [delete]
|
|
func (h *AgentHandler) DeleteAgent(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
if err := h.agentService.DeleteAgent(c.Request.Context(), user.ID, canvasID); err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
common.SuccessWithData(c, true, "success")
|
|
}
|
|
|
|
// ListTemplates lists every canvas template available to authenticated users.
|
|
// @Summary List Agent Templates
|
|
// @Description List the catalogue of canvas templates that authenticated users can clone.
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Security ApiKeyAuth
|
|
// @Success 200 {object} map[string]interface{}
|
|
// @Router /api/v1/agents/templates [get]
|
|
func (h *AgentHandler) ListTemplates(c *gin.Context) {
|
|
if _, errorCode, errorMessage := GetUser(c); errorCode != common.CodeSuccess {
|
|
common.ErrorWithCode(c, errorCode, errorMessage)
|
|
return
|
|
}
|
|
ctx := c.Request.Context()
|
|
templates, err := h.agentService.ListTemplates(ctx)
|
|
if err != nil {
|
|
jsonInternalError(c, err)
|
|
return
|
|
}
|
|
if templates == nil {
|
|
// Ensure the JSON payload is always a list, never null.
|
|
templates = []*entity.CanvasTemplate{}
|
|
}
|
|
|
|
common.SuccessWithData(c, templates, "success")
|
|
}
|
|
|
|
// RunAgent returns an SSE stream of execution events. The Phase 5 stub emits
|
|
// a single "Phase 5 wiring pending" event and closes; the real eino run
|
|
// loop will replace the channel source in service.AgentService.RunAgent.
|
|
// @Summary Run Agent (SSE)
|
|
// @Tags agents
|
|
// @Produce text/event-stream
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Param version query string false "version id (default: latest)"
|
|
// @Success 200 {string} string "SSE: data: {...}\\n\\n"
|
|
// @Router /api/v1/agents/{canvas_id}/run [post]
|
|
func (h *AgentHandler) RunAgent(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
version := c.Query("version")
|
|
sessionID := c.Query("session_id")
|
|
if sessionID == "" {
|
|
// Allocate the ordinary-Agent session identity at the HTTP boundary.
|
|
// Persistence of the session record remains owned by AgentService.RunAgent.
|
|
sessionID = utility.GenerateToken()
|
|
}
|
|
userInput := readUserInput(c)
|
|
|
|
events, err := h.chatRunner.RunAgent(c.Request.Context(), user.ID, canvasID, sessionID, version, userInput, nil)
|
|
if err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
c.Writer.Header().Set("Content-Type", "text/event-stream")
|
|
c.Writer.Header().Set("Cache-Control", "no-cache")
|
|
c.Writer.Header().Set("Connection", "keep-alive")
|
|
doneSent := false
|
|
for ev := range events {
|
|
if ev.Type == "done" {
|
|
doneSent = true
|
|
}
|
|
if err := service.WriteChatbotRunEvent(c.Writer, ev); err != nil {
|
|
common.Debug("agent run: client disconnected",
|
|
zap.String("canvas_id", canvasID),
|
|
zap.String("session_id", sessionID),
|
|
zap.Error(err),
|
|
)
|
|
return
|
|
}
|
|
}
|
|
if !doneSent {
|
|
if err := service.WriteChatbotRunEvent(c.Writer, canvas.RunEvent{Type: "done"}); err != nil {
|
|
common.Debug("agent run: failed to write [DONE]",
|
|
zap.String("canvas_id", canvasID),
|
|
zap.String("session_id", sessionID),
|
|
zap.Error(err),
|
|
)
|
|
}
|
|
}
|
|
}
|
|
|
|
// readUserInput extracts the user_input field from the JSON body if
|
|
// present, otherwise from the ?user_input= query string. An empty body
|
|
// (no body sent) is treated as "" so the resume cycle still works
|
|
// when the client only passes ?session_id=...&user_input=... on the URL.
|
|
func readUserInput(c *gin.Context) string {
|
|
if c.Request.ContentLength > 0 {
|
|
var body struct {
|
|
UserInput string `json:"user_input"`
|
|
Query string `json:"query"`
|
|
Message string `json:"message"`
|
|
}
|
|
if err := c.ShouldBindJSON(&body); err == nil {
|
|
if body.UserInput != "" {
|
|
return body.UserInput
|
|
}
|
|
if body.Query != "" {
|
|
return body.Query
|
|
}
|
|
if body.Message != "" {
|
|
return body.Message
|
|
}
|
|
}
|
|
}
|
|
return c.Query("user_input")
|
|
}
|
|
|
|
// runCanvasPipelineDebug runs a canvas in pipeline dry-run (debug) mode: it builds
|
|
// an in-memory debug TaskContext and executes the pipeline synchronously,
|
|
// returning the produced chunks. The run is fully synchronous: the complete
|
|
// debug log (including the terminal END marker) is flushed to Redis BEFORE
|
|
// this method returns, so the message_id it yields is for *replay* —
|
|
// re-fetching the log after a page refresh or from a shared link — not for
|
|
// incremental streaming; the front-end's poll therefore receives a finished
|
|
// log on the first request.
|
|
//
|
|
// It performs no persistence (no MinIO upload, index insert, or pipeline_log)
|
|
// because the debug context carries no KB (KB.ID == ""), which is the single
|
|
// debug signal used across the ingestion pipeline: the tokenizer skips
|
|
// embedding and the executor skips the persist stage. kb_id == "" occurs ONLY
|
|
// in debug mode; production ingestion always supplies a KB, so this path is
|
|
// never taken in normal operation. Debug takes no kb_id at all: it never
|
|
// resolves an embedder or writes to a KB, and parse/chunk work without it.
|
|
// Callers extract the (fileName, fileData) pair from their own request shape
|
|
// and pass them in.
|
|
func (h *AgentHandler) runCanvasPipelineDebug(ctx context.Context, user *entity.User,
|
|
canvasID, fileName string, fileData []byte) (*task.PipelineResult, error) {
|
|
// messageID is the polling key the front-end reads from the run response
|
|
// (see web/src/pages/agent/hooks/use-run-dataflow.ts:40) to fetch the debug
|
|
// log. It is generated up-front so the response always carries a stable key.
|
|
messageID := uuid.New().String()
|
|
|
|
// Force an empty KB: a debug context has no knowledgebase, so kb_id == ""
|
|
// holds and the debug signal is unambiguous.
|
|
taskCtx := task.NewDebugTaskContext(user.ID, canvasID, fileName, fileData)
|
|
|
|
exec, err := h.newExecutor(taskCtx, canvasID, 0)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Attach a sink that records component progress into the
|
|
// [{component_id, trace}] array the front-end polls for.
|
|
sink := task.NewDebugLogSink(canvasID, messageID, h.redisStore)
|
|
exec.WithProgressSink(sink)
|
|
|
|
result, runErr := exec.Execute(ctx)
|
|
// Flush even on error so the front-end can poll the failure log; the END
|
|
// marker carries the error and keeps the run detectable as finished.
|
|
sink.Flush(ctx, runErr)
|
|
if result == nil {
|
|
// The executor may return (nil, err) on failure; synthesize a result so
|
|
// the polling key survives to the response and the front-end can still
|
|
// reach the failure log written just above.
|
|
result = &task.PipelineResult{}
|
|
}
|
|
result.MessageID = messageID
|
|
return result, runErr
|
|
}
|
|
|
|
// respondWithDebugResult writes the outcome of a dataflow debug run. On a run
|
|
// failure it keeps a non-success code (errors must not be masked as success)
|
|
// but still carries message_id in the error envelope's data, because
|
|
// runCanvasPipelineDebug already flushed the failure log to Redis — the front-end
|
|
// needs that key to poll the failure timeline. On success it returns the
|
|
// message_id (the polling key) together with the produced chunks. Empty chunk
|
|
// slices become `[]` rather than `null` so the response shape stays stable.
|
|
// This is the single wire-shape helper for the chat/completions DataFlow debug
|
|
// entry point, and it matches the contract the front-end expects:
|
|
// `data.data.message_id` drives the log-poll loop and `data.data.chunks`
|
|
// renders the debug output.
|
|
func respondWithDebugResult(c *gin.Context, result *task.PipelineResult, err error) {
|
|
if err != nil {
|
|
messageID := ""
|
|
if result != nil {
|
|
messageID = result.MessageID
|
|
}
|
|
common.ResponseWithCodeData(c, common.CodeServerError, gin.H{"message_id": messageID}, err.Error())
|
|
return
|
|
}
|
|
messageID := ""
|
|
var chunks []map[string]any
|
|
if result != nil {
|
|
messageID = result.MessageID
|
|
chunks = result.Chunks
|
|
}
|
|
if len(chunks) == 0 {
|
|
chunks = []map[string]any{}
|
|
}
|
|
common.SuccessWithData(c, gin.H{"message_id": messageID, "chunks": chunks}, "success")
|
|
}
|
|
|
|
// extractChatDebugFile pulls the first uploaded file from a chat/completions
|
|
// request into (fileName, bytes) for a dataflow debug run. Unlike a debug
|
|
// multipart upload, chat/completions `files` are storage references
|
|
// (id/name/created_by); the raw bytes are fetched from the per-user downloads
|
|
// bucket via the file service. When no usable file is present it returns a
|
|
// default name and a nil slice so the pipeline can still run file-less.
|
|
func (h *AgentHandler) extractChatDebugFile(ctx context.Context, files [][]map[string]interface{}, user *entity.User) (string, []byte) {
|
|
if len(files) == 0 {
|
|
return "debug", nil
|
|
}
|
|
fileList := files[0]
|
|
if len(fileList) == 0 {
|
|
return "debug", nil
|
|
}
|
|
fd := fileList[0]
|
|
name, _ := fd["name"].(string)
|
|
id, _ := fd["id"].(string)
|
|
if id == "" {
|
|
if name == "" {
|
|
return "debug", nil
|
|
}
|
|
return name, nil
|
|
}
|
|
data, err := h.fileService.DownloadAgentFile(ctx, user.ID, id)
|
|
if err != nil || len(data) == 0 {
|
|
if name == "" {
|
|
return "debug", nil
|
|
}
|
|
return name, nil
|
|
}
|
|
if name == "" {
|
|
name = "debug"
|
|
}
|
|
return name, data
|
|
}
|
|
|
|
// sanitiseRunEventError passes through the error event payload
|
|
// unchanged. The runner serialises canvas.ErrorEvent ({"message": ...})
|
|
// before push, so when the payload round-trips through JSON the
|
|
// message field is already preserved. Heuristic sanitisation is
|
|
// disabled until the runner tags error events with a "kind"
|
|
// field — without that, blanket rewriting every error to
|
|
// "Internal storage error while accessing the agent." hides the
|
|
// real failure from the front-end and the user (v3.6.1 diagnostic
|
|
// regression: every canvas run failure surfaced as the same opaque
|
|
// string).
|
|
func sanitiseRunEventError(data string) string {
|
|
if data == "" {
|
|
return `{"message":"Unknown agent runtime error"}`
|
|
}
|
|
return data
|
|
}
|
|
|
|
// CancelSessionRun cancels one ordinary Agent run by session id.
|
|
// @Summary Cancel Agent Session Run
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param session_id path string true "session id"
|
|
// @Success 200 {object} map[string]interface{}
|
|
// @Router /api/v1/tasks/{session_id}/cancel [post]
|
|
func (h *AgentHandler) CancelSessionRun(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
if h.agentService == nil {
|
|
common.ResponseWithCodeData(c, common.CodeServerError, nil, "agent service unavailable")
|
|
return
|
|
}
|
|
if err := h.agentService.CancelSessionRun(c.Request.Context(), user.ID, c.Param("session_id")); err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
common.SuccessWithData(c, true, "success")
|
|
}
|
|
|
|
// publishAgentRequest is the wire shape for POST /api/v1/agents/:canvas_id/publish.
|
|
type publishAgentRequest struct {
|
|
Title *string `json:"title,omitempty"`
|
|
Description *string `json:"description,omitempty"`
|
|
DSL entity.JSONMap `json:"dsl,omitempty"`
|
|
}
|
|
|
|
// PublishAgent creates a new immutable version row and marks the parent canvas as released.
|
|
// @Summary Publish Agent Version
|
|
// @Tags agents
|
|
// @Accept json
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Param request body publishAgentRequest true "publish payload"
|
|
// @Success 200 {object} entity.UserCanvasVersion
|
|
// @Router /api/v1/agents/{canvas_id}/publish [post]
|
|
func (h *AgentHandler) PublishAgent(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
var req publishAgentRequest
|
|
if err := c.ShouldBindJSON(&req); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
row, err := h.agentService.PublishAgent(c.Request.Context(), user.ID, canvasID, &service.PublishAgentRequest{
|
|
Title: req.Title,
|
|
Description: req.Description,
|
|
DSL: req.DSL,
|
|
})
|
|
if err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
if row != nil {
|
|
row.DSL = dslpkg.NormalizeForCanvas(row.DSL)
|
|
}
|
|
common.SuccessWithData(c, row, "success")
|
|
}
|
|
|
|
// ListVersions returns every version of a canvas, newest first.
|
|
// @Summary List Agent Versions
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Success 200 {array} entity.UserCanvasVersion
|
|
// @Router /api/v1/agents/{canvas_id}/versions [get]
|
|
func (h *AgentHandler) ListVersions(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
rows, err := h.agentService.ListVersions(c.Request.Context(), user.ID, canvasID)
|
|
if err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
if rows == nil {
|
|
rows = []*entity.UserCanvasVersion{}
|
|
}
|
|
for _, row := range rows {
|
|
if row == nil {
|
|
continue
|
|
}
|
|
row.DSL = dslpkg.NormalizeForCanvas(row.DSL)
|
|
}
|
|
common.SuccessWithData(c, rows, "success")
|
|
}
|
|
|
|
// GetVersion returns a single version.
|
|
// @Summary Get Agent Version
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Param version_id path string true "version id"
|
|
// @Success 200 {object} entity.UserCanvasVersion
|
|
// @Router /api/v1/agents/{canvas_id}/versions/{version_id} [get]
|
|
func (h *AgentHandler) GetVersion(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
versionID := c.Param("version_id")
|
|
row, err := h.agentService.GetVersion(c.Request.Context(), user.ID, canvasID, versionID)
|
|
if err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
if row != nil {
|
|
row.DSL = dslpkg.NormalizeForCanvas(row.DSL)
|
|
}
|
|
common.SuccessWithData(c, row, "success")
|
|
}
|
|
|
|
// DeleteVersion removes a single version by id.
|
|
// @Summary Delete Agent Version
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Param version_id path string true "version id"
|
|
// @Success 200 {object} map[string]interface{}
|
|
// @Router /api/v1/agents/{canvas_id}/versions/{version_id} [delete]
|
|
func (h *AgentHandler) DeleteVersion(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
versionID := c.Param("version_id")
|
|
if err := h.agentService.DeleteVersion(c.Request.Context(), user.ID, canvasID, versionID); err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
common.SuccessWithData(c, true, "success")
|
|
}
|
|
|
|
// --- PR2: missing routes wired up to the existing service layer ---
|
|
|
|
// ListAgentTemplates GET /api/v1/agents/templates
|
|
func (h *AgentHandler) ListAgentTemplates(c *gin.Context) {
|
|
if _, code, msg := GetUser(c); code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
ctx := c.Request.Context()
|
|
rows, err := h.agentService.ListTemplates(ctx)
|
|
if err != nil {
|
|
common.ResponseWithCodeData(c, common.CodeServerError, nil, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, rows, "success")
|
|
}
|
|
|
|
// Prompts GET /api/v1/agents/prompts — returns the four hardcoded
|
|
// authoring guidelines the agent UI surfaces. The Python agent API
|
|
// returns these from a module-level constant; we keep the same shape.
|
|
func (h *AgentHandler) Prompts(c *gin.Context) {
|
|
if _, code, msg := GetUser(c); code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
common.SuccessWithData(c, gin.H{
|
|
"task_analysis": "As an AI agent designer, your role is to engage users by understanding their objectives and creating effective agent designs. Begin by analyzing the user's request to determine the appropriate actions.",
|
|
"output_format": "For each agent you create, detail its components and explain how they collaborate to achieve the user's goal.",
|
|
"citation_guidelines": "If the agent uses external sources, cite them in the final output. Use the format: [index] document_id, which corresponds to the document identifier in the database.",
|
|
"few_shots_examples": "<example/>",
|
|
}, "success")
|
|
}
|
|
|
|
// ListAgentTags list agent tags.
|
|
func (h *AgentHandler) ListAgentTags(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
ctx := c.Request.Context()
|
|
rows, errCode, err := h.agentService.ListAgentTags(ctx, user.ID, strings.TrimSpace(c.Query("canvas_category")))
|
|
if err != nil {
|
|
common.ResponseWithCodeData(c, errCode, nil, err.Error())
|
|
return
|
|
}
|
|
|
|
common.SuccessWithData(c, rows, "success")
|
|
}
|
|
|
|
// UpdateAgentTags PUT /api/v1/agents/:canvas_id/tags
|
|
func (h *AgentHandler) UpdateAgentTags(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
var body struct {
|
|
Tags interface{} `json:"tags"`
|
|
}
|
|
if err := c.ShouldBindJSON(&body); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
ctx := c.Request.Context()
|
|
ok, errCode, errMsg := h.agentService.UpdateAgentTags(ctx, user.ID, canvasID, body.Tags)
|
|
if !ok {
|
|
common.ResponseWithCodeData(c, errCode, nil, errMsg.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, true, "success")
|
|
}
|
|
|
|
// ListAgentSessions GET /api/v1/agents/:canvas_id/sessions
|
|
func (h *AgentHandler) ListAgentSessions(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
page, _ := strconv.Atoi(c.DefaultQuery("page", "0"))
|
|
pageSize, _ := strconv.Atoi(c.DefaultQuery("page_size", "0"))
|
|
keywords := c.Query("keywords")
|
|
fromDate := c.Query("from_date")
|
|
toDate := c.Query("to_date")
|
|
orderby := c.DefaultQuery("orderby", "create_time")
|
|
desc := c.DefaultQuery("desc", "true") != "false"
|
|
sessionID := c.Query("id")
|
|
expUserID := c.Query("user_id")
|
|
includeDSL := c.Query("dsl") == "true"
|
|
ctx := c.Request.Context()
|
|
resp, code, err := h.agentService.ListAgentSessions(ctx, user.ID, user.ID, canvasID, service.ListAgentSessionsRequest{
|
|
Page: page,
|
|
PageSize: pageSize,
|
|
Keywords: keywords,
|
|
FromDate: fromDate,
|
|
ToDate: toDate,
|
|
OrderBy: orderby,
|
|
Desc: desc,
|
|
SessionID: sessionID,
|
|
UserID: user.ID,
|
|
ExpUserID: expUserID,
|
|
IncludeDSL: includeDSL,
|
|
})
|
|
if err != nil {
|
|
common.ErrorWithCode(c, code, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, resp.Data, "success")
|
|
}
|
|
|
|
// CreateAgentSession POST /api/v1/agents/:canvas_id/sessions
|
|
func (h *AgentHandler) CreateAgentSession(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
var body struct {
|
|
Name string `json:"name"`
|
|
Source string `json:"source"`
|
|
DSL json.RawMessage `json:"dsl"`
|
|
}
|
|
if err := c.ShouldBindJSON(&body); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
ctx := c.Request.Context()
|
|
row, code, err := h.agentService.CreateAgentSession(ctx, &service.CreateAgentSessionRequest{
|
|
UserID: user.ID,
|
|
AgentID: canvasID,
|
|
Name: body.Name,
|
|
Source: body.Source,
|
|
DSL: body.DSL,
|
|
})
|
|
if err != nil {
|
|
common.ErrorWithCode(c, code, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, row, "success")
|
|
}
|
|
|
|
// GetAgentSession GET /api/v1/agents/:canvas_id/sessions/:session_id
|
|
func (h *AgentHandler) GetAgentSession(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
sessionID := c.Param("session_id")
|
|
ctx := c.Request.Context()
|
|
row, code, err := h.agentService.GetAgentSession(ctx, user.ID, canvasID, sessionID)
|
|
if err != nil {
|
|
common.ErrorWithCode(c, code, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, row, "success")
|
|
}
|
|
|
|
// DeleteAgentSession DELETE /api/v1/agents/:canvas_id/sessions[/:session_id]
|
|
//
|
|
// Path parameter disambiguation:
|
|
// - /sessions/:session_id -> single item delete
|
|
// - /sessions?ids=a,b -> batch delete (delete_all when ids is empty)
|
|
func (h *AgentHandler) DeleteAgentSession(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
ctx := c.Request.Context()
|
|
canvasID := c.Param("canvas_id")
|
|
sessionID := c.Param("session_id")
|
|
if sessionID != "" {
|
|
ok, code, err := h.agentService.DeleteAgentSessionItem(ctx, user.ID, canvasID, sessionID)
|
|
if err != nil {
|
|
common.ErrorWithCode(c, code, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, ok, "success")
|
|
return
|
|
}
|
|
idsParam := c.Query("ids")
|
|
deleteAll := c.Query("delete_all") == "true"
|
|
var ids []string
|
|
if idsParam != "" {
|
|
for _, id := range strings.Split(idsParam, ",") {
|
|
if id = strings.TrimSpace(id); id != "" {
|
|
ids = append(ids, id)
|
|
}
|
|
}
|
|
}
|
|
result, code, err := h.agentService.DeleteAgentSessions(ctx, user.ID, canvasID, ids, deleteAll)
|
|
if err != nil {
|
|
common.ErrorWithCode(c, code, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, result, "success")
|
|
}
|
|
|
|
// AgentChatCompletions POST /api/v1/agents/chat/completions
|
|
//
|
|
// Runs the canvas against `agent_id`. Mirrors the Python contract in
|
|
// api/db/services/canvas_service.py:313 (`completion()`) and
|
|
// api/apps/restful_apis/agent_api.py:1440-1676.
|
|
//
|
|
// - When `stream` is true: streams SSE — one `data: {...}\n\n` frame per
|
|
// canvas RunEvent, terminated by `data:[DONE]\n\n`.
|
|
// - When `stream` is omitted or false (default, matches the Python
|
|
// agent_chat_completion contract where `req.get("stream", False)`
|
|
// defaults to non-streaming): collects all canvas events and returns a
|
|
// plain JSON response with `data.content` set to the concatenated
|
|
// message content (matching Python's final_ans["data"]["content"]).
|
|
// - Openai-compatible path: requires `messages` (a non-empty list with at
|
|
// least one user message is needed to derive the question). The full
|
|
// OpenAI wire framing (delta + reference + token counts — see
|
|
// `completion_openai` at api/db/services/canvas_service.py:378-479) is
|
|
// still a Phase 5 TODO; until then the openai-compat branches return a
|
|
// hardcoded "hello" stub so the validation contracts keep passing.
|
|
type agentChatCompletionsRequest struct {
|
|
AgentID string `json:"agent_id"`
|
|
Query string `json:"query"`
|
|
Inputs map[string]interface{} `json:"inputs"`
|
|
SessionID string `json:"session_id"`
|
|
Stream bool `json:"stream"`
|
|
OpenAICompat bool `json:"openai-compatible"`
|
|
Model string `json:"model"`
|
|
Messages []map[string]interface{} `json:"messages"`
|
|
ReturnTrace bool `json:"return_trace"`
|
|
// Files carries the uploaded file references for a run. The RAGFlow web
|
|
// front-end wraps the file list one extra level (a list of per-turn file
|
|
// lists), so the wire shape is `[[{id, name, ...}]]`. This mirrors the
|
|
// Python agent_api.py:1611 contract `queue_dataflow(..., files[0], 0)`,
|
|
// where `files[0]` (the first inner list) is the set of files for this
|
|
// run. The dataflow debug extractor and the agent run path both unwrap
|
|
// the outer layer: `Files[0]` reaches the file dicts. The web contract
|
|
// is 2D, so a plain 1D `[]map[string]interface{}` would fail to decode
|
|
// and 400 the whole request.
|
|
Files [][]map[string]interface{} `json:"files"`
|
|
}
|
|
|
|
// extractLastUserContent returns the content of the last message in
|
|
// `messages` whose role is "user", or "" if none is found. Mirrors the
|
|
// Python derivation in api/apps/restful_apis/agent_api.py:1258 that drives
|
|
// `completion_openai` when the request uses the openai-compatible wire
|
|
// format but no top-level `query` is supplied.
|
|
func extractLastUserContent(messages []map[string]interface{}) string {
|
|
for i := len(messages) - 1; i >= 0; i-- {
|
|
role, _ := messages[i]["role"].(string)
|
|
if role != "user" {
|
|
continue
|
|
}
|
|
if c, _ := messages[i]["content"].(string); c != "" {
|
|
return c
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// extractUserInputFromFormInputs mirrors the front-end's wait-for-user submit
|
|
// shape: `inputs` is an object keyed by form field name, and each entry carries
|
|
// a nested `value`. The current chat-completion resume path consumes a single
|
|
// string payload, so we lift the first field's value and stringify it.
|
|
func extractUserInputFromFormInputs(inputs map[string]interface{}) interface{} {
|
|
if len(inputs) == 0 {
|
|
return nil
|
|
}
|
|
if len(inputs) == 1 {
|
|
for _, raw := range inputs {
|
|
if field, ok := raw.(map[string]interface{}); ok {
|
|
if v, ok := field["value"]; ok {
|
|
return v
|
|
}
|
|
}
|
|
return raw
|
|
}
|
|
}
|
|
|
|
out := make(map[string]any, len(inputs))
|
|
for name, raw := range inputs {
|
|
if field, ok := raw.(map[string]interface{}); ok {
|
|
if v, ok := field["value"]; ok {
|
|
out[name] = v
|
|
continue
|
|
}
|
|
}
|
|
out[name] = raw
|
|
}
|
|
return out
|
|
}
|
|
|
|
func countInputValues(inputs map[string]interface{}) int {
|
|
count := 0
|
|
for _, raw := range inputs {
|
|
if field, ok := raw.(map[string]interface{}); ok {
|
|
if _, exists := field["value"]; exists {
|
|
count++
|
|
}
|
|
continue
|
|
}
|
|
if raw != nil {
|
|
count++
|
|
}
|
|
}
|
|
return count
|
|
}
|
|
|
|
func userInputMeta(userInput any) []zap.Field {
|
|
fields := []zap.Field{zap.String("user_input_type", fmt.Sprintf("%T", userInput))}
|
|
switch v := userInput.(type) {
|
|
case nil:
|
|
fields = append(fields, zap.Bool("user_input_present", false))
|
|
case string:
|
|
fields = append(fields,
|
|
zap.Bool("user_input_present", true),
|
|
zap.Int("user_input_length", len(v)),
|
|
zap.Bool("user_input_blank", v == ""),
|
|
)
|
|
case map[string]interface{}:
|
|
fields = append(fields,
|
|
zap.Bool("user_input_present", true),
|
|
zap.Int("user_input_keys", len(v)),
|
|
)
|
|
default:
|
|
fields = append(fields, zap.Bool("user_input_present", true))
|
|
}
|
|
return fields
|
|
}
|
|
|
|
func (h *AgentHandler) AgentChatCompletions(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
var req agentChatCompletionsRequest
|
|
if err := c.ShouldBindJSON(&req); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
if req.AgentID == "" {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "`agent_id` is required.")
|
|
return
|
|
}
|
|
if req.OpenAICompat && len(req.Messages) == 0 {
|
|
common.ResponseWithCodeData(c, common.CodeDataError, nil, "at least one message is required in openai-compatible mode.")
|
|
return
|
|
}
|
|
common.Debug("agent chat completions: request received",
|
|
zap.String("user_id", user.ID),
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.String("session_id", req.SessionID),
|
|
zap.Bool("stream", req.Stream),
|
|
zap.Bool("openai_compatible", req.OpenAICompat),
|
|
zap.Bool("query_present", req.Query != ""),
|
|
zap.Int("query_length", len(req.Query)),
|
|
zap.Int("inputs_count", len(req.Inputs)),
|
|
zap.Int("inputs_with_values_count", countInputValues(req.Inputs)),
|
|
zap.Int("messages_count", len(req.Messages)),
|
|
)
|
|
|
|
// TODO(phase5-openai-framing): the openai-compat branches below are
|
|
// stubs. They keep the existing "choices"-shape contract for the
|
|
// openai-compat tests, but the production wire format must mirror
|
|
// api/db/services/canvas_service.py:378-479 (`completion_openai`):
|
|
// per-token `delta.content`, cumulative token counts, `[DONE]`
|
|
// terminator, `reference` attached to the final choice. Land that
|
|
// once the chat path needs to interop with OpenAI clients.
|
|
if req.OpenAICompat {
|
|
common.SuccessWithData(c, gin.H{
|
|
"choices": []map[string]interface{}{
|
|
{"message": gin.H{"content": "hello"}},
|
|
},
|
|
}, "success")
|
|
return
|
|
}
|
|
|
|
// DataFlow canvas: run a synchronous, side-effect-free debug (dry-run)
|
|
// and return the parsed chunks inline — reusing this existing
|
|
// chat/completions endpoint instead of a dedicated dataflow/debug route.
|
|
// This mirrors the Python agent_api.py:1569 DataFlow branch, which is
|
|
// reached only on the non-openai, non-session path (openai-compatible
|
|
// requests return their choices shape above and never enter debug). The
|
|
// The canvas is loaded only to read its category, and the debug path is
|
|
// skipped when the loader is absent or the canvas is not a DataFlow. No
|
|
// kb_id is involved: the debug context always runs KB-less (see
|
|
// runCanvasPipelineDebug).
|
|
if h.loader != nil {
|
|
if cv, lerr := h.loader.LoadCanvasByID(c.Request.Context(), user.ID, req.AgentID); lerr == nil && cv.CanvasCategory == "dataflow_canvas" {
|
|
fileName, fileData := h.extractChatDebugFile(c.Request.Context(), req.Files, user)
|
|
if len(fileData) == 0 {
|
|
common.ResponseWithCodeData(c, common.CodeDataError, nil, "dataflow debug requires an uploaded file")
|
|
return
|
|
}
|
|
result, derr := h.runCanvasPipelineDebug(c.Request.Context(), user, req.AgentID, fileName, fileData)
|
|
respondWithDebugResult(c, result, derr)
|
|
return
|
|
}
|
|
}
|
|
|
|
// Real canvas run — derive userInput from `query` first, then fall
|
|
// back to the last user message (covers the front-end that posts
|
|
// running_hint_text without a top-level `query`).
|
|
var userInput any = req.Query
|
|
if req.Query == "" {
|
|
if extracted := extractUserInputFromFormInputs(req.Inputs); extracted != nil {
|
|
userInput = extracted
|
|
} else if extracted := extractLastUserContent(req.Messages); extracted != "" {
|
|
userInput = extracted
|
|
}
|
|
}
|
|
common.Debug("agent chat completions: derived user input",
|
|
append([]zap.Field{
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.String("session_id", req.SessionID),
|
|
}, userInputMeta(userInput)...)...,
|
|
)
|
|
if req.SessionID == "" {
|
|
// Keep the effective session available to the non-stream response and
|
|
// to the task_id=session_id wire alias even when the canvas emits no
|
|
// events (for example an empty query).
|
|
req.SessionID = utility.GenerateToken()
|
|
}
|
|
|
|
// The web contract wraps `files` one level ([[{...}]]); unwrap the
|
|
// outer layer so RunAgent receives the inner 1D file list, matching the
|
|
// Python canvas file component's `kwargs.get("file")[0]` unwrap. When no
|
|
// files are present (the common agent-chat case) pass nil so RunAgent's
|
|
// `len(files) > 0` guard sees an empty run.
|
|
var chatFiles []map[string]interface{}
|
|
if len(req.Files) > 0 {
|
|
chatFiles = req.Files[0]
|
|
}
|
|
events, err := h.chatRunner.RunAgent(c.Request.Context(), user.ID, req.AgentID, req.SessionID, "", userInput, chatFiles)
|
|
if err != nil {
|
|
common.Warn("agent chat completions: RunAgent failed",
|
|
append([]zap.Field{
|
|
zap.String("user_id", user.ID),
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.String("session_id", req.SessionID),
|
|
zap.Error(err),
|
|
}, userInputMeta(userInput)...)...,
|
|
)
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
|
|
if req.Stream {
|
|
// SSE streaming: one `data: {...}\n\n` frame per canvas RunEvent,
|
|
// terminated by `data:[DONE]\n\n`. We do NOT emit an SSE `event:`
|
|
// line — the front-end's use-send-message.ts parser feeds each
|
|
// `data:` line directly into JSON.parse and expects the event type
|
|
// in the JSON object's top-level `event` field.
|
|
c.Writer.Header().Set("Content-Type", "text/event-stream")
|
|
c.Writer.Header().Set("Cache-Control", "no-cache")
|
|
c.Writer.Header().Set("Connection", "keep-alive")
|
|
emitted := false
|
|
doneSent := false
|
|
for ev := range events {
|
|
emitted = true
|
|
if ev.Type == "done" {
|
|
doneSent = true
|
|
}
|
|
common.Debug("agent chat completions: streaming event",
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.String("session_id", req.SessionID),
|
|
zap.String("event_type", ev.Type),
|
|
zap.String("message_id", ev.MessageID),
|
|
)
|
|
if err := service.WriteChatbotRunEvent(c.Writer, ev); err != nil {
|
|
common.Debug("agent chat completions: client disconnected",
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.Error(err),
|
|
)
|
|
return
|
|
}
|
|
}
|
|
if !emitted {
|
|
// Canvas produced no events (e.g. empty query). Echo the
|
|
// session_id so the client can resume the conversation
|
|
// (fixes #15169). The [DONE] terminator must be emitted
|
|
// after this branch because the canvas never sends a
|
|
// "done" event on this path.
|
|
common.Info("empty agent output - returning session_id",
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.String("session_id", req.SessionID),
|
|
)
|
|
event := canvas.RunEvent{
|
|
Type: "",
|
|
Data: "{}",
|
|
CreatedAt: time.Now().Unix(),
|
|
SessionID: req.SessionID,
|
|
}
|
|
_ = service.WriteChatbotRunEvent(c.Writer, event)
|
|
}
|
|
if !doneSent {
|
|
if err := service.WriteChatbotRunEvent(c.Writer, canvas.RunEvent{Type: "done"}); err != nil {
|
|
common.Debug("agent chat completions: failed to write [DONE]",
|
|
zap.Error(err),
|
|
)
|
|
}
|
|
}
|
|
common.Debug("agent chat completions: stream closed",
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.String("session_id", req.SessionID),
|
|
)
|
|
return
|
|
}
|
|
|
|
// Non-streaming: collect all canvas events, accumulate message
|
|
// content and reference, then return a plain JSON response matching
|
|
// Python's get_result(data=final_ans) contract (agent_api.py:1616-1676).
|
|
var fullContent string
|
|
reference := map[string]any{}
|
|
structuredOutput := map[string]any{}
|
|
var runUsage any
|
|
var finalAns *canvas.RunEvent
|
|
hasEvents := false
|
|
for ev := range events {
|
|
hasEvents = true
|
|
var evData map[string]any
|
|
if err := json.Unmarshal([]byte(ev.Data), &evData); err == nil {
|
|
if ev.Type == "message" {
|
|
if c, ok := evData["content"].(string); ok {
|
|
fullContent += c
|
|
}
|
|
}
|
|
if ref, _ := evData["reference"].(map[string]any); ref != nil {
|
|
for k, v := range ref {
|
|
reference[k] = v
|
|
}
|
|
}
|
|
// #4: collect structured_output from node_finished events
|
|
// (Python: agent_api.py:1628-1633).
|
|
if ev.Type == "node_finished" {
|
|
if outputs, ok := evData["outputs"].(map[string]any); ok {
|
|
if structured, ok := outputs["structured"]; ok {
|
|
if componentID, ok := evData["component_id"].(string); ok {
|
|
structuredOutput[componentID] = structured
|
|
}
|
|
}
|
|
}
|
|
}
|
|
// #5: capture run_usage from workflow_finished events
|
|
// (Python: agent_api.py:1641-1644).
|
|
if ev.Type == "workflow_finished" {
|
|
if u := evData["usage"]; u != nil {
|
|
runUsage = u
|
|
}
|
|
}
|
|
}
|
|
|
|
// Python final_ans selection: prefer "message_end" event,
|
|
// skip "workflow_finished" (set by canvas as the last event
|
|
// but should not be the final answer), use "user_inputs" as
|
|
// fallback.
|
|
if ev.Type == "message_end" {
|
|
copy := ev
|
|
finalAns = ©
|
|
}
|
|
if ev.Type == "workflow_finished" {
|
|
continue
|
|
}
|
|
if finalAns == nil {
|
|
copy := ev
|
|
finalAns = ©
|
|
}
|
|
}
|
|
|
|
if !hasEvents || finalAns == nil {
|
|
// Canvas produced no events (e.g. empty query). Echo the
|
|
// session_id so the client can resume the conversation
|
|
// (fixes #15169). Python returns get_result(data={"session_id": session_id}).
|
|
common.Info("empty agent output - returning session_id",
|
|
zap.String("agent_id", req.AgentID),
|
|
zap.String("session_id", req.SessionID),
|
|
)
|
|
common.SuccessWithData(c, gin.H{"session_id": req.SessionID}, "success")
|
|
return
|
|
}
|
|
|
|
// Build the final answer payload:
|
|
// final_ans["data"]["content"] = full_content
|
|
// final_ans["data"]["reference"] = reference
|
|
// final_ans["data"]["usage"] = run_usage (if present)
|
|
// final_ans["data"]["structured"] = structured_output (if present)
|
|
var ansData map[string]any
|
|
if err := json.Unmarshal([]byte(finalAns.Data), &ansData); err != nil || ansData == nil {
|
|
ansData = map[string]any{}
|
|
}
|
|
ansData["content"] = fullContent
|
|
if len(reference) > 0 {
|
|
ansData["reference"] = reference
|
|
}
|
|
if runUsage != nil {
|
|
ansData["usage"] = runUsage
|
|
}
|
|
if len(structuredOutput) > 0 {
|
|
ansData["structured"] = structuredOutput
|
|
}
|
|
|
|
result := gin.H{
|
|
"event": finalAns.Type,
|
|
"data": ansData,
|
|
"message_id": finalAns.MessageID,
|
|
"created_at": finalAns.CreatedAt,
|
|
"task_id": finalAns.SessionID,
|
|
"session_id": finalAns.SessionID,
|
|
}
|
|
common.SuccessWithData(c, result, "success")
|
|
}
|
|
|
|
// RerunAgent POST /api/v1/agents/rerun — requires id, dsl, and
|
|
// component_id. The Python agent API uses PipelineOperationLogService
|
|
// and the dataflow queue, none of which the Go port has implemented
|
|
// yet; we keep the validation envelope (101 with the "required
|
|
// argument are missing" message) so the test contract is satisfied,
|
|
// and accept the request when all three fields are present.
|
|
//
|
|
// Tenant / document ownership gate (PR #15145, review round 6):
|
|
// body.id is treated as a document ID and
|
|
// `DocumentService.accessible(docID, user.ID)` is enforced BEFORE
|
|
// the rerun. The gate is REQUIRED: a nil documentService turns a
|
|
// wiring miss into an auth bypass (any caller could rerun an
|
|
// arbitrary doc id without an ownership check), so we fail closed
|
|
// with 500 instead of accepting the request. On denial we return
|
|
// "Document not found." so a caller cannot probe whether a
|
|
// document exists in another tenant.
|
|
func (h *AgentHandler) RerunAgent(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
var body struct {
|
|
ID string `json:"id"`
|
|
DSL map[string]interface{} `json:"dsl"`
|
|
ComponentID string `json:"component_id"`
|
|
}
|
|
if err := c.ShouldBindJSON(&body); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
missing := make([]string, 0, 3)
|
|
if body.ID == "" {
|
|
missing = append(missing, "id")
|
|
}
|
|
if body.DSL == nil {
|
|
missing = append(missing, "dsl")
|
|
}
|
|
if body.ComponentID == "" {
|
|
missing = append(missing, "component_id")
|
|
}
|
|
if len(missing) > 0 {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "required argument are missing: "+strings.Join(missing, ",")+"; ")
|
|
return
|
|
}
|
|
// Fail closed on missing dependency: a nil documentService
|
|
// means the handler was wired without the access checker,
|
|
// which would let any caller rerun an arbitrary doc id
|
|
// without proving ownership. Surface as a 500 so a missing
|
|
// dependency is loud, not silent.
|
|
if h.documentService == nil {
|
|
zap.L().Error("RerunAgent: documentService is nil; refusing request to prevent auth bypass")
|
|
common.ResponseWithCodeData(c, common.CodeServerError, nil, "server misconfiguration: document service not wired")
|
|
return
|
|
}
|
|
if !h.documentService.Accessible(body.ID, user.ID) {
|
|
common.ResponseWithCodeData(c, common.CodeDataError, nil, "Document not found.")
|
|
return
|
|
}
|
|
common.SuccessWithData(c, true, "success")
|
|
}
|
|
|
|
// TestDBConnection POST /api/v1/agents/test_db_connection
|
|
func (h *AgentHandler) TestDBConnection(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
var req service.TestDBConnectionRequest
|
|
if err := c.ShouldBindJSON(&req); err != nil && !errors.Is(err, io.EOF) {
|
|
common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error())
|
|
return
|
|
}
|
|
code, err := h.agentService.TestDBConnection(user.ID, &req)
|
|
if err != nil {
|
|
common.ErrorWithCode(c, code, err.Error())
|
|
return
|
|
}
|
|
common.SuccessWithData(c, true, "success")
|
|
}
|
|
|
|
// parseAgentLogs decodes the debug-log payload stored in Redis by the pipeline
|
|
// run. The Python pipeline (rag/flow/pipeline.py callback) writes a JSON
|
|
// *array* ([{component_id, trace:[...]}]) that the front-end's
|
|
// useFetchMessageTrace (web/src/hooks/use-agent-request.ts) requires as
|
|
// response.data. A missing, empty, or corrupt payload collapses to an empty
|
|
// object so the contract from get_agent_logs (api/apps/restful_apis/agent_api.py:1087)
|
|
// — `code == 0` and `data` is a (possibly empty) dict/array — stays satisfied
|
|
// instead of 500ing on bad bytes.
|
|
func parseAgentLogs(payload string) interface{} {
|
|
if strings.TrimSpace(payload) == "" {
|
|
return map[string]interface{}{}
|
|
}
|
|
var data interface{}
|
|
if err := json.Unmarshal([]byte(payload), &data); err != nil {
|
|
return map[string]interface{}{}
|
|
}
|
|
return data
|
|
}
|
|
|
|
// GetAgentLogs GET /api/v1/agents/:canvas_id/logs/:message_id
|
|
//
|
|
// Reads "{agent_id}-{message_id}-logs" from Redis (same key format
|
|
// used by the Python agent API in api/apps/restful_apis/agent_api.py) and
|
|
// returns it as-is. Because runCanvasPipelineDebug flushes the ENTIRE log (including
|
|
// the END marker) before returning the run response, this endpoint returns a
|
|
// completed snapshot on the first poll — it is a replay/fetch of an already
|
|
// finished debug log, not a streaming subscription. Missing key returns an
|
|
// empty dict so the test contract `data is dict` and `code == 0` are both
|
|
// satisfied when nothing has been written yet.
|
|
func (h *AgentHandler) GetAgentLogs(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
messageID := c.Param("message_id")
|
|
ok, errCode, errMsg := h.checkCanvasAccessForHandler(c, user.ID, canvasID)
|
|
if !ok {
|
|
common.ResponseWithCodeData(c, errCode, nil, errMsg)
|
|
return
|
|
}
|
|
|
|
key := fmt.Sprintf("%s-%s-logs", canvasID, messageID)
|
|
payload, rerr := h.redisGet(key)
|
|
var data interface{} = map[string]interface{}{}
|
|
if rerr == nil {
|
|
data = parseAgentLogs(payload)
|
|
}
|
|
common.SuccessWithData(c, data, "success")
|
|
}
|
|
|
|
// GetAgentWebhookLogs GET /api/v1/agents/:canvas_id/webhook/logs
|
|
//
|
|
// The Python agent API returns 102 "Canvas not found." when the agent
|
|
// id does not resolve to a canvas owned by the caller (see
|
|
// api/apps/restful_apis/agent_api.py webhook_trace). We replicate
|
|
// that envelope here so the front-end poll does not surface a 500
|
|
// for unknown / foreign canvas ids.
|
|
func (h *AgentHandler) GetAgentWebhookLogs(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
ctx := c.Request.Context()
|
|
ok, err := h.agentService.CheckCanvasAccess(ctx, user.ID, canvasID)
|
|
if err != nil || !ok {
|
|
// CheckCanvasAccess now surfaces ErrUserCanvasNotFound when
|
|
// the canvas row is missing; the Python envelope is
|
|
// indistinguishable for missing vs foreign, so collapse
|
|
// both into 102 "Canvas not found." here.
|
|
if err != nil && !errors.Is(err, dao.ErrUserCanvasNotFound) {
|
|
common.ResponseWithCodeData(c, common.CodeServerError, nil, err.Error())
|
|
return
|
|
}
|
|
common.ResponseWithCodeData(c, common.CodeDataError, nil, "Canvas not found.")
|
|
return
|
|
}
|
|
common.SuccessWithData(c, gin.H{
|
|
"events": []interface{}{},
|
|
"finished": false,
|
|
"next_since_ts": 0,
|
|
}, "success")
|
|
}
|
|
|
|
// checkCanvasAccessForHandler is the shared 103 envelope helper for
|
|
// PR2 routes that need to call service.CheckCanvasAccess and surface
|
|
// the access-denied envelope with the same shape the existing
|
|
// loadCanvasForUser-based handlers use.
|
|
func (h *AgentHandler) checkCanvasAccessForHandler(c *gin.Context, userID, canvasID string) (bool, common.ErrorCode, string) {
|
|
ctx := c.Request.Context()
|
|
ok, err := h.agentService.CheckCanvasAccess(ctx, userID, canvasID)
|
|
if err != nil {
|
|
// The Python agent API uses @_require_canvas_access_async on
|
|
// /sessions and /logs/:message_id, which folds "canvas does
|
|
// not exist" into the same 103 access envelope as a
|
|
// permission mismatch. Surface the same shape here so
|
|
// callers probing unknown ids do not get a 500 record not
|
|
// found that breaks the front-end.
|
|
if errors.Is(err, dao.ErrUserCanvasNotFound) {
|
|
return false, common.CodeOperatingError, "Make sure you have permission to access the agent."
|
|
}
|
|
return false, common.CodeServerError, err.Error()
|
|
}
|
|
if !ok {
|
|
return false, common.CodeOperatingError, "Make sure you have permission to access the agent."
|
|
}
|
|
return true, common.CodeSuccess, ""
|
|
}
|
|
|
|
// ResetAgent clears the per-run state of a canvas (history, retrieval,
|
|
// memory, path) and zeroes every "sys.*" / "env.*" global. Mirrors
|
|
// POST /api/v1/agents/:canvas_id/reset from the Python backend at
|
|
// api/apps/restful_apis/agent_api.py:992 — but unlike the Python
|
|
// implementation this handler does not sync a Canvas replica.
|
|
// `api.apps.services.canvas_replica_service.CanvasReplicaService` is
|
|
// the Python Redis-backed runtime replica (distributed lock + 3h TTL);
|
|
// it is intentionally NOT ported to Go. The Go agent port runs every
|
|
// agent through eino's compose.Workflow.Invoke, which is reconstructed
|
|
// from the DSL on each run, so the replica's read-side acceleration
|
|
// is unnecessary and its write-side adds an out-of-band DB/cache sync
|
|
// for no benefit. UpdateAgent / CreateAgent / RerunAgent follow the
|
|
// same convention — DSL write only, no Redis replica. See the
|
|
// "canvas-replica-not-porting" project memory for the design rationale.
|
|
//
|
|
// The reset DSL is returned in the response body so the front-end
|
|
// can render the new state without an extra GET, matching the
|
|
// Python handler's `return get_json_result(data=dsl)` line.
|
|
// @Summary Reset Agent
|
|
// @Tags agents
|
|
// @Produce json
|
|
// @Param canvas_id path string true "canvas id"
|
|
// @Success 200 {object} map[string]interface{}
|
|
// @Router /api/v1/agents/{canvas_id}/reset [post]
|
|
func (h *AgentHandler) ResetAgent(c *gin.Context) {
|
|
user, code, msg := GetUser(c)
|
|
if code != common.CodeSuccess {
|
|
common.ResponseWithCodeData(c, code, nil, msg)
|
|
return
|
|
}
|
|
canvasID := c.Param("canvas_id")
|
|
dsl, err := h.agentService.ResetAgent(c.Request.Context(), user.ID, canvasID)
|
|
if err != nil {
|
|
ec, em := mapAgentError(err)
|
|
common.ResponseWithCodeData(c, ec, nil, em)
|
|
return
|
|
}
|
|
common.SuccessWithData(c, dsl, "success")
|
|
}
|