Merge pull request #1942 from satti-hari-krishna-reddy/nightly
feat: added runMCPAction to onprem
This commit is contained in:
@@ -4083,6 +4083,437 @@ func remoteOrgJobHandler(org shuffle.Org, interval int) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// /api/v1/mcp
|
||||
// /api/v1/agent
|
||||
// /api/v1/apps/{appid}/mcp
|
||||
func runMCPAction(resp http.ResponseWriter, request *http.Request) {
|
||||
cors := shuffle.HandleCors(resp, request)
|
||||
if cors {
|
||||
return
|
||||
}
|
||||
|
||||
ctx := shuffle.GetContext(request)
|
||||
parentExec := shuffle.WorkflowExecution{}
|
||||
user, err := shuffle.HandleApiAuthentication(resp, request)
|
||||
if err != nil {
|
||||
// Look for org_id query as app may be private
|
||||
// No validation is done here, as it's just running the app
|
||||
// to find a user
|
||||
orgId := request.URL.Query().Get("org_id")
|
||||
if len(orgId) > 0 {
|
||||
user.ActiveOrg.Id = orgId
|
||||
} else {
|
||||
executionId := request.URL.Query().Get("execution_id")
|
||||
authorization := request.URL.Query().Get("authorization")
|
||||
if len(executionId) == 0 || len(authorization) == 0 {
|
||||
log.Printf("[WARNING] Bad execution id/auth in single action validate (1): %#v, %#v. Continuing with the 'public' org id", executionId, authorization)
|
||||
err := shuffle.ValidateRequestOverload(resp, request)
|
||||
if err != nil {
|
||||
log.Printf("[INFO] Request overload for IP %s in single action execution", shuffle.GetRequestIp(request))
|
||||
resp.WriteHeader(429)
|
||||
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Too many requests. Please try again in 30 seconds."}`)))
|
||||
return
|
||||
}
|
||||
|
||||
user.Username = shuffle.GetRequestIp(request)
|
||||
user.ActiveOrg.Name = shuffle.GetRequestIp(request)
|
||||
user.ActiveOrg.Id = "public"
|
||||
|
||||
} else {
|
||||
// Find the execution
|
||||
exec, err := shuffle.GetWorkflowExecution(ctx, executionId)
|
||||
if err != nil {
|
||||
log.Printf("[WARNING] Bad execution id in single action validate (2): %s", err)
|
||||
resp.WriteHeader(401)
|
||||
resp.Write([]byte(`{"success": false, "reason": "Bad execution mapping (1)"}`))
|
||||
return
|
||||
}
|
||||
|
||||
if exec.Authorization != authorization {
|
||||
log.Printf("[WARNING] Bad execution auth in single action validate (3): %#v, %#v", exec.Authorization, authorization)
|
||||
resp.WriteHeader(403)
|
||||
resp.Write([]byte(`{"success": false, "reason": "Bad execution mapping (2)"}`))
|
||||
return
|
||||
}
|
||||
|
||||
parentExec = *exec
|
||||
user.ActiveOrg.Id = exec.OrgId
|
||||
if len(user.ActiveOrg.Id) == 0 {
|
||||
user.ActiveOrg.Id = exec.ExecutionOrg
|
||||
}
|
||||
|
||||
user.Username = fmt.Sprintf("org %s", user.ActiveOrg.Id)
|
||||
}
|
||||
}
|
||||
|
||||
if len(user.ActiveOrg.Id) == 0 {
|
||||
resp.WriteHeader(401)
|
||||
resp.Write([]byte(`{"success": false, "reason": "No org_id found to map back to"}`))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
location := strings.Split(request.URL.Path, "/")
|
||||
var fileId string
|
||||
if location[1] == "api" {
|
||||
if len(location) <= 3 {
|
||||
resp.WriteHeader(400)
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
|
||||
if len(location) == 4 {
|
||||
fileId = location[3] // /api/v1/agent or /api/v1/mcp
|
||||
} else {
|
||||
fileId = location[4] // /api/v1/apps/{appid}/mcp
|
||||
}
|
||||
}
|
||||
|
||||
//log.Printf("[AUDIT] User Authentication failed in execute SINGLE action - CONTINUING ANYWAY: %s. Found OrgID: %#v", err, user.ActiveOrg.Id)
|
||||
log.Printf("[AUDIT] User %s (%s) in org %s (%s) is running SINGLE App run for App ID '%s'", user.Username, user.Id, user.ActiveOrg.Name, user.ActiveOrg.Id, fileId)
|
||||
|
||||
body, err := ioutil.ReadAll(request.Body)
|
||||
if err != nil {
|
||||
log.Printf("[INFO] Failed single execution POST body read: %s", err)
|
||||
resp.WriteHeader(401)
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
|
||||
foundRequest := shuffle.MCPRequest{}
|
||||
//func HandleAiAgentExecutionStart(execution WorkflowExecution, startNode Action, createNextActions bool) (Action, error) {
|
||||
// Unmarshal it
|
||||
err = json.Unmarshal(body, &foundRequest)
|
||||
if err != nil {
|
||||
log.Printf("[INFO] Failed single execution POST body unmarshal: %s", err)
|
||||
resp.WriteHeader(400)
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
|
||||
foundEnvironment := "Shuffle"
|
||||
if len(foundRequest.Params.Environment) > 0 {
|
||||
foundEnvironment = foundRequest.Params.Environment
|
||||
}
|
||||
|
||||
if len(foundRequest.Params.Input.Text) < 5 {
|
||||
resp.WriteHeader(400)
|
||||
resp.Write([]byte(`{"success": false, "reason": "Input text is required and must be at least 5 characters"}`))
|
||||
return
|
||||
}
|
||||
|
||||
var newAction shuffle.Action
|
||||
|
||||
if fileId == "agent" {
|
||||
newAction = shuffle.Action{
|
||||
Name: "agent",
|
||||
AppName: "AI Agent",
|
||||
AppID: "shuffle_agent",
|
||||
AppVersion: "1.0.0",
|
||||
Environment: foundEnvironment,
|
||||
Parameters: []shuffle.WorkflowAppActionParameter{
|
||||
{
|
||||
Name: "app_name",
|
||||
Value: "openai",
|
||||
},
|
||||
{
|
||||
Name: "input",
|
||||
Value: foundRequest.Params.Input.Text,
|
||||
},
|
||||
},
|
||||
}
|
||||
} else {
|
||||
foundId := ""
|
||||
if len(foundRequest.Params.ToolID) > 0 {
|
||||
foundId = foundRequest.Params.ToolID
|
||||
} else {
|
||||
if len(foundRequest.Params.ToolName) == 32 {
|
||||
foundId = foundRequest.Params.ToolName
|
||||
} else {
|
||||
foundApps, err := shuffle.FindWorkflowAppByName(ctx, foundRequest.Params.ToolName)
|
||||
if err != nil || len(foundApps) == 0 {
|
||||
log.Printf("[INFO] Failed to find app by name '%s' in single execution: %s", foundRequest.Params.ToolName, err)
|
||||
resp.WriteHeader(400)
|
||||
resp.Write([]byte(`{"success": false, "reason": "Valid param.tool_id (app ID) is required"}`))
|
||||
return
|
||||
}
|
||||
|
||||
for _, app := range foundApps {
|
||||
if app.Name == foundRequest.Params.ToolName {
|
||||
foundId = app.ID
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
app, err := shuffle.GetApp(ctx, foundId, shuffle.User{}, false)
|
||||
if err != nil {
|
||||
log.Printf("[INFO] Failed to find app by id '%s' in single execution: %s", foundId, err)
|
||||
resp.WriteHeader(400)
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
|
||||
if !app.Public {
|
||||
if user.Id == app.Owner || user.ActiveOrg.Id == app.ReferenceOrg || shuffle.ArrayContains(app.Contributors, user.Id) {
|
||||
log.Printf("[AUDIT] Support & Admin user %s (%s) got access to app %s (MCP)", user.Username, user.Id, app.ID)
|
||||
} else if user.Role == "admin" && app.Owner == "" {
|
||||
log.Printf("[AUDIT] Any admin can GET %s (%s), since it doesn't have an owner (GET - MCP).", app.Name, app.ID)
|
||||
} else {
|
||||
log.Printf("[AUDIT] User %s (%s) in org %s (%s) was denied access to app %s (MCP)", user.Username, user.Id, user.ActiveOrg.Name, user.ActiveOrg.Id, app.ID)
|
||||
resp.WriteHeader(403)
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
} else {
|
||||
log.Printf("[AUDIT] User %s (%s) in org %s (%s) got access to public app %s (MCP)", user.Username, user.Id, user.ActiveOrg.Name, user.ActiveOrg.Id, app.ID)
|
||||
}
|
||||
|
||||
// Check permissions
|
||||
parsedName := strings.ToLower(strings.ReplaceAll(app.Name, " ", "_"))
|
||||
parsedApp := fmt.Sprintf("app:%s:%s", app.ID, parsedName)
|
||||
|
||||
// Run the action
|
||||
newAction = shuffle.Action{
|
||||
Name: "agent",
|
||||
AppName: "AI Agent",
|
||||
AppID: "shuffle_agent",
|
||||
AppVersion: "1.0.0",
|
||||
Environment: foundEnvironment,
|
||||
Parameters: []shuffle.WorkflowAppActionParameter{
|
||||
{
|
||||
Name: "app_name",
|
||||
Value: "openai",
|
||||
},
|
||||
{
|
||||
Name: "input",
|
||||
Value: foundRequest.Params.Input.Text,
|
||||
},
|
||||
{
|
||||
Name: "app_name",
|
||||
Value: parsedApp,
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
marshalledAction, err := json.Marshal(newAction)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed to marshal single action body: %s", err)
|
||||
resp.WriteHeader(500)
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
|
||||
// Special handling for agent calls
|
||||
if fileId == "agent" {
|
||||
if len(parentExec.ExecutionId) > 0 {
|
||||
targetActionId := request.URL.Query().Get("action_id")
|
||||
agentNodeFound := false
|
||||
var agentNode shuffle.Action
|
||||
|
||||
for i, action := range parentExec.Workflow.Actions {
|
||||
if len(targetActionId) > 0 && action.ID != targetActionId {
|
||||
continue
|
||||
}
|
||||
if len(targetActionId) == 0 && action.AppName != "AI Agent" {
|
||||
continue
|
||||
}
|
||||
|
||||
if len(foundRequest.Params.Input.Text) > 0 {
|
||||
for j, param := range parentExec.Workflow.Actions[i].Parameters {
|
||||
if param.Name == "input" {
|
||||
parentExec.Workflow.Actions[i].Parameters[j].Value = foundRequest.Params.Input.Text
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
agentNode = parentExec.Workflow.Actions[i]
|
||||
agentNodeFound = true
|
||||
log.Printf("[DEBUG][%s] AI Agent: Running AI Agent node '%s' (ID: %s)", parentExec.ExecutionId, agentNode.Label, agentNode.ID)
|
||||
break
|
||||
}
|
||||
|
||||
if !agentNodeFound {
|
||||
log.Printf("[ERROR][%s] No AI Agent node found in parent workflow (Target ID: %s)", parentExec.ExecutionId, targetActionId)
|
||||
resp.WriteHeader(400)
|
||||
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "No AI Agent node found in parent workflow matching request"}`)))
|
||||
return
|
||||
}
|
||||
|
||||
go shuffle.HandleAiAgentExecutionStart(parentExec, agentNode, false)
|
||||
|
||||
resp.WriteHeader(200)
|
||||
resp.Write([]byte(fmt.Sprintf(`{"success": true, "execution_id": "%s", "authorization": "%s", "mode": "hybrid"}`, parentExec.ExecutionId, parentExec.Authorization)))
|
||||
return
|
||||
|
||||
} else {
|
||||
//standalone mode
|
||||
workflowExecution, err := shuffle.PrepareSingleAction(ctx, user, "agent_starter", marshalledAction, false, "")
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed to prepare standalone agent execution: %s", err)
|
||||
resp.WriteHeader(500)
|
||||
resp.Write([]byte(`{"success": false, "reason": "Failed to create agent execution"}`))
|
||||
return
|
||||
}
|
||||
|
||||
log.Printf("[INFO] Standalone agent execution created: %s", workflowExecution.ExecutionId)
|
||||
resp.WriteHeader(200)
|
||||
resp.Write([]byte(fmt.Sprintf(`{"success": true, "execution_id": "%s", "authorization": "%s", "mode": "standalone"}`, workflowExecution.ExecutionId, workflowExecution.Authorization)))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// For non-agent calls (MCP or other apps): continue with existing flow
|
||||
workflowExecution, err := shuffle.PrepareSingleAction(ctx, user, "agent", marshalledAction, false, "")
|
||||
if fileId == "agent_starter" {
|
||||
log.Printf("[INFO] Returning early for agent_starter single action execution: %s", workflowExecution.ExecutionId)
|
||||
resp.WriteHeader(200)
|
||||
resp.Write([]byte(fmt.Sprintf(`{"success": true, "execution_id": "%s", "authorization": "%s"}`, workflowExecution.ExecutionId, workflowExecution.Authorization)))
|
||||
return
|
||||
}
|
||||
|
||||
debugUrl := fmt.Sprintf("/workflows/%s?execution_id=%s", workflowExecution.Workflow.ID, workflowExecution.ExecutionId)
|
||||
resp.Header().Add("X-Debug-Url", debugUrl)
|
||||
|
||||
if err != nil {
|
||||
returndata := shuffle.ResultChecker{
|
||||
Success: false,
|
||||
Reason: fmt.Sprintf("%s", err),
|
||||
}
|
||||
|
||||
// Special handler for decision reruns~
|
||||
if strings.Contains(err.Error(), "Successfully") {
|
||||
returndata.Success = true
|
||||
resp.WriteHeader(200)
|
||||
} else {
|
||||
log.Printf("[INFO] Failed workflowrequest POST read in single action (4): %s", err)
|
||||
resp.WriteHeader(400)
|
||||
}
|
||||
|
||||
respBytes, err := json.Marshal(returndata)
|
||||
if err != nil {
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
|
||||
resp.Write(respBytes)
|
||||
return
|
||||
}
|
||||
|
||||
foundEnv := ""
|
||||
params := []string{}
|
||||
for _, action := range workflowExecution.Workflow.Actions {
|
||||
for _, param := range action.Parameters {
|
||||
params = append(params, param.Name)
|
||||
}
|
||||
|
||||
if len(action.Environment) > 0 {
|
||||
foundEnv = action.Environment
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
go shuffle.IncrementCache(ctx, workflowExecution.OrgId, "workflow_executions")
|
||||
|
||||
executionRequest := shuffle.ExecutionRequest{
|
||||
ExecutionId: workflowExecution.ExecutionId,
|
||||
WorkflowId: workflowExecution.Workflow.ID,
|
||||
Authorization: workflowExecution.Authorization,
|
||||
Environments: []string{foundEnv},
|
||||
}
|
||||
|
||||
parsedEnv := fmt.Sprintf("%s_%s", strings.ToLower(strings.ReplaceAll(strings.ReplaceAll(foundEnv, " ", "-"), "_", "-")), workflowExecution.ExecutionOrg)
|
||||
|
||||
// Check if environment is distributed from parent org
|
||||
if len(workflowExecution.ExecutionOrg) > 0 {
|
||||
environments, err := shuffle.GetEnvironments(ctx, workflowExecution.ExecutionOrg)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed getting environments for org %s in single action. May fail to verify env.: %s", workflowExecution.ExecutionOrg, err)
|
||||
} else {
|
||||
for _, env := range environments {
|
||||
if env.Archived {
|
||||
continue
|
||||
}
|
||||
|
||||
if env.Name != foundEnv {
|
||||
continue
|
||||
}
|
||||
|
||||
if env.OrgId != workflowExecution.ExecutionOrg && len(env.OrgId) > 0 {
|
||||
if debug {
|
||||
log.Printf("[DEBUG][%s] Found suborg environment %s for org %s in single action. Re-mapping it to org-id %s", workflowExecution.ExecutionId, env.Name, env.OrgId, env.OrgId)
|
||||
}
|
||||
|
||||
parsedEnv = fmt.Sprintf("%s_%s", strings.ToLower(strings.ReplaceAll(strings.ReplaceAll(foundEnv, " ", "-"), "_", "-")), env.OrgId)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
log.Printf("[INFO][%s] Adding new single-action job to env queue (4 - MCP): %s", workflowExecution.ExecutionId, parsedEnv)
|
||||
err = shuffle.SetWorkflowQueue(ctx, executionRequest, parsedEnv)
|
||||
if err != nil {
|
||||
log.Printf("[WARNING][%s] Failed adding %s to db (single action queue): %s", workflowExecution.ExecutionId, parsedEnv, err)
|
||||
}
|
||||
|
||||
actionId := ""
|
||||
if len(workflowExecution.Workflow.Actions) == 1 {
|
||||
actionId = workflowExecution.Workflow.Actions[0].ID
|
||||
}
|
||||
|
||||
mappedResponse := shuffle.MCPResponse{
|
||||
Jsonrpc: foundRequest.Jsonrpc,
|
||||
ID: foundRequest.ID,
|
||||
}
|
||||
|
||||
singleResult := shuffle.HandleRetValidation(ctx, workflowExecution, 1, 45, actionId)
|
||||
agentOutput := shuffle.AgentOutput{}
|
||||
err = json.Unmarshal([]byte(singleResult.Result), &agentOutput)
|
||||
if err == nil && len(agentOutput.Output) > 0 {
|
||||
log.Printf("[INFO] Returning agent output in MCP response: %s", agentOutput.Output)
|
||||
|
||||
marshalledAgentOutput, err := json.Marshal(agentOutput)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed to marshal agent output in MCP response: %s", err)
|
||||
} else {
|
||||
marshalledOutput := map[string]interface{}{}
|
||||
err = json.Unmarshal(marshalledAgentOutput, &marshalledOutput)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed to unmarshal agent output in MCP response: %s", err)
|
||||
}
|
||||
|
||||
marshalledOutput["message"] = agentOutput.Output
|
||||
mappedResponse.Result = marshalledOutput
|
||||
}
|
||||
} else {
|
||||
marshalledSingleResult, err := json.Marshal(singleResult)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed to marshal single result in MCP response: %s", err)
|
||||
}
|
||||
|
||||
marshalledResult := map[string]interface{}{}
|
||||
err = json.Unmarshal(marshalledSingleResult, &marshalledResult)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed to unmarshal single result in MCP response: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
marshalledMappedResponse, err := json.Marshal(mappedResponse)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed to marshal mapped response in MCP response: %s", err)
|
||||
resp.WriteHeader(500)
|
||||
resp.Write([]byte(`{"success": false}`))
|
||||
return
|
||||
}
|
||||
|
||||
resp.WriteHeader(200)
|
||||
resp.Write(marshalledMappedResponse)
|
||||
}
|
||||
|
||||
func runInitCloudSetup() {
|
||||
action := shuffle.CloudSyncJob{
|
||||
Type: "setup",
|
||||
@@ -5456,6 +5887,12 @@ func initHandlers() {
|
||||
r.HandleFunc("/api/v1/apps/{key}/execute", executeSingleAction).Methods("POST", "OPTIONS")
|
||||
r.HandleFunc("/api/v1/apps/{key}/run", executeSingleAction).Methods("POST", "OPTIONS")
|
||||
|
||||
// Agent / MCP actions
|
||||
r.HandleFunc("/api/v1/apps/{key}/mcp", runMCPAction).Methods("POST", "OPTIONS")
|
||||
r.HandleFunc("/api/v1/mcp", runMCPAction).Methods("POST", "OPTIONS")
|
||||
r.HandleFunc("/api/v1/agent", runMCPAction).Methods("POST", "OPTIONS")
|
||||
|
||||
|
||||
//r.HandleFunc("/api/v1/apps/categories/run", shuffle.RunCategoryAction).Methods("POST", "OPTIONS")
|
||||
r.HandleFunc("/api/v1/apps/upload", handleAppZipUpload).Methods("POST", "OPTIONS")
|
||||
r.HandleFunc("/api/v1/apps/{appId}/activate", activateWorkflowAppDocker).Methods("GET", "OPTIONS")
|
||||
|
||||
Reference in New Issue
Block a user