diff --git a/py/semoss_automation_runtime.py b/py/semoss_automation_runtime.py index cc05e4e4d2..4eda8ca24f 100644 --- a/py/semoss_automation_runtime.py +++ b/py/semoss_automation_runtime.py @@ -99,7 +99,13 @@ def resolve_config(self, value: Any) -> Any: return self.resolve(value) -def execute_node(encoded_scope: str, encoded_source: str, max_output_bytes: int) -> Any: +def execute_node( + encoded_scope: str, + encoded_source: str, + max_output_bytes: int, + output_variable: str, + session_globals: dict[str, Any], +) -> Any: """Execute one persisted node module with a fresh module namespace.""" scope = _decode_scope(encoded_scope) source = _decode(encoded_source) @@ -110,7 +116,9 @@ def execute_node(encoded_scope: str, encoded_source: str, max_output_bytes: int) run = module.get("run") if not callable(run): raise ValueError("Automation node source must define callable run(scope).") - return _json_result(run(scope), max_output_bytes) + result = _json_result(run(scope), max_output_bytes) + _prepare_frame(result, output_variable, session_globals) + return result def execute_trigger( @@ -161,6 +169,25 @@ def _is_json_compatible(value: Any) -> bool: return False +def _prepare_frame( + value: Any, output_variable: str, session_globals: dict[str, Any] +) -> None: + """Retain row output where the standard SEMOSS Python-frame bridge expects it. + + The JSON result remains the Automation value used by history and downstream + scope. The DataFrame is the run-Insight view used for paged UI inspection. + """ + session_globals.pop(output_variable, None) + if not value or not isinstance(value, list): + return + if not all(isinstance(row, dict) for row in value): + return + + import pandas as pd + + session_globals[output_variable] = pd.DataFrame.from_records(value) + + def _json_result(value: Any, max_bytes: int) -> Any: try: serialized = json.dumps(value, allow_nan=False) diff --git a/src/prerna/reactor/automation/AGENTS.md b/src/prerna/reactor/automation/AGENTS.md index 825787f51b..279a8dba32 100644 --- a/src/prerna/reactor/automation/AGENTS.md +++ b/src/prerna/reactor/automation/AGENTS.md @@ -1,4 +1,4 @@ -# Automation Python — Agent Guide +# Automation — Agent Guide Automation projects persist a typed graph and one Python source file per Python-backed node. The graph is canonical: Java traverses its control edges, and Python executes each selected node module @@ -102,6 +102,13 @@ database writes use the database SDK's `ExecQuery` path, which retains edit auth behavior, and configured `insertData` guardrails. Generated updates always require a `WHERE` clause; use custom Python for an intentionally unbounded operation. Return a value so Java can store it under the node's `outputVar`. +Node output, run inputs, and aggregate scope remain bounded by `AutomationConstants`. Row-shaped node output is +also registered under its `outputVar` as a standard SEMOSS Python frame in that run's execution Insight, matching +Notebook's named-frame convention. `GetAutomationRun` returns the ordinary `FRAME_MAP` noun while the execution +Insight remains live; the UI reads it with the existing +`Frame | QueryAll | Offset | Limit | Collect` path. The frame is a display boundary, not a durable-data contract: +when the run Insight has closed, callers fall back to the persisted output preview. + The bridge reloads the Java-bound node from the immutable run snapshot and retains the callback insight's user/security context. It does not accept an arbitrary node definition, node id, engine id, or Java object from Python. Cancellation sets the DB flag, signals the same-pod Python socket @@ -111,9 +118,27 @@ job when possible, and is checked before each node and during waits. | Class | Purpose | | --- | --- | -| `AutomationDatabaseUtility` | Physical run records, node outputs, per-run claiming, and stale-run recovery in the scheduler DB. | +| `AutomationRunStore` | Physical run records, node outputs, per-run claiming, and stale-run recovery in the scheduler DB. | | `SchedulerOwlCreator` | Authoritative OWL schema for both scheduler-owned and automation-owned tables in that DB. | -| `AutomationPythonRunRegistry` | Same-pod Python socket interruption, heartbeat, and cancellation state. | +| `AutomationRunRegistry` | Same-pod Python socket interruption, heartbeat, and cancellation state. | +| `AutomationProjectService` | Project permissions, definition locking/persistence, reference validation, and derived-asset synchronization. | | `AutomationRuntimeUtils` | JSON serialization, scope construction, and output previews. | Do not bypass the per-run database claim, run snapshot, or DB/in-memory cancellation signal. + +## Java package layout + +Automation backend code is grouped by owned capability: + +| Package | Ownership | +| --- | --- | +| `prerna.reactor.automation` | Shared constants and graph/Python invocation support. | +| `prerna.reactor.automation.definition` | Canonical graph validation, node catalog, generated source, definition persistence, and authoring reactors. | +| `prerna.reactor.automation.project` | `AutomationProjectService`, project creation, permission-aware project coordination, and derived MCP assets. | +| `prerna.reactor.automation.run` | Durable run storage, execution lifecycle, cancellation registry, history, and run reactors. | +| `prerna.reactor.automation.agent` | Trace-linked child-agent access and action reactors. | +| `prerna.reactor.automation.utils` | Cross-boundary runtime JSON, scope, and preview helpers only. | + +Keep additions with the capability that owns their lifecycle. Do not recreate a +flat package, add generic `service` or `helpers` buckets, or create a subpackage +for one speculative abstraction. diff --git a/src/prerna/reactor/automation/AutomationRuntime.java b/src/prerna/reactor/automation/AutomationRuntime.java index 0ca069bb12..5025c99ec0 100644 --- a/src/prerna/reactor/automation/AutomationRuntime.java +++ b/src/prerna/reactor/automation/AutomationRuntime.java @@ -38,6 +38,7 @@ import java.util.List; import java.util.Map; +import prerna.reactor.automation.definition.AutomationDefinitionValidator; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.util.Constants; import prerna.util.Utility; @@ -50,12 +51,12 @@ * with only the source selected for the current node and a bounded JSON scope. * This class never discovers or executes arbitrary project files. */ -final class AutomationRuntime { +public final class AutomationRuntime { private AutomationRuntime() { } - static List> nodesForRun(AutomationDefinitionValidator.ValidatedDefinition definition) { + public static List> nodesForRun(AutomationDefinitionValidator.ValidatedDefinition definition) { List> nodes = new ArrayList<>(); for (Map original : controlOrderedNodes(definition)) { Map node = new LinkedHashMap<>(original); @@ -86,7 +87,8 @@ private static String defaultOutputVariable(Map node) { * initialization. Runtime traversal still selects only one condition path and * remains Java-owned. */ - static List> controlOrderedNodes(AutomationDefinitionValidator.ValidatedDefinition definition) { + static List> controlOrderedNodes( + AutomationDefinitionValidator.ValidatedDefinition definition) { Map> nodes = new LinkedHashMap<>(); Map incoming = new LinkedHashMap<>(); for (Map node : definition.nodes()) { @@ -128,7 +130,7 @@ static List> controlOrderedNodes(AutomationDefinitionValidat return ordered; } - static String startNodeId(AutomationDefinitionValidator.ValidatedDefinition definition) { + public static String startNodeId(AutomationDefinitionValidator.ValidatedDefinition definition) { for (Map node : definition.nodes()) { if (AutomationConstants.NODE_START.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) { return (String) node.get(AutomationConstants.NODE_FIELD_ID); @@ -137,7 +139,7 @@ static String startNodeId(AutomationDefinitionValidator.ValidatedDefinition defi throw new IllegalArgumentException("Automation definition has no trigger.start node."); } - static Map> controlTargets( + public static Map> controlTargets( AutomationDefinitionValidator.ValidatedDefinition definition) { Map> targets = new LinkedHashMap<>(); for (Map edge : definition.edges()) { @@ -155,8 +157,8 @@ static Map> controlTargets( /** * Runs one node module with the workflow scope supplied by the Java scheduler. */ - static String buildNodeInvocationScript(String source, Map scope) { - return buildPythonInvocation("execute_node", source, scope); + public static String buildNodeInvocationScript(String source, Map scope, String outputVariable) { + return buildPythonInvocation("execute_node", source, scope, outputVariable); } /** @@ -164,34 +166,37 @@ static String buildNodeInvocationScript(String source, Map scope * JSON-compatible globals. A trigger may also return a map from * {@code run(scope)} to define computed globals. */ - static String buildTriggerInvocationScript(String source, Map scope) { - return buildPythonInvocation("execute_trigger", source, scope); + public static String buildTriggerInvocationScript(String source, Map scope) { + return buildPythonInvocation("execute_trigger", source, scope, null); } - private static String buildPythonInvocation(String function, String source, Map scope) { + private static String buildPythonInvocation(String function, String source, Map scope, + String outputVariable) { Path runtimePath = Path.of(Utility.getBaseFolder(), Constants.PY_BASE_FOLDER, "semoss_automation_runtime.py") .toAbsolutePath().normalize(); if (!Files.isRegularFile(runtimePath)) { throw new IllegalStateException("Automation Python runtime is unavailable: " + runtimePath); } + String frameArguments = outputVariable == null ? "" + : ", " + AutomationRuntimeUtils.GSON.toJson(outputVariable) + ", globals()"; return """ import importlib.util as _automation_importlib _automation_spec = _automation_importlib.spec_from_file_location( "_semoss_automation_runtime", %s) _automation_runtime = _automation_importlib.module_from_spec(_automation_spec) _automation_spec.loader.exec_module(_automation_runtime) - _automation_runtime.%s("%s", "%s", %d) + _automation_runtime.%s("%s", "%s", %d%s) """.formatted(AutomationRuntimeUtils.GSON.toJson(runtimePath.toString()), function, encode(AutomationRuntimeUtils.toBoundedRuntimeJson(scope != null ? scope : Map.of(), AutomationConstants.RUN_SCOPE_MAX_BYTES, "Automation run scope")), - encode(source != null ? source : ""), AutomationConstants.NODE_OUTPUT_MAX_BYTES); + encode(source != null ? source : ""), AutomationConstants.NODE_OUTPUT_MAX_BYTES, frameArguments); } /** * Returns each trigger global's declared default for Get/Save responses and * playground defaults. */ - static Map declaredGlobals(AutomationDefinitionValidator.ValidatedDefinition definition) { + public static Map declaredGlobals(AutomationDefinitionValidator.ValidatedDefinition definition) { for (Map node : definition.nodes()) { if (AutomationConstants.NODE_START.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) { return triggerGlobalDefaults(node); @@ -204,7 +209,7 @@ static Map declaredGlobals(AutomationDefinitionValidator.Validat * Returns the trigger declarations held in * {@code trigger.start.config.globals}. */ - static List> triggerGlobalDefinitions( + public static List> triggerGlobalDefinitions( AutomationDefinitionValidator.ValidatedDefinition definition) { for (Map node : definition.nodes()) { if (AutomationConstants.NODE_START.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) { @@ -235,7 +240,7 @@ private static List> triggerGlobalDefinitions(Map node) { + public static String triggerSource(Map node) { Object rawConfig = node.get(AutomationConstants.NODE_FIELD_CONFIG); if (rawConfig instanceof Map raw) { return sourceValue(((Map) raw).get(AutomationConstants.CONFIG_PYTHON_SOURCE)); @@ -247,7 +252,7 @@ static String triggerSource(Map node) { * Returns the declared default for each trigger global the run should seed into * scope. */ - static Map triggerGlobalDefaults(Map node) { + public static Map triggerGlobalDefaults(Map node) { Map globals = new LinkedHashMap<>(); for (Map global : triggerGlobalDefinitions(node)) { if (global.get("name") instanceof String name @@ -264,7 +269,7 @@ static Map triggerGlobalDefaults(Map node) { * run-history reader compare a child agent's returned workspace against this * value, so they resolve it the same way. */ - static String configuredAgentWorkspaceId(Map node) { + public static String configuredAgentWorkspaceId(Map node) { if (!AutomationConstants.NODE_AGENT_RUN.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) { return null; } @@ -283,7 +288,7 @@ private static String sourceValue(Object value) { return value instanceof String source && !source.isBlank() ? source : null; } - static Object normalizeNodeResult(Object output) { + public static Object normalizeNodeResult(Object output) { Object value = output; if (value instanceof String string) { try { diff --git a/src/prerna/reactor/automation/agent/AGENTS.md b/src/prerna/reactor/automation/agent/AGENTS.md new file mode 100644 index 0000000000..d972d38f0f --- /dev/null +++ b/src/prerna/reactor/automation/agent/AGENTS.md @@ -0,0 +1,19 @@ +# Automation Agent Package + +This package exposes only trace-linked agent-run operations for Automation. +Follow the parent `../AGENTS.md` for the complete Automation contract. + +## Invariants + +- Apply the Automation project ACL before accessing a child agent run. +- Require the exact persisted project/run/node/agent-run relationship; do not + accept a caller-provided room or agent run as sufficient authorization. +- Use the standard Agent run service for retrieval and actions. +- Keep view and control authorization distinct. +- Do not move general agent execution into this package; it remains owned by + `prerna.reactor.agent`. + +## Verification + +Cover missing authentication, project access, trace mismatch, view-only access, +and edit/control access whenever these reactors change. diff --git a/src/prerna/reactor/automation/AutomationAgentRunAccess.java b/src/prerna/reactor/automation/agent/AutomationAgentRunAccess.java similarity index 84% rename from src/prerna/reactor/automation/AutomationAgentRunAccess.java rename to src/prerna/reactor/automation/agent/AutomationAgentRunAccess.java index 2f644d680e..9426f051c1 100644 --- a/src/prerna/reactor/automation/AutomationAgentRunAccess.java +++ b/src/prerna/reactor/automation/agent/AutomationAgentRunAccess.java @@ -25,11 +25,13 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.agent; import prerna.auth.User; import prerna.om.Insight; import prerna.project.api.IProject; +import prerna.reactor.automation.project.AutomationProjectService; +import prerna.reactor.automation.run.AutomationRunStore; /** * Permission boundary for Automation's trace-linked agent-run APIs. @@ -39,7 +41,7 @@ * the generic agent-run owner's identity. It first uses the normal project ACL, * then requires an exact persisted Automation run/node/agent-run relationship. */ -final class AutomationAgentRunAccess { +public final class AutomationAgentRunAccess { private AutomationAgentRunAccess() { } @@ -47,22 +49,21 @@ private AutomationAgentRunAccess() { static Access authorizeView(Insight insight, String projectId, String automationRunId, String nodeId, String agentRunId) { User user = authenticatedUser(insight); - IProject project = AutomationProjectUtils.getViewableAutomationProject(user, required(projectId, "project")); + IProject project = AutomationProjectService.getViewableAutomationProject(user, required(projectId, "project")); requireTrace(project.getProjectId(), automationRunId, nodeId, agentRunId); return new Access(project.getProjectId(), canControl(user, project.getProjectId())); } - static Access authorizeEdit(Insight insight, String projectId, String automationRunId, String nodeId, + public static void authorizeEdit(Insight insight, String projectId, String automationRunId, String nodeId, String agentRunId) { User user = authenticatedUser(insight); - IProject project = AutomationProjectUtils.getEditableAutomationProject(user, required(projectId, "project")); + IProject project = AutomationProjectService.getEditableAutomationProject(user, required(projectId, "project")); requireTrace(project.getProjectId(), automationRunId, nodeId, agentRunId); - return new Access(project.getProjectId(), true); } private static boolean canControl(User user, String projectId) { try { - AutomationProjectUtils.getEditableAutomationProject(user, projectId); + AutomationProjectService.getEditableAutomationProject(user, projectId); return true; } catch (IllegalArgumentException e) { return false; @@ -70,7 +71,7 @@ private static boolean canControl(User user, String projectId) { } private static void requireTrace(String projectId, String automationRunId, String nodeId, String agentRunId) { - if (!AutomationDatabaseUtility.hasAgentRunTrace(projectId, required(automationRunId, "automation run"), + if (!AutomationRunStore.hasAgentRunTrace(projectId, required(automationRunId, "automation run"), required(nodeId, "node"), required(agentRunId, "agent run"))) { throw new SecurityException("The requested agent run is not trace-linked to this Automation run and node."); } diff --git a/src/prerna/reactor/automation/GetAutomationAgentRunReactor.java b/src/prerna/reactor/automation/agent/GetAutomationAgentRunReactor.java similarity index 98% rename from src/prerna/reactor/automation/GetAutomationAgentRunReactor.java rename to src/prerna/reactor/automation/agent/GetAutomationAgentRunReactor.java index b8e4606150..f6ab058a4c 100644 --- a/src/prerna/reactor/automation/GetAutomationAgentRunReactor.java +++ b/src/prerna/reactor/automation/agent/GetAutomationAgentRunReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.agent; import java.util.Map; diff --git a/src/prerna/reactor/automation/ResolveAutomationAgentRunActionReactor.java b/src/prerna/reactor/automation/agent/ResolveAutomationAgentRunActionReactor.java similarity index 99% rename from src/prerna/reactor/automation/ResolveAutomationAgentRunActionReactor.java rename to src/prerna/reactor/automation/agent/ResolveAutomationAgentRunActionReactor.java index 0e00e00772..f4e0400d73 100644 --- a/src/prerna/reactor/automation/ResolveAutomationAgentRunActionReactor.java +++ b/src/prerna/reactor/automation/agent/ResolveAutomationAgentRunActionReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.agent; import java.util.HashMap; import java.util.List; diff --git a/src/prerna/reactor/automation/StopAutomationAgentRunReactor.java b/src/prerna/reactor/automation/agent/StopAutomationAgentRunReactor.java similarity index 98% rename from src/prerna/reactor/automation/StopAutomationAgentRunReactor.java rename to src/prerna/reactor/automation/agent/StopAutomationAgentRunReactor.java index be35635b04..519d686377 100644 --- a/src/prerna/reactor/automation/StopAutomationAgentRunReactor.java +++ b/src/prerna/reactor/automation/agent/StopAutomationAgentRunReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.agent; import java.util.Map; diff --git a/src/prerna/reactor/automation/definition/AGENTS.md b/src/prerna/reactor/automation/definition/AGENTS.md new file mode 100644 index 0000000000..f6381d5a35 --- /dev/null +++ b/src/prerna/reactor/automation/definition/AGENTS.md @@ -0,0 +1,30 @@ +# Automation Definition Package + +This package owns the canonical workflow definition and authoring API. Follow +the parent `../AGENTS.md` for the complete Automation contract. + +## Ownership + +- Validate and normalize the typed graph. +- Define the server-owned node catalog and configuration schema. +- Persist the graph and per-node Python source as one definition aggregate. +- Render generated node source and expose authoring reactors. +- Evaluate bounded control conditions; never execute arbitrary expressions. + +## Invariants + +- The persisted graph is canonical. Do not create a second node contract in the + UI, MCP metadata, or generated Python. +- Authoring may save an incomplete acyclic draft; execution validation remains + stricter and requires a connected runnable graph. +- Node types, ports, configuration fields, and output fields come from + `AutomationNodeCatalog` and stable enums rather than string-prefix logic. +- Definition mutation must enter through `AutomationProjectService` so project + permissions, locking, revision checks, asset sync, and MCP regeneration stay + together. + +## Verification + +Keep focused tests under `test/prerna/reactor/automation/definition`. When a +node contract changes, cover validation, catalog serialization, generated +source, and scope descriptors as applicable. diff --git a/src/prerna/reactor/automation/AddAutomationStepReactor.java b/src/prerna/reactor/automation/definition/AddAutomationStepReactor.java similarity index 97% rename from src/prerna/reactor/automation/AddAutomationStepReactor.java rename to src/prerna/reactor/automation/definition/AddAutomationStepReactor.java index c8da531468..9b8365913d 100644 --- a/src/prerna/reactor/automation/AddAutomationStepReactor.java +++ b/src/prerna/reactor/automation/definition/AddAutomationStepReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.ArrayList; import java.util.LinkedHashMap; @@ -36,6 +36,8 @@ import prerna.ds.py.PyUtils; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; @@ -77,7 +79,7 @@ public NounMetadata execute() { String branchPort = this.keyValue.get(BRANCH_PORT_KEY); Map config = AutomationRuntimeUtils.parseJsonObject(required(CONFIG_KEY), "config"); String customSource = customSource(nodeType, config); - return AutomationProjectUtils.withLockedDefinition(projectId, files -> addStep(projectId, files, nodeType, + return AutomationProjectService.withLockedDefinition(projectId, files -> addStep(projectId, files, nodeType, label, outputVar, afterNodeId, branchPort, config, customSource)); } @@ -133,7 +135,7 @@ private NounMetadata addStep(String projectId, AutomationDefinitionService.Defin if (customSource != null) { nodeSources.put(nodeId, customSource); } - AutomationDefinitionService.DefinitionFiles saved = AutomationProjectUtils.saveDefinition(projectId, json, + AutomationDefinitionService.DefinitionFiles saved = AutomationProjectService.saveDefinition(projectId, json, nodeSources, this.insight.getUser()); Map result = new LinkedHashMap<>(); result.put("node", node); @@ -144,7 +146,7 @@ private NounMetadata addStep(String projectId, AutomationDefinitionService.Defin } private String editableProjectId() { - return AutomationProjectUtils + return AutomationProjectService .getEditableAutomationProject(this.insight.getUser(), required(ReactorKeysEnum.PROJECT.getKey())) .getProjectId(); } diff --git a/src/prerna/reactor/automation/AutomationConditionEvaluator.java b/src/prerna/reactor/automation/definition/AutomationConditionEvaluator.java similarity index 99% rename from src/prerna/reactor/automation/AutomationConditionEvaluator.java rename to src/prerna/reactor/automation/definition/AutomationConditionEvaluator.java index 917c5c2829..4ea972c806 100644 --- a/src/prerna/reactor/automation/AutomationConditionEvaluator.java +++ b/src/prerna/reactor/automation/definition/AutomationConditionEvaluator.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.math.BigDecimal; import java.util.ArrayList; @@ -46,7 +46,7 @@ * parentheses. Evaluation never invokes Python, Pixel, reflection, or user * code. */ -final class AutomationConditionEvaluator { +public final class AutomationConditionEvaluator { private static final int MAX_EXPRESSION_LENGTH = 4_096; private static final int MAX_TOKENS = 256; @@ -62,7 +62,7 @@ static void validate(String expression) { } /** Evaluates an expression against the current run's read-only scope. */ - static boolean evaluate(String expression, Map scope) { + public static boolean evaluate(String expression, Map scope) { Object result = parse(expression).evaluate(scope != null ? scope : Map.of()); if (result instanceof Boolean value) { return value; diff --git a/src/prerna/reactor/automation/AutomationDefinitionService.java b/src/prerna/reactor/automation/definition/AutomationDefinitionService.java similarity index 99% rename from src/prerna/reactor/automation/AutomationDefinitionService.java rename to src/prerna/reactor/automation/definition/AutomationDefinitionService.java index 19482f91c2..d35132f5fa 100644 --- a/src/prerna/reactor/automation/AutomationDefinitionService.java +++ b/src/prerna/reactor/automation/definition/AutomationDefinitionService.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.io.IOException; import java.nio.charset.StandardCharsets; @@ -42,13 +42,14 @@ import java.util.TreeMap; import java.util.stream.Stream; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - import com.google.gson.Gson; import com.google.gson.GsonBuilder; import com.google.gson.JsonParser; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.util.AssetUtility; @@ -328,7 +329,7 @@ private static void validateNodeSource(String nodeId, Map node, * @param source persisted node source * @return {@code true} when the source binds a top-level {@code run} */ - static boolean definesRunEntryPoint(String source) { + public static boolean definesRunEntryPoint(String source) { return source.lines().anyMatch(line -> { if (line.isEmpty() || Character.isWhitespace(line.charAt(0))) { return false; diff --git a/src/prerna/reactor/automation/AutomationDefinitionValidator.java b/src/prerna/reactor/automation/definition/AutomationDefinitionValidator.java similarity index 91% rename from src/prerna/reactor/automation/AutomationDefinitionValidator.java rename to src/prerna/reactor/automation/definition/AutomationDefinitionValidator.java index 4564c7a23e..e24cde4bdb 100644 --- a/src/prerna/reactor/automation/AutomationDefinitionValidator.java +++ b/src/prerna/reactor/automation/definition/AutomationDefinitionValidator.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; @@ -41,7 +41,6 @@ import java.util.TreeMap; import com.google.gson.JsonParseException; - import net.sf.jsqlparser.JSQLParserException; import net.sf.jsqlparser.parser.CCJSqlParserUtil; import net.sf.jsqlparser.statement.Statement; @@ -49,7 +48,9 @@ import net.sf.jsqlparser.statement.insert.Insert; import net.sf.jsqlparser.statement.select.Select; import net.sf.jsqlparser.statement.update.Update; + import prerna.ds.py.PyUtils; +import prerna.reactor.automation.AutomationConstants; import prerna.reactor.automation.utils.AutomationRuntimeUtils; /** @@ -236,22 +237,51 @@ private static void validateNodeConfig(String nodeId, AutomationNodeType nodeTyp requireConfigString(nodeId, config, "prompt"); validateOptionalConfigObject(nodeId, config, AutomationConstants.CONFIG_PARAM_VALUES); } - case MODEL_EMBEDDINGS -> requireConfigString(nodeId, config, "text"); + case MODEL_EMBEDDINGS -> { + requireConfigString(nodeId, config, "text"); + validateOptionalConfigObject(nodeId, config, AutomationConstants.CONFIG_PARAM_VALUES); + } case MODEL_NER -> { requireConfigString(nodeId, config, "text"); requireConfigStringListOrPlaceholder(nodeId, config, "entities"); + validateOptionalStringListOrPlaceholder(nodeId, config, "maskEntities"); } case MODEL_VISION -> { requireConfigString(nodeId, config, "prompt"); - requireConfigString(nodeId, config, "image"); + validateOptionalConfigString(nodeId, config, "image"); + validateOptionalStringListOrPlaceholder(nodeId, config, "urls"); + if (!hasNonblankConfigString(config, "image") && !hasNonemptyConfigList(config, "urls")) { + throw new IllegalArgumentException( + "Node '" + nodeId + "' must provide config.image or config.urls for a vision model."); + } + validateOptionalConfigObject(nodeId, config, AutomationConstants.CONFIG_PARAM_VALUES); + } + case STORAGE_READ -> { + requireConfigString(nodeId, config, "path"); + validateOptionalConfigBoolean(nodeId, config, "convertToPdf"); + } + case STORAGE_DELETE -> { + requireConfigString(nodeId, config, "path"); + validateOptionalConfigBoolean(nodeId, config, "leaveFolderStructure"); } - case STORAGE_READ, STORAGE_DELETE -> requireConfigString(nodeId, config, "path"); case STORAGE_UPLOAD -> { requireConfigString(nodeId, config, "path"); requireConfigString(nodeId, config, "destination"); + validateOptionalConfigObject(nodeId, config, "metadata"); + } + case STORAGE_DOWNLOAD -> { + requireConfigString(nodeId, config, "path"); + validateOptionalConfigString(nodeId, config, "version"); + } + case VECTOR_SEARCH -> { + requireConfigString(nodeId, config, "value"); + validateOptionalConfigObject(nodeId, config, "filters"); + validateOptionalConfigObject(nodeId, config, AutomationConstants.CONFIG_PARAM_VALUES); + } + case VECTOR_ADD, VECTOR_DELETE -> { + requireConfigString(nodeId, config, "value"); + validateOptionalConfigObject(nodeId, config, AutomationConstants.CONFIG_PARAM_VALUES); } - case STORAGE_DOWNLOAD -> requireConfigString(nodeId, config, "path"); - case VECTOR_SEARCH, VECTOR_ADD, VECTOR_DELETE -> requireConfigString(nodeId, config, "value"); case FUNCTION_EXECUTE -> requireConfigObject(nodeId, config, "arguments"); case APP_PIXEL -> requireConfigString(nodeId, config, "pixel"); case AGENT_RUN -> { @@ -474,6 +504,39 @@ private static void validateOptionalConfigObject(String nodeId, Map config, String key) { + if (!config.containsKey(key) || config.get(key) == null || "".equals(config.get(key))) { + return; + } + requireConfigString(nodeId, config, key); + } + + private static void validateOptionalConfigBoolean(String nodeId, Map config, String key) { + Object value = config.get(key); + if (value != null && !(value instanceof Boolean)) { + throw new IllegalArgumentException("Node '" + nodeId + "' config." + key + " must be a boolean."); + } + } + + private static void validateOptionalStringListOrPlaceholder(String nodeId, Map config, + String key) { + if (!config.containsKey(key) || config.get(key) == null + || config.get(key) instanceof List list && list.isEmpty()) { + return; + } + requireConfigStringListOrPlaceholder(nodeId, config, key); + } + + private static boolean hasNonblankConfigString(Map config, String key) { + return config.get(key) instanceof String value && !value.isBlank(); + } + + private static boolean hasNonemptyConfigList(Map config, String key) { + Object value = config.get(key); + return value instanceof List list && !list.isEmpty() + || value instanceof String placeholder && isScopePlaceholder(placeholder); + } + private static void requireConfigObject(String nodeId, Map config, String key) { Object value = config.get(key); if (value instanceof Map) { @@ -495,7 +558,7 @@ private static void requireConfigObject(String nodeId, Map confi private static void requireConfigStringListOrPlaceholder(String nodeId, Map config, String key) { Object value = config.get(key); - if (value instanceof String placeholder && placeholder.matches("\\$\\{[A-Za-z_][A-Za-z0-9_]*\\}")) { + if (value instanceof String placeholder && isScopePlaceholder(placeholder)) { return; } if (value instanceof List values && !values.isEmpty() @@ -506,6 +569,10 @@ private static void requireConfigStringListOrPlaceholder(String nodeId, Map config) { if (!AutomationConstants.NODE_DATABASE_QUERY.equals(nodeType) && !AutomationConstants.NODE_DATABASE_INSERT.equals(nodeType) diff --git a/src/prerna/reactor/automation/AutomationNodeCatalog.java b/src/prerna/reactor/automation/definition/AutomationNodeCatalog.java similarity index 85% rename from src/prerna/reactor/automation/AutomationNodeCatalog.java rename to src/prerna/reactor/automation/definition/AutomationNodeCatalog.java index 28056e0b00..8b69b52fc4 100644 --- a/src/prerna/reactor/automation/AutomationNodeCatalog.java +++ b/src/prerna/reactor/automation/definition/AutomationNodeCatalog.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.ArrayList; import java.util.Collections; @@ -36,13 +36,14 @@ import java.util.Set; import prerna.engine.api.IEngine; -import prerna.reactor.automation.AutomationNodeDefinition.ConfigField; -import prerna.reactor.automation.AutomationNodeDefinition.ConfigFieldType; -import prerna.reactor.automation.AutomationNodeDefinition.OutputField; -import prerna.reactor.automation.AutomationNodeDefinition.OutputFieldType; -import prerna.reactor.automation.AutomationNodeDefinition.Port; -import prerna.reactor.automation.AutomationNodeDefinition.PortDirection; -import prerna.reactor.automation.AutomationNodeDefinition.PortKind; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.definition.AutomationNodeDefinition.ConfigField; +import prerna.reactor.automation.definition.AutomationNodeDefinition.ConfigFieldType; +import prerna.reactor.automation.definition.AutomationNodeDefinition.OutputField; +import prerna.reactor.automation.definition.AutomationNodeDefinition.OutputFieldType; +import prerna.reactor.automation.definition.AutomationNodeDefinition.Port; +import prerna.reactor.automation.definition.AutomationNodeDefinition.PortDirection; +import prerna.reactor.automation.definition.AutomationNodeDefinition.PortKind; /** * Server-owned catalog of Automation node authoring contracts. @@ -121,21 +122,29 @@ private static List createDefinitions() { field("paramValues", ConfigFieldType.JSON, "Model parameters", false, Map.of())), controlInputs(), controlAndResultOutputs())); definitions.add(definition(AutomationNodeType.MODEL_EMBEDDINGS, "Create embeddings", - "Turn text into vectors with a model engine.", orderedConfig("engineId", "", "text", ""), + "Turn text into vectors with a model engine.", + orderedConfig("engineId", "", "text", "", "paramValues", Map.of()), List.of(engineField(IEngine.CATALOG_TYPE.MODEL), - field("text", ConfigFieldType.TEXT, "Text to embed", true, "")), + field("text", ConfigFieldType.TEXT, "Text to embed", true, ""), + field("paramValues", ConfigFieldType.JSON, "Model parameters", false, Map.of())), controlInputs(), controlAndResultOutputs())); definitions.add(definition(AutomationNodeType.MODEL_VISION, "Analyze image", - "Ask a vision model about an image.", orderedConfig("engineId", "", "image", "", "prompt", ""), + "Ask a multimodal model about one or more files or URLs.", + orderedConfig("engineId", "", "systemPrompt", "", "image", "", "urls", List.of(), "prompt", "", + "paramValues", Map.of()), List.of(engineField(IEngine.CATALOG_TYPE.MODEL), - field("image", ConfigFieldType.STRING, "Image path or URL", true, ""), - field("prompt", ConfigFieldType.TEXT, "Question about the image", true, "")), + field("systemPrompt", ConfigFieldType.TEXT, "Instructions for the model", false, ""), + field("image", ConfigFieldType.STRING, "Media path", false, ""), + field("urls", ConfigFieldType.STRING_LIST, "Media URLs", false, List.of()), + field("prompt", ConfigFieldType.TEXT, "Question about the media", true, ""), + field("paramValues", ConfigFieldType.JSON, "Model parameters", false, Map.of())), controlInputs(), controlAndResultOutputs())); definitions.add(definition(AutomationNodeType.MODEL_NER, "Extract entities", "Find named entities in text.", - orderedConfig("engineId", "", "text", "", "entities", List.of()), + orderedConfig("engineId", "", "text", "", "entities", List.of(), "maskEntities", List.of()), List.of(engineField(IEngine.CATALOG_TYPE.MODEL), field("text", ConfigFieldType.TEXT, "Text to analyze", true, ""), - field("entities", ConfigFieldType.STRING_LIST, "Entity types", true, List.of())), + field("entities", ConfigFieldType.STRING_LIST, "Entity types", true, List.of()), + field("maskEntities", ConfigFieldType.STRING_LIST, "Entity types to mask", false, List.of())), controlInputs(), controlAndResultOutputs())); definitions.add(storageDefinition(AutomationNodeType.STORAGE_LIST, "List files", @@ -233,6 +242,19 @@ private static AutomationNodeDefinition storageDefinition(AutomationNodeType nod fields.add(field("destination", ConfigFieldType.STRING, "Destination", requireDestination, destinationDefault)); } + if (nodeType == AutomationNodeType.STORAGE_READ) { + defaultConfig.put("convertToPdf", false); + fields.add(field("convertToPdf", ConfigFieldType.BOOLEAN, "Convert supported files to PDF", false, false)); + } else if (nodeType == AutomationNodeType.STORAGE_UPLOAD) { + defaultConfig.put("metadata", Map.of()); + fields.add(field("metadata", ConfigFieldType.JSON, "Metadata", false, Map.of())); + } else if (nodeType == AutomationNodeType.STORAGE_DOWNLOAD) { + defaultConfig.put("version", ""); + fields.add(field("version", ConfigFieldType.STRING, "Version", false, "")); + } else if (nodeType == AutomationNodeType.STORAGE_DELETE) { + defaultConfig.put("leaveFolderStructure", false); + fields.add(field("leaveFolderStructure", ConfigFieldType.BOOLEAN, "Keep empty folders", false, false)); + } List outputFields = nodeType == AutomationNodeType.STORAGE_DOWNLOAD ? List.of(outputField("success", OutputFieldType.BOOLEAN, "Succeeded", "Whether the storage transfer completed.", true), @@ -254,15 +276,21 @@ private static AutomationNodeDefinition storageDefinition(AutomationNodeType nod private static AutomationNodeDefinition vectorDefinition(AutomationNodeType nodeType, String label, String description, boolean includeLimit) { Map defaultConfig = includeLimit - ? orderedConfig("engineId", "", "value", "", "limit", AutomationConstants.DEFAULT_VECTOR_SEARCH_LIMIT) - : orderedConfig("engineId", "", "value", ""); + ? orderedConfig("engineId", "", "value", "", "limit", AutomationConstants.DEFAULT_VECTOR_SEARCH_LIMIT, + "filters", Map.of(), "paramValues", Map.of()) + : orderedConfig("engineId", "", "value", "", "paramValues", Map.of()); List fields = new ArrayList<>(); fields.add(engineField(IEngine.CATALOG_TYPE.VECTOR)); fields.add(field("value", ConfigFieldType.TEXT, includeLimit ? "Search query" : "Records", true, "")); if (includeLimit) { fields.add(boundedIntegerField("limit", "Result limit", AutomationConstants.DEFAULT_VECTOR_SEARCH_LIMIT, 1, null)); + fields.add(field("filters", ConfigFieldType.JSON, "Filters", false, Map.of())); + } else if (nodeType == AutomationNodeType.VECTOR_ADD) { + defaultConfig.put("space", ""); + fields.add(field("space", ConfigFieldType.STRING, "Source space", false, "")); } + fields.add(field("paramValues", ConfigFieldType.JSON, "Engine parameters", false, Map.of())); return definition(nodeType, label, description, defaultConfig, fields, controlInputs(), controlAndResultOutputs()); } diff --git a/src/prerna/reactor/automation/AutomationNodeDefinition.java b/src/prerna/reactor/automation/definition/AutomationNodeDefinition.java similarity index 98% rename from src/prerna/reactor/automation/AutomationNodeDefinition.java rename to src/prerna/reactor/automation/definition/AutomationNodeDefinition.java index 2c3ee94c1c..cba80053ce 100644 --- a/src/prerna/reactor/automation/AutomationNodeDefinition.java +++ b/src/prerna/reactor/automation/definition/AutomationNodeDefinition.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.ArrayList; import java.util.Collections; @@ -63,7 +63,7 @@ public record AutomationNodeDefinition(AutomationNodeType nodeType, String label /** Supported configuration value shapes exposed to authoring clients. */ public enum ConfigFieldType { - ENGINE("engine"), STRING("string"), STRING_LIST("string[]"), TEXT("textarea"), CODE("code"), + BOOLEAN("boolean"), ENGINE("engine"), STRING("string"), STRING_LIST("string[]"), TEXT("textarea"), CODE("code"), INTEGER("number"), NUMBER("number"), JSON("json"), GLOBALS("globals"), BRANCH_CLAUSES("branch-clauses"); private final String value; diff --git a/src/prerna/reactor/automation/AutomationNodeType.java b/src/prerna/reactor/automation/definition/AutomationNodeType.java similarity index 98% rename from src/prerna/reactor/automation/AutomationNodeType.java rename to src/prerna/reactor/automation/definition/AutomationNodeType.java index 7f9fe5af37..7b994b0380 100644 --- a/src/prerna/reactor/automation/AutomationNodeType.java +++ b/src/prerna/reactor/automation/definition/AutomationNodeType.java @@ -25,13 +25,14 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.Collections; import java.util.LinkedHashMap; import java.util.Map; import prerna.engine.api.IEngine; +import prerna.reactor.automation.AutomationConstants; /** * Stable node identifiers and server-owned capabilities for Automation graphs. diff --git a/src/prerna/reactor/automation/AutomationScopeCatalog.java b/src/prerna/reactor/automation/definition/AutomationScopeCatalog.java similarity index 98% rename from src/prerna/reactor/automation/AutomationScopeCatalog.java rename to src/prerna/reactor/automation/definition/AutomationScopeCatalog.java index 07fc9b62bb..3f84cc4ad0 100644 --- a/src/prerna/reactor/automation/AutomationScopeCatalog.java +++ b/src/prerna/reactor/automation/definition/AutomationScopeCatalog.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.ArrayList; import java.util.HashMap; @@ -35,6 +35,9 @@ import java.util.Map; import java.util.Set; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; + /** * Describes the run-scope values that are guaranteed to exist when each node * executes. diff --git a/src/prerna/reactor/automation/AutomationSourceRenderer.java b/src/prerna/reactor/automation/definition/AutomationSourceRenderer.java similarity index 81% rename from src/prerna/reactor/automation/AutomationSourceRenderer.java rename to src/prerna/reactor/automation/definition/AutomationSourceRenderer.java index 3ab641a303..dba70e238a 100644 --- a/src/prerna/reactor/automation/AutomationSourceRenderer.java +++ b/src/prerna/reactor/automation/definition/AutomationSourceRenderer.java @@ -25,10 +25,12 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; +import java.util.List; import java.util.Map; +import prerna.reactor.automation.AutomationConstants; import prerna.reactor.automation.utils.AutomationRuntimeUtils; /** @@ -74,9 +76,9 @@ public static String renderNode(Map node) { case MODEL_NER -> modelNerSource(config); case STORAGE_LIST -> storageSource(config, "list", "STORAGE_PATH"); case STORAGE_READ -> storageReadSource(config); - case STORAGE_UPLOAD -> storageTransferSource(config, "copyToStorage"); + case STORAGE_UPLOAD -> storageUploadSource(config); case STORAGE_DOWNLOAD -> storageDownloadSource(config); - case STORAGE_DELETE -> storageSource(config, "deleteFromStorage", "STORAGE_PATH"); + case STORAGE_DELETE -> storageDeleteSource(config); case VECTOR_SEARCH -> vectorSearchSource(config); case VECTOR_ADD -> vectorAddSource(config); case VECTOR_DELETE -> vectorDeleteSource(config); @@ -136,10 +138,10 @@ raise ValueError("SEMOSS SQL query returned an invalid row shape.") def run(scope): query = scope.resolve(QUERY) pixel = "SqlQuery(" + ", ".join([ - _pixel_value("database", scope.resolve(ENGINE_ID)), - 'query=["' + query + '"]', - _pixel_value("limit", int(scope.resolve(LIMIT))), - ]) + ");" + _pixel_value("database", scope.resolve(ENGINE_ID)), + 'query=["' + query + '"]', + _pixel_value("limit", int(scope.resolve(LIMIT))), + ]) + ");" response = Insight().run_pixel(pixel, raw=True) result = response[0]["pixelReturn"][-1] if "ERROR" in result.get("operationType", []): @@ -212,6 +214,7 @@ private static String modelEmbeddingsSource(Map config) { ENGINE_ID = %s VALUES = %s + PARAMETERS_JSON = %s def _automation_model_response(result): if not isinstance(result, list) or len(result) != 1 or not isinstance(result[0], dict): @@ -223,9 +226,13 @@ raise ValueError("SEMOSS embeddings response is missing its vector payload.") def run(scope): model = ModelEngine(engine_id=scope.resolve(ENGINE_ID)) - result = model.embeddings(strings_to_embed=scope.resolve(VALUES)) + result = model.embeddings( + strings_to_embed=scope.resolve(VALUES), + param_dict=scope.resolve_config(PARAMETERS_JSON), + ) return _automation_model_response(result) - """.formatted(value(config, "engineId"), value(config, "text")); + """.formatted(value(config, "engineId"), value(config, "text"), + valueOrDefault(config, "paramValues", Map.of())); } private static String modelVisionSource(Map config) { @@ -235,7 +242,10 @@ private static String modelVisionSource(Map config) { ENGINE_ID = %s PROMPT = %s + SYSTEM_PROMPT = %s MEDIA = %s + URLS = %s + PARAMETERS_JSON = %s ROOM_ID = "${_automation_room_id}" RESULT_VALUE_KEY = %s RESULT_METADATA_KEY = %s @@ -254,22 +264,28 @@ raise ValueError("SEMOSS model response is missing content.") }, } - def _automation_media(value): + def _automation_values(value): values = value if isinstance(value, list) else [value] - media = [str(item).strip() for item in values if item is not None and str(item).strip()] - if not media: - raise ValueError("SEMOSS model media path is empty.") - return media + return [str(item).strip() for item in values if item is not None and str(item).strip()] def run(scope): model = ModelEngine(engine_id=scope.resolve(ENGINE_ID)) + media = _automation_values(scope.resolve(MEDIA)) + urls = _automation_values(scope.resolve(URLS)) + if not media and not urls: + raise ValueError("SEMOSS model media path or URL is required.") result = model.ask( command=scope.resolve(PROMPT), - media=_automation_media(scope.resolve(MEDIA)), + context=scope.resolve(SYSTEM_PROMPT), + media=media or None, + url=urls or None, + param_dict=scope.resolve_config(PARAMETERS_JSON), room_id=scope.resolve(ROOM_ID), ) return _automation_model_result(result) - """.formatted(value(config, "engineId"), value(config, "prompt"), value(config, "image"), + """.formatted(value(config, "engineId"), value(config, "prompt"), value(config, "systemPrompt"), + valueOrDefault(config, "image", ""), valueOrDefault(config, "urls", List.of()), + valueOrDefault(config, "paramValues", Map.of()), pythonValue(AutomationConstants.INTERNAL_RESULT_VALUE), pythonValue(AutomationConstants.INTERNAL_RESULT_METADATA)); } @@ -282,6 +298,7 @@ private static String modelNerSource(Map config) { ENGINE_ID = %s TEXT = %s ENTITIES = %s + MASK_ENTITIES = %s def _automation_model_response(result): if not isinstance(result, dict) or "response" not in result: @@ -293,9 +310,14 @@ raise ValueError(str(response.get("message") or "SEMOSS NER model execution fail def run(scope): model = ModelEngine(engine_id=scope.resolve(ENGINE_ID)) - result = model.ner(text=scope.resolve(TEXT), entities=scope.resolve(ENTITIES)) + result = model.ner( + text=scope.resolve(TEXT), + entities=scope.resolve(ENTITIES), + mask_entities=scope.resolve(MASK_ENTITIES), + ) return _automation_model_response(result) - """.formatted(value(config, "engineId"), value(config, "text"), value(config, "entities")); + """.formatted(value(config, "engineId"), value(config, "text"), value(config, "entities"), + valueOrDefault(config, "maskEntities", List.of())); } private static String storageSource(Map config, String method, String argument) { @@ -320,6 +342,7 @@ private static String storageReadSource(Map config) { ENGINE_ID = %s STORAGE_PATH = %s + CONVERT_TO_PDF = %s def _pixel_value(name, value): return name + "=[" + json.dumps(value) + "]" @@ -328,12 +351,14 @@ def run(scope): pixel = "GetStorageFileAsBase64(" + ", ".join([ _pixel_value("storage", scope.resolve(ENGINE_ID)), _pixel_value("storagePath", scope.resolve(STORAGE_PATH)), + _pixel_value("convertToPdf", bool(scope.resolve(CONVERT_TO_PDF))), ]) + ");" return Insight().run_pixel(pixel, raw=False) - """.formatted(value(config, "engineId"), value(config, "path")); + """.formatted(value(config, "engineId"), value(config, "path"), + valueOrDefault(config, "convertToPdf", false)); } - private static String storageTransferSource(Map config, String method) { + private static String storageUploadSource(Map config) { return """ # Transfer files with SEMOSS storage through the Python SDK. from ai_server import StorageEngine @@ -341,11 +366,36 @@ private static String storageTransferSource(Map config, String m ENGINE_ID = %s STORAGE_PATH = %s FILE_PATH = %s + METADATA = %s def run(scope): storage = StorageEngine(engine_id=scope.resolve(ENGINE_ID)) - return storage.%s(storagePath=scope.resolve(STORAGE_PATH), localPath=scope.resolve(FILE_PATH)) - """.formatted(value(config, "engineId"), value(config, "path"), value(config, "destination"), method); + return storage.copyToStorage( + storagePath=scope.resolve(STORAGE_PATH), + localPath=scope.resolve(FILE_PATH), + metadata=scope.resolve_config(METADATA), + ) + """.formatted(value(config, "engineId"), value(config, "path"), value(config, "destination"), + valueOrDefault(config, "metadata", Map.of())); + } + + private static String storageDeleteSource(Map config) { + return """ + # Delete storage content through the SEMOSS Python SDK. + from ai_server import StorageEngine + + ENGINE_ID = %s + STORAGE_PATH = %s + LEAVE_FOLDER_STRUCTURE = %s + + def run(scope): + storage = StorageEngine(engine_id=scope.resolve(ENGINE_ID)) + return storage.deleteFromStorage( + storagePath=scope.resolve(STORAGE_PATH), + leaveFolderStructure=bool(scope.resolve(LEAVE_FOLDER_STRUCTURE)), + ) + """.formatted(value(config, "engineId"), value(config, "path"), + valueOrDefault(config, "leaveFolderStructure", false)); } private static String storageDownloadSource(Map config) { @@ -358,6 +408,7 @@ private static String storageDownloadSource(Map config) { ENGINE_ID = %s STORAGE_PATH = %s DESTINATION = %s + VERSION = %s def _pixel_value(name, value): return name + "=[" + json.dumps(value) + "]" @@ -383,10 +434,20 @@ def run(scope): storage = StorageEngine(engine_id=scope.resolve(ENGINE_ID)) storage_path = str(scope.resolve(STORAGE_PATH)) destination = str(scope.resolve(DESTINATION) or "").strip() or "/" - copied = storage.copyToLocal( - storagePath=storage_path, - localPath=destination, - ) + version = str(scope.resolve(VERSION) or "").strip() + if version: + pixel = "PullFromStorage(" + ", ".join([ + _pixel_value("storage", scope.resolve(ENGINE_ID)), + _pixel_value("storagePath", storage_path), + _pixel_value("filePath", destination), + _pixel_value("version", version), + ]) + ");" + copied = Insight().run_pixel(pixel, raw=False) + else: + copied = storage.copyToLocal( + storagePath=storage_path, + localPath=destination, + ) files = _destination_files(destination) return { "success": copied, @@ -396,7 +457,8 @@ def run(scope): "files": files, "filePath": files[0] if len(files) == 1 else None, } - """.formatted(value(config, "engineId"), value(config, "path"), value(config, "destination")); + """.formatted(value(config, "engineId"), value(config, "path"), value(config, "destination"), + valueOrDefault(config, "version", "")); } private static String vectorSearchSource(Map config) { @@ -407,12 +469,20 @@ private static String vectorSearchSource(Map config) { ENGINE_ID = %s QUERY = %s LIMIT = %s + FILTERS = %s + PARAMETERS_JSON = %s def run(scope): vector = VectorEngine(engine_id=scope.resolve(ENGINE_ID)) - return vector.nearestNeighbor(search_statement=scope.resolve(QUERY), limit=scope.resolve(LIMIT)) + return vector.nearestNeighbor( + search_statement=scope.resolve(QUERY), + limit=scope.resolve(LIMIT), + filters=scope.resolve_config(FILTERS), + param_dict=scope.resolve_config(PARAMETERS_JSON), + ) """.formatted(value(config, "engineId"), value(config, "value"), - valueOrDefault(config, "limit", AutomationConstants.DEFAULT_VECTOR_SEARCH_LIMIT)); + valueOrDefault(config, "limit", AutomationConstants.DEFAULT_VECTOR_SEARCH_LIMIT), + valueOrDefault(config, "filters", Map.of()), valueOrDefault(config, "paramValues", Map.of())); } private static String vectorAddSource(Map config) { @@ -422,11 +492,23 @@ private static String vectorAddSource(Map config) { ENGINE_ID = %s FILE_PATHS = %s + SPACE = %s + PARAMETERS_JSON = %s + + def _automation_list(value): + if isinstance(value, list): + return value + return [item.strip() for item in str(value).split(",") if item.strip()] def run(scope): vector = VectorEngine(engine_id=scope.resolve(ENGINE_ID)) - return vector.addDocument(file_paths=scope.resolve(FILE_PATHS).split(",")) - """.formatted(value(config, "engineId"), value(config, "value")); + return vector.addDocument( + file_paths=_automation_list(scope.resolve(FILE_PATHS)), + space=str(scope.resolve(SPACE) or "").strip() or None, + param_dict=scope.resolve_config(PARAMETERS_JSON), + ) + """.formatted(value(config, "engineId"), value(config, "value"), valueOrDefault(config, "space", ""), + valueOrDefault(config, "paramValues", Map.of())); } private static String vectorDeleteSource(Map config) { @@ -436,11 +518,21 @@ private static String vectorDeleteSource(Map config) { ENGINE_ID = %s FILE_NAMES = %s + PARAMETERS_JSON = %s + + def _automation_list(value): + if isinstance(value, list): + return value + return [item.strip() for item in str(value).split(",") if item.strip()] def run(scope): vector = VectorEngine(engine_id=scope.resolve(ENGINE_ID)) - return vector.removeDocument(file_names=scope.resolve(FILE_NAMES).split(",")) - """.formatted(value(config, "engineId"), value(config, "value")); + return vector.removeDocument( + file_names=_automation_list(scope.resolve(FILE_NAMES)), + param_dict=scope.resolve_config(PARAMETERS_JSON), + ) + """.formatted(value(config, "engineId"), value(config, "value"), + valueOrDefault(config, "paramValues", Map.of())); } private static String functionSource(Map config) { diff --git a/src/prerna/reactor/automation/GetAutomationNodeDefinitionsReactor.java b/src/prerna/reactor/automation/definition/GetAutomationNodeDefinitionsReactor.java similarity index 97% rename from src/prerna/reactor/automation/GetAutomationNodeDefinitionsReactor.java rename to src/prerna/reactor/automation/definition/GetAutomationNodeDefinitionsReactor.java index 7fb62c9e83..9d3580536e 100644 --- a/src/prerna/reactor/automation/GetAutomationNodeDefinitionsReactor.java +++ b/src/prerna/reactor/automation/definition/GetAutomationNodeDefinitionsReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import prerna.reactor.AbstractReactor; import prerna.sablecc2.om.PixelDataType; diff --git a/src/prerna/reactor/automation/GetAutomationReactor.java b/src/prerna/reactor/automation/definition/GetAutomationReactor.java similarity index 92% rename from src/prerna/reactor/automation/GetAutomationReactor.java rename to src/prerna/reactor/automation/definition/GetAutomationReactor.java index 07e8cf8743..f28a77c9a8 100644 --- a/src/prerna/reactor/automation/GetAutomationReactor.java +++ b/src/prerna/reactor/automation/definition/GetAutomationReactor.java @@ -25,12 +25,15 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.LinkedHashMap; import java.util.Map; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; @@ -54,7 +57,7 @@ public GetAutomationReactor() { @Override public NounMetadata execute() { organizeKeys(); - String projectId = AutomationProjectUtils.getViewableAutomationProject(this.insight.getUser(), + String projectId = AutomationProjectService.getViewableAutomationProject(this.insight.getUser(), this.keyValue.get(ReactorKeysEnum.PROJECT.getKey())).getProjectId(); AutomationDefinitionService.DefinitionFiles files = AutomationDefinitionService.load(projectId); diff --git a/src/prerna/reactor/automation/RemoveAutomationStepReactor.java b/src/prerna/reactor/automation/definition/RemoveAutomationStepReactor.java similarity index 96% rename from src/prerna/reactor/automation/RemoveAutomationStepReactor.java rename to src/prerna/reactor/automation/definition/RemoveAutomationStepReactor.java index 11b89151b1..3e7139c5a6 100644 --- a/src/prerna/reactor/automation/RemoveAutomationStepReactor.java +++ b/src/prerna/reactor/automation/definition/RemoveAutomationStepReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.ArrayDeque; import java.util.ArrayList; @@ -37,6 +37,8 @@ import prerna.reactor.AbstractReactor; import prerna.reactor.agent.mcp.MCPUtility; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; @@ -64,12 +66,12 @@ public RemoveAutomationStepReactor() { @Override public NounMetadata execute() { organizeKeys(); - String projectId = AutomationProjectUtils + String projectId = AutomationProjectService .getEditableAutomationProject(this.insight.getUser(), required(ReactorKeysEnum.PROJECT.getKey())) .getProjectId(); String nodeId = required(NODE_ID_KEY); boolean removeDownstream = optionalBoolean(REMOVE_DOWNSTREAM_KEY, false); - return AutomationProjectUtils.withLockedDefinition(projectId, + return AutomationProjectService.withLockedDefinition(projectId, files -> removeStep(projectId, files, nodeId, removeDownstream)); } @@ -108,7 +110,7 @@ private NounMetadata removeStep(String projectId, AutomationDefinitionService.De Map updatedSources = new LinkedHashMap<>(files.nodeSources()); removedNodeIds.forEach(updatedSources::remove); - AutomationDefinitionService.DefinitionFiles saved = AutomationProjectUtils.saveDefinition(projectId, + AutomationDefinitionService.DefinitionFiles saved = AutomationProjectService.saveDefinition(projectId, AutomationRuntimeUtils.GSON.toJson(updatedDefinition), updatedSources, this.insight.getUser()); Map result = new LinkedHashMap<>(); result.put("removed", true); diff --git a/src/prerna/reactor/automation/SaveAutomationReactor.java b/src/prerna/reactor/automation/definition/SaveAutomationReactor.java similarity index 93% rename from src/prerna/reactor/automation/SaveAutomationReactor.java rename to src/prerna/reactor/automation/definition/SaveAutomationReactor.java index 9c337371a4..94dc38d091 100644 --- a/src/prerna/reactor/automation/SaveAutomationReactor.java +++ b/src/prerna/reactor/automation/definition/SaveAutomationReactor.java @@ -25,12 +25,15 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.LinkedHashMap; import java.util.Map; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; @@ -66,10 +69,10 @@ public NounMetadata execute() { throw new IllegalArgumentException("Must provide a project id."); } - projectId = AutomationProjectUtils.getEditableAutomationProject(this.insight.getUser(), projectId) + projectId = AutomationProjectService.getEditableAutomationProject(this.insight.getUser(), projectId) .getProjectId(); - AutomationDefinitionService.DefinitionFiles files = AutomationProjectUtils.saveDefinition(projectId, definition, + AutomationDefinitionService.DefinitionFiles files = AutomationProjectService.saveDefinition(projectId, definition, nodeSources, this.keyValue.get(AutomationConstants.EXPECTED_REVISION_KEY), this.insight.getUser()); Map result = new LinkedHashMap<>(); diff --git a/src/prerna/reactor/automation/UpdateAutomationCustomStepReactor.java b/src/prerna/reactor/automation/definition/UpdateAutomationCustomStepReactor.java similarity index 95% rename from src/prerna/reactor/automation/UpdateAutomationCustomStepReactor.java rename to src/prerna/reactor/automation/definition/UpdateAutomationCustomStepReactor.java index 405f65099d..ba0d19ac74 100644 --- a/src/prerna/reactor/automation/UpdateAutomationCustomStepReactor.java +++ b/src/prerna/reactor/automation/definition/UpdateAutomationCustomStepReactor.java @@ -25,13 +25,15 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; @@ -65,7 +67,7 @@ public NounMetadata execute() { String nodeId = required(NODE_ID_KEY); String source = required(SOURCE_KEY); String expectedHash = required(EXPECTED_SOURCE_HASH_KEY); - return AutomationProjectUtils.withLockedDefinition(projectId, + return AutomationProjectService.withLockedDefinition(projectId, files -> updateSource(projectId, files, nodeId, source, expectedHash)); } @@ -79,7 +81,7 @@ private NounMetadata updateSource(String projectId, AutomationDefinitionService. Map updatedSources = new LinkedHashMap<>(files.nodeSources()); updatedSources.put(nodeId, source); - AutomationDefinitionService.DefinitionFiles saved = AutomationProjectUtils.saveDefinition(projectId, + AutomationDefinitionService.DefinitionFiles saved = AutomationProjectService.saveDefinition(projectId, files.definition(), updatedSources, this.insight.getUser()); Map result = new LinkedHashMap<>(); result.put("nodeId", nodeId); @@ -90,7 +92,7 @@ private NounMetadata updateSource(String projectId, AutomationDefinitionService. } private String editableProjectId() { - return AutomationProjectUtils + return AutomationProjectService .getEditableAutomationProject(this.insight.getUser(), required(ReactorKeysEnum.PROJECT.getKey())) .getProjectId(); } diff --git a/src/prerna/reactor/automation/UpdateAutomationStepReactor.java b/src/prerna/reactor/automation/definition/UpdateAutomationStepReactor.java similarity index 95% rename from src/prerna/reactor/automation/UpdateAutomationStepReactor.java rename to src/prerna/reactor/automation/definition/UpdateAutomationStepReactor.java index 0982975809..883241b5c2 100644 --- a/src/prerna/reactor/automation/UpdateAutomationStepReactor.java +++ b/src/prerna/reactor/automation/definition/UpdateAutomationStepReactor.java @@ -25,13 +25,15 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; @@ -64,7 +66,7 @@ public NounMetadata execute() { String nodeId = required(NODE_ID_KEY); Map config = AutomationRuntimeUtils.parseJsonObject(required(CONFIG_KEY), "config"); String label = this.keyValue.get(LABEL_KEY); - return AutomationProjectUtils.withLockedDefinition(projectId, + return AutomationProjectService.withLockedDefinition(projectId, files -> updateStep(projectId, files, nodeId, config, label)); } @@ -107,7 +109,7 @@ private NounMetadata updateStep(String projectId, AutomationDefinitionService.De Map sources = new LinkedHashMap<>(files.nodeSources()); sources.remove(nodeId); - AutomationDefinitionService.DefinitionFiles saved = AutomationProjectUtils.saveDefinition(projectId, + AutomationDefinitionService.DefinitionFiles saved = AutomationProjectService.saveDefinition(projectId, AutomationRuntimeUtils.GSON.toJson(document), sources, this.insight.getUser()); Map result = new LinkedHashMap<>(); result.put("node", updatedNode); @@ -118,7 +120,7 @@ private NounMetadata updateStep(String projectId, AutomationDefinitionService.De } private String editableProjectId() { - return AutomationProjectUtils + return AutomationProjectService .getEditableAutomationProject(this.insight.getUser(), required(ReactorKeysEnum.PROJECT.getKey())) .getProjectId(); } diff --git a/src/prerna/reactor/automation/project/AGENTS.md b/src/prerna/reactor/automation/project/AGENTS.md new file mode 100644 index 0000000000..bef87c00be --- /dev/null +++ b/src/prerna/reactor/automation/project/AGENTS.md @@ -0,0 +1,27 @@ +# Automation Project Package + +This package owns the Automation project boundary. Follow the parent +`../AGENTS.md` for the complete Automation contract. + +## Ownership + +- Resolve Automation projects through standard SEMOSS project permissions. +- Serialize definition mutations under the project lock. +- Coordinate definition persistence, reference validation, edit timestamps, + cluster synchronization, and derived MCP assets. +- Create new Automation projects and their starter definition. + +## Invariants + +- View and edit operations must use the corresponding project ACL check. +- Keep graph and node-source publication atomic from the caller's perspective. +- MCP files are derived project assets, not a second workflow definition. +- Validate referenced engines, apps, and agent workspaces against the acting + user before publishing the definition. +- Keep `AutomationProjectService` concrete; split it only when a responsibility + gains an independently owned lifecycle. + +## Verification + +Test permission failures, revision conflicts, rollback behavior, reference +validation, and MCP regeneration when this boundary changes. diff --git a/src/prerna/reactor/automation/AutomationMcpSync.java b/src/prerna/reactor/automation/project/AutomationMcpSync.java similarity index 97% rename from src/prerna/reactor/automation/AutomationMcpSync.java rename to src/prerna/reactor/automation/project/AutomationMcpSync.java index 5dae3fdc95..4cf1e70162 100644 --- a/src/prerna/reactor/automation/AutomationMcpSync.java +++ b/src/prerna/reactor/automation/project/AutomationMcpSync.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.project; import java.nio.charset.StandardCharsets; import java.nio.file.Files; @@ -44,9 +44,17 @@ import prerna.auth.User; import prerna.engine.api.IEngine; import prerna.project.api.IProject; -import prerna.reactor.agent.mcp.MCPUtility; import prerna.reactor.agent.mcp.MCPUtility.MCPDisplayOption; import prerna.reactor.agent.mcp.MCPUtility.MCPExecution; +import prerna.reactor.agent.mcp.MCPUtility; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; +import prerna.reactor.automation.definition.AutomationDefinitionService; +import prerna.reactor.automation.definition.AutomationDefinitionValidator; +import prerna.reactor.automation.definition.AutomationNodeCatalog; +import prerna.reactor.automation.definition.AutomationNodeDefinition; +import prerna.reactor.automation.definition.AutomationNodeType; +import prerna.reactor.automation.definition.GetAutomationNodeDefinitionsReactor; import prerna.reactor.function.GetFunctionEngineDefinitionReactor; import prerna.reactor.project.GetProjectAvailableReactorsReactor; import prerna.reactor.project.GetProjectReactorSignatureReactor; diff --git a/src/prerna/reactor/automation/AutomationProjectUtils.java b/src/prerna/reactor/automation/project/AutomationProjectService.java similarity index 97% rename from src/prerna/reactor/automation/AutomationProjectUtils.java rename to src/prerna/reactor/automation/project/AutomationProjectService.java index d60a31cb5a..5b207ff3f3 100644 --- a/src/prerna/reactor/automation/AutomationProjectUtils.java +++ b/src/prerna/reactor/automation/project/AutomationProjectService.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.project; import java.nio.file.Files; import java.nio.file.Path; @@ -45,6 +45,10 @@ import prerna.engine.api.IEngine; import prerna.engine.impl.model.inferencetracking.ModelInferenceLogsUtils; import prerna.project.api.IProject; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.definition.AutomationDefinitionService; +import prerna.reactor.automation.definition.AutomationDefinitionValidator; +import prerna.reactor.automation.definition.AutomationNodeType; import prerna.util.AssetUtility; import prerna.util.ProjectSyncUtility; import prerna.util.Utility; @@ -60,11 +64,11 @@ * regenerating MCP metadata, synchronizing cluster assets, and updating the * project edit timestamp. */ -public final class AutomationProjectUtils { +public final class AutomationProjectService { - private static final Logger classLogger = LogManager.getLogger(AutomationProjectUtils.class); + private static final Logger classLogger = LogManager.getLogger(AutomationProjectService.class); - private AutomationProjectUtils() { + private AutomationProjectService() { } /** diff --git a/src/prerna/reactor/automation/CreateAutomationReactor.java b/src/prerna/reactor/automation/project/CreateAutomationReactor.java similarity index 98% rename from src/prerna/reactor/automation/CreateAutomationReactor.java rename to src/prerna/reactor/automation/project/CreateAutomationReactor.java index efd32c856e..c1b37f2f3f 100644 --- a/src/prerna/reactor/automation/CreateAutomationReactor.java +++ b/src/prerna/reactor/automation/project/CreateAutomationReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.project; import java.io.IOException; import java.util.Map; @@ -94,7 +94,7 @@ public NounMetadata execute() { null, user, classLogger); String projectId = project.getProjectId(); try { - AutomationProjectUtils.createStarterDefinition(project, user); + AutomationProjectService.createStarterDefinition(project, user); } catch (RuntimeException e) { classLogger.error("Failed to scaffold automation project '{}' (id {})", projectName, projectId, e); try { diff --git a/src/prerna/reactor/automation/run/AGENTS.md b/src/prerna/reactor/automation/run/AGENTS.md new file mode 100644 index 0000000000..5332c1553e --- /dev/null +++ b/src/prerna/reactor/automation/run/AGENTS.md @@ -0,0 +1,31 @@ +# Automation Run Package + +This package owns durable run state and the lifecycle of one Automation +execution. Follow the parent `../AGENTS.md` for the complete Automation +contract. + +## Ownership + +- `AutomationRunStore` is the concrete scheduler-database persistence owner. +- `AutomationRunExecutionService` claims and executes an immutable run + snapshot in one run-owned Insight. +- `AutomationRunRegistry` is only the same-JVM cancellation and heartbeat fast + path; the database remains authoritative. +- Run reactors trigger, inspect, list, resume, and cancel executions. + +## Invariants + +- Preserve project permission checks and the atomic submitted-to-running claim. +- Execute only the definition, inputs, and source captured for that run. +- Keep node transitions, cancellation, waits, and terminal status durable. +- Keep frame display data scoped to the execution Insight. A frame is not a + durable-history contract. +- Do not pass engine objects, arbitrary Java objects, or unbounded payloads + through Python scope or browser responses. +- Do not add a persistence interface until there is a real second store. + +## Verification + +Keep focused tests under `test/prerna/reactor/automation/run`. Exercise claim +behavior, status transitions, cancellation, waits/resume, run snapshots, and +cleanup whenever their lifecycle changes. diff --git a/src/prerna/reactor/automation/AutomationRunExecutionService.java b/src/prerna/reactor/automation/run/AutomationRunExecutionService.java similarity index 90% rename from src/prerna/reactor/automation/AutomationRunExecutionService.java rename to src/prerna/reactor/automation/run/AutomationRunExecutionService.java index d17a1ab19b..fc7eb019a8 100644 --- a/src/prerna/reactor/automation/AutomationRunExecutionService.java +++ b/src/prerna/reactor/automation/run/AutomationRunExecutionService.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import java.io.File; import java.sql.Timestamp; @@ -42,26 +42,34 @@ import java.util.concurrent.TimeUnit; import java.util.function.LongSupplier; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - import com.google.re2j.Matcher; import com.google.re2j.Pattern; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; import prerna.auth.utils.SecurityEngineUtils; import prerna.ds.py.PyTranslator; import prerna.engine.api.IEngine; import prerna.engine.api.IModelEngine; import prerna.engine.api.ITypeSafeEngine; -import prerna.engine.impl.model.responses.TypeSafeModelEngineResponse; import prerna.engine.impl.model.RoomUtils; +import prerna.engine.impl.model.responses.TypeSafeModelEngineResponse; import prerna.om.Insight; import prerna.om.InsightStore; import prerna.om.ThreadStore; import prerna.project.api.IProject; import prerna.reactor.agent.run.AgentRunService; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; +import prerna.reactor.automation.definition.AutomationConditionEvaluator; +import prerna.reactor.automation.definition.AutomationDefinitionValidator; import prerna.reactor.automation.utils.AutomationRuntimeUtils; +import prerna.reactor.frame.py.GenerateFrameFromPyVariableReactor; import prerna.sablecc2.comm.PixelJobManager; +import prerna.sablecc2.om.NounStore; +import prerna.sablecc2.om.PixelDataType; +import prerna.sablecc2.om.ReactorKeysEnum; +import prerna.sablecc2.om.nounmeta.NounMetadata; import prerna.util.EngineUtility; import prerna.util.Utility; import prerna.util.insight.InsightUtility; @@ -127,7 +135,7 @@ final class AutomationRunExecutionService { Map executeInitializedRun(String runId, String projectId, AutomationDefinitionValidator.ValidatedDefinition definition, List> runNodes, Map traceRoomIds) { - if (!AutomationDatabaseUtility.claimRun(runId)) { + if (!AutomationRunStore.claimRun(runId)) { return buildCurrentRunResult(runId, projectId); } streamRunStarted(runId, definition); @@ -142,11 +150,11 @@ Map executeInitializedRun(String runId, String projectId, if (translator == null) { throw new IllegalStateException("Python runtime is not available for this insight."); } - AutomationPythonRunRegistry.register(runId, translator, executionInsight, streamJobId); + AutomationRunRegistry.register(runId, translator, executionInsight, streamJobId); Map scope = AutomationRuntimeUtils.buildInitialScope(runId, executionInsight.getUser()); - scope.putAll(AutomationDatabaseUtility.getRunInputs(runId)); - Map runNodeSources = AutomationDatabaseUtility.getRunNodeSources(runId); + scope.putAll(AutomationRunStore.getRunInputs(runId)); + Map runNodeSources = AutomationRunStore.getRunNodeSources(runId); result = executeInControlOrder(executionInsight, projectId, runId, definition, runNodes, runNodeSources, scope, traceRoomIds, AutomationRuntime.startNodeId(definition)); if (!Boolean.TRUE.equals(result.get("waitingForInput"))) { @@ -157,7 +165,7 @@ Map executeInitializedRun(String runId, String projectId, finishFailedRun(runId, projectId, e); result = Map.of("error", safeMessage(e)); } finally { - AutomationPythonRunRegistry.unregister(runId); + AutomationRunRegistry.unregister(runId); if (executionInsightLease == null || executionInsightLease.cleanupOnCompletion()) { cleanupExecutionInsight(executionInsight); } @@ -172,19 +180,19 @@ Map executeInitializedRun(String runId, String projectId, * rather than an error. */ private static Map buildCurrentRunResult(String runId, String projectId) { - Map persisted = AutomationDatabaseUtility.getRunDetail(runId); + Map persisted = AutomationRunStore.getRunDetail(runId); if (persisted == null) { throw new IllegalStateException( "Automation run '" + runId + "' no longer exists for project '" + projectId + "'."); } Map result = new LinkedHashMap<>(persisted); result.put(AutomationConstants.RESULT_NODE_RESULTS, - AutomationDatabaseUtility.buildNodeResults(AutomationDatabaseUtility.getNodeOutputsForRun(runId))); + AutomationRunStore.buildNodeResults(AutomationRunStore.getNodeOutputsForRun(runId))); String executionInsightId = getAvailableExecutionInsightId(runId); if (executionInsightId != null) { result.put(AutomationConstants.RESULT_EXECUTION_INSIGHT_ID, executionInsightId); } - Map wait = AutomationDatabaseUtility.getActiveWait(runId); + Map wait = AutomationRunStore.getActiveWait(runId); if (wait != null) { result.put("wait", wait); } @@ -216,7 +224,7 @@ private Map executeInControlOrder(Insight executionInsight, Stri String currentNodeId = initialNodeId; boolean pathCompleted = true; while (currentNodeId != null) { - if (AutomationPythonRunRegistry.isCancellationRequested(runId)) { + if (AutomationRunRegistry.isCancellationRequested(runId)) { pathCompleted = false; break; } @@ -270,7 +278,7 @@ private Map executeInControlOrder(Insight executionInsight, Stri currentNodeId = controlTargets.getOrDefault(nodeId, Map.of()).get(selectedPort); } if (pathCompleted) { - AutomationDatabaseUtility.skipPendingNodes(runId, "Control branch was not selected"); + AutomationRunStore.skipPendingNodes(runId, "Control branch was not selected"); } result.put("scope", scope); return result; @@ -287,7 +295,7 @@ private Map executeJevDecisionNode(Insight executionInsight, Str String nodeId = (String) node.get(AutomationConstants.NODE_FIELD_ID); Timestamp started = Utility.getSqlTimestampUTC(LocalDateTime.ofInstant(Instant.now(), ZoneOffset.UTC)); long startedMs = System.currentTimeMillis(); - AutomationDatabaseUtility.markNodeRunning(runId, nodeId); + AutomationRunStore.markNodeRunning(runId, nodeId); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_RUNNING, null, null, null); try { Map config = (Map) node.get(AutomationConstants.NODE_FIELD_CONFIG); @@ -328,15 +336,15 @@ private Map executeJevDecisionNode(Insight executionInsight, Str AutomationConstants.NODE_OUTPUT_MAX_BYTES, "Automation Jev decision '" + nodeId + "' output"); long duration = System.currentTimeMillis() - startedMs; String preview = AutomationRuntimeUtils.generatePreview(output); - AutomationDatabaseUtility.updateNodeSuccess(runId, nodeId, started, duration, null, output, preview, null, + AutomationRunStore.updateNodeSuccess(runId, nodeId, started, duration, null, output, preview, null, null); - AutomationPythonRunRegistry.nodeCompleted(runId); + AutomationRunRegistry.nodeCompleted(runId); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_SUCCESS, duration, preview, null); return nodeResult(nodeId, AutomationConstants.NODE_STATUS_SUCCESS, decision, null); } catch (Exception e) { long duration = System.currentTimeMillis() - startedMs; String message = safeMessage(e); - AutomationDatabaseUtility.updateNodeFailed(runId, nodeId, started, duration, message); + AutomationRunStore.updateNodeFailed(runId, nodeId, started, duration, message); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_FAILED, duration, null, message); throw e instanceof RuntimeException runtimeException ? runtimeException : new RuntimeException(e); } @@ -471,7 +479,7 @@ private Map executeConditionNode(String runId, Map config = (Map) node.get(AutomationConstants.NODE_FIELD_CONFIG); @@ -493,15 +501,15 @@ private Map executeConditionNode(String runId, Map executeStartNode(Insight executionInsight, String pr String nodeId = (String) node.get(AutomationConstants.NODE_FIELD_ID); Timestamp started = Utility.getSqlTimestampUTC(LocalDateTime.ofInstant(Instant.now(), ZoneOffset.UTC)); long startedMs = System.currentTimeMillis(); - AutomationDatabaseUtility.markNodeRunning(runId, nodeId); + AutomationRunStore.markNodeRunning(runId, nodeId); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_RUNNING, null, null, null); try { Map declaredGlobals = AutomationRuntime.triggerGlobalDefaults(node); @@ -551,9 +559,9 @@ private Map executeStartNode(Insight executionInsight, String pr "Automation run scope"); long duration = System.currentTimeMillis() - startedMs; String preview = AutomationRuntimeUtils.generatePreview(output); - AutomationDatabaseUtility.updateNodeSuccess(runId, nodeId, started, duration, null, output, preview, null, + AutomationRunStore.updateNodeSuccess(runId, nodeId, started, duration, null, output, preview, null, null); - AutomationPythonRunRegistry.nodeCompleted(runId); + AutomationRunRegistry.nodeCompleted(runId); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_SUCCESS, duration, preview, null); Map result = nodeResult(nodeId, AutomationConstants.NODE_STATUS_SUCCESS, output, null); result.put(AutomationConstants.RESULT_GLOBALS, globals); @@ -561,7 +569,7 @@ private Map executeStartNode(Insight executionInsight, String pr } catch (Exception e) { long duration = System.currentTimeMillis() - startedMs; String message = safeMessage(e); - AutomationDatabaseUtility.updateNodeFailed(runId, nodeId, started, duration, message); + AutomationRunStore.updateNodeFailed(runId, nodeId, started, duration, message); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_FAILED, duration, null, message); throw e instanceof RuntimeException runtimeException ? runtimeException : new RuntimeException(e); } @@ -590,33 +598,64 @@ private Map executeNodeSource(Insight executionInsight, String p String nodeId = (String) node.get(AutomationConstants.NODE_FIELD_ID); Timestamp started = Utility.getSqlTimestampUTC(LocalDateTime.ofInstant(Instant.now(), ZoneOffset.UTC)); long startedMs = System.currentTimeMillis(); - AutomationDatabaseUtility.markNodeRunning(runId, nodeId); + AutomationRunStore.markNodeRunning(runId, nodeId); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_RUNNING, null, null, null, traceForNode(node, traceRoomId, null, null)); try { prepareGeneratedAgentRoom(executionInsight, node, traceRoomId); + String outputVariable = (String) node.get(AutomationConstants.NODE_FIELD_OUTPUT_VAR); Map nodeScope = scope; if (traceRoomId != null) { nodeScope = new LinkedHashMap<>(scope); nodeScope.put(AutomationConstants.SCOPE_ROOM_ID, traceRoomId); } Object raw = translator.runScriptWithExplicitAssetPaths(executionInsight, - AutomationRuntime.buildNodeInvocationScript(source, nodeScope), getProjectAssetsFolder(projectId), + AutomationRuntime.buildNodeInvocationScript(source, nodeScope, outputVariable), + getProjectAssetsFolder(projectId), new String[] { getProjectPyFolder(projectId) }); Object value = AutomationRuntime.normalizeNodeResult(raw); + registerNodeFrame(executionInsight, outputVariable, value); value = awaitGeneratedAgentRun(executionInsight, runId, node, value, traceRoomId, scope); return persistNativeNodeResult(runId, projectId, node, value, started, startedMs, traceRoomId, resumeNodeId, scope); } catch (Exception e) { long duration = System.currentTimeMillis() - startedMs; String message = safeMessage(e); - AutomationDatabaseUtility.updateNodeFailed(runId, nodeId, started, duration, message); + AutomationRunStore.updateNodeFailed(runId, nodeId, started, duration, message); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_FAILED, duration, null, message, traceForNode(node, traceRoomId, null, null)); throw e instanceof RuntimeException runtimeException ? runtimeException : new RuntimeException(e); } } + /** + * Promotes row-shaped output through the same Python-frame reactor used by + * Notebook. The bounded JSON value remains the Automation contract; this named + * frame only lets the UI inspect it through standard paged frame queries. + */ + private static void registerNodeFrame(Insight executionInsight, String outputVariable, Object value) { + if (!(value instanceof List rows) || rows.isEmpty() + || rows.stream().anyMatch(row -> !(row instanceof Map))) { + return; + } + + NounStore nounStore = new NounStore("GenerateFrameFromPyVariable"); + nounStore.makeGenRowStruct(ReactorKeysEnum.VARIABLE.getKey()) + .add(new NounMetadata(outputVariable, PixelDataType.CONST_STRING)); + nounStore.makeGenRowStruct(ReactorKeysEnum.OVERRIDE.getKey()) + .add(new NounMetadata(false, PixelDataType.BOOLEAN)); + + GenerateFrameFromPyVariableReactor reactor = new GenerateFrameFromPyVariableReactor(); + reactor.setInsight(executionInsight); + reactor.setNounStore(nounStore); + try { + reactor.execute(); + } catch (RuntimeException e) { + classLogger.warn("Unable to register tabular output '{}' in Automation execution insight '{}'", + outputVariable, executionInsight.getInsightId(), e); + } + } + /** * The validator treats every codeMode except the literal "custom" as generated, * including an absent value. Matching that convention here (rather than @@ -682,7 +721,7 @@ Object awaitGeneratedAgentRun(Insight executionInsight, String automationRunId, throw missingAgentTrace(node); } long waitStartedAt = monotonicTime.getAsLong(); - AutomationDatabaseUtility.updateNodeAgentRunTrace(automationRunId, + AutomationRunStore.updateNodeAgentRunTrace(automationRunId, (String) node.get(AutomationConstants.NODE_FIELD_ID), agentRunId); String previousStatus = stringValue(submittedRun.get("status")); @@ -693,7 +732,7 @@ Object awaitGeneratedAgentRun(Insight executionInsight, String automationRunId, streamNodeProgress(automationRunId, node, AutomationConstants.NODE_STATUS_RUNNING, null, null, null, traceForNode(node, traceRoomId, null, agentRunId, previousStatus)); while (true) { - if (AutomationPythonRunRegistry.isCancellationRequested(automationRunId) && !cancellationSignalled) { + if (AutomationRunRegistry.isCancellationRequested(automationRunId) && !cancellationSignalled) { cancellationSignalled = true; AgentRunService.get().cancelRun(agentRunId, "Automation run cancelled"); } @@ -979,6 +1018,17 @@ static String getAvailableExecutionInsightId(String runId) { return InsightStore.getInstance().containsKey(executionInsightId) ? executionInsightId : null; } + /** + * Returns the live run-owned Insight without exposing its deterministic naming + * convention to callers. + * + * @param runId durable run identifier + * @return live execution Insight, or {@code null} after its workspace closes + */ + static Insight getAvailableExecutionInsight(String runId) { + return InsightStore.getInstance().get(executionInsightId(runId)); + } + private static String executionInsightId(String runId) { return EXECUTION_INSIGHT_PREFIX + runId; } @@ -1048,10 +1098,10 @@ private Map persistCancelledNodeResult(String runId, Map persistNativeNodeResult(String runId, String project Object value, Timestamp started, long startedMs, String traceRoomId, String resumeNodeId, Map scope) { String nodeId = (String) node.get(AutomationConstants.NODE_FIELD_ID); - if (AutomationPythonRunRegistry.isCancellationRequested(runId)) { + if (AutomationRunRegistry.isCancellationRequested(runId)) { return persistCancelledNodeResult(runId, node, nodeId, value, started, startedMs, traceRoomId); } GeneratedNodeResult generatedResult = splitGeneratedNodeResult(node, value); @@ -1091,7 +1141,7 @@ private Map persistNativeNodeResult(String runId, String project AutomationConstants.NODE_OUTPUT_MAX_BYTES, "Automation node '" + nodeId + "' waiting output"); long duration = System.currentTimeMillis() - startedMs; String preview = AutomationRuntimeUtils.generatePreview(output); - AutomationDatabaseUtility.persistAgentWait(runId, projectId, nodeId, + AutomationRunStore.persistAgentWait(runId, projectId, nodeId, (String) node.get(AutomationConstants.NODE_FIELD_OUTPUT_VAR), output, preview, agentRunId, traceRoomId, resumeNodeId, currentUserId(), Instant.now().plus(DEFAULT_AGENT_APPROVAL_TIMEOUT_HOURS, ChronoUnit.HOURS), duration); @@ -1116,17 +1166,17 @@ traceRoomId, resumeNodeId, currentUserId(), String preview = AutomationRuntimeUtils.generatePreview(output); String modelMessageId = generatedAgentNode ? null : extractModelMessageId(node, traceMetadata, traceRoomId); if (agentFailure != null) { - AutomationDatabaseUtility.updateNodeFailedWithResult(runId, nodeId, started, duration, + AutomationRunStore.updateNodeFailedWithResult(runId, nodeId, started, duration, (String) node.get(AutomationConstants.NODE_FIELD_OUTPUT_VAR), output, preview, agentRunId, agentFailure); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_FAILED, duration, preview, agentFailure, traceForNode(node, traceRoomId, null, agentRunId)); return nodeResult(nodeId, AutomationConstants.NODE_STATUS_FAILED, persistedValue, agentFailure); } - AutomationDatabaseUtility.updateNodeSuccess(runId, nodeId, started, duration, + AutomationRunStore.updateNodeSuccess(runId, nodeId, started, duration, (String) node.get(AutomationConstants.NODE_FIELD_OUTPUT_VAR), output, preview, modelMessageId, agentRunId); - AutomationPythonRunRegistry.nodeCompleted(runId); + AutomationRunRegistry.nodeCompleted(runId); streamNodeProgress(runId, node, AutomationConstants.NODE_STATUS_SUCCESS, duration, preview, null, traceForNode(node, traceRoomId, modelMessageId, agentRunId)); return nodeResult(nodeId, AutomationConstants.NODE_STATUS_SUCCESS, persistedValue, null); @@ -1148,14 +1198,14 @@ traceRoomId, resumeNodeId, currentUserId(), * @return cancelled or current run detail */ Map cancelWaitingRun(String runId, String projectId) { - Map run = AutomationDatabaseUtility.getRunDetail(runId); + Map run = AutomationRunStore.getRunDetail(runId); if (run == null || !projectId.equals(run.get(AutomationConstants.PROJECT_ID))) { throw new IllegalArgumentException("Automation run not found: " + runId); } if (!AutomationConstants.STATUS_WAITING_FOR_INPUT.equals(run.get(AutomationConstants.STATUS))) { return buildCurrentRunResult(runId, projectId); } - Map wait = AutomationDatabaseUtility.claimWaitingRun(runId, projectId); + Map wait = AutomationRunStore.claimWaitingRun(runId, projectId); if (wait == null) { // A concurrent resume owns the run. It observes the durable cancellation // request before its next node and completes the run as CANCELLED. @@ -1163,11 +1213,11 @@ Map cancelWaitingRun(String runId, String projectId) { } String waitingNodeId = stringValue(wait.get(AutomationConstants.NODE_ID)); String message = "Run cancelled by user"; - AutomationDatabaseUtility.updateNodeFailed(runId, waitingNodeId, utcNow(), 0, message); - AutomationDatabaseUtility.resolveWait(runId, stringValue(wait.get(AutomationConstants.WAIT_ID)), + AutomationRunStore.updateNodeFailed(runId, waitingNodeId, utcNow(), 0, message); + AutomationRunStore.resolveWait(runId, stringValue(wait.get(AutomationConstants.WAIT_ID)), currentUserId()); - AutomationDatabaseUtility.skipPendingNodes(runId, message); - AutomationDatabaseUtility.completeRun(runId, projectId, AutomationConstants.STATUS_CANCELLED, waitingNodeId, + AutomationRunStore.skipPendingNodes(runId, message); + AutomationRunStore.completeRun(runId, projectId, AutomationConstants.STATUS_CANCELLED, waitingNodeId, message); return buildResult(runId, projectId, Map.of("error", message)); } @@ -1183,7 +1233,7 @@ Map cancelWaitingRun(String runId, String projectId) { * @return current or resumed run detail */ Map resumeWaitingRun(String runId, String projectId) { - Map run = AutomationDatabaseUtility.getRunDetail(runId); + Map run = AutomationRunStore.getRunDetail(runId); if (run == null || !projectId.equals(run.get(AutomationConstants.PROJECT_ID))) { throw new IllegalArgumentException("Automation run not found: " + runId); } @@ -1191,7 +1241,7 @@ Map resumeWaitingRun(String runId, String projectId) { return buildCurrentRunResult(runId, projectId); } - Map wait = AutomationDatabaseUtility.getActiveWait(runId); + Map wait = AutomationRunStore.getActiveWait(runId); if (wait == null) { throw new IllegalStateException("Automation run has no active input boundary: " + runId); } @@ -1207,7 +1257,7 @@ Map resumeWaitingRun(String runId, String projectId) { "Agent run '" + agentRunId + "' has not completed its input flow (" + agentStatus + ")."); } - wait = AutomationDatabaseUtility.claimWaitingRun(runId, projectId); + wait = AutomationRunStore.claimWaitingRun(runId, projectId); if (wait == null) { return buildCurrentRunResult(runId, projectId); } @@ -1233,10 +1283,10 @@ Map resumeWaitingRun(String runId, String projectId) { } String output = AutomationRuntimeUtils.toBoundedRuntimeJson(finalText, AutomationConstants.NODE_OUTPUT_MAX_BYTES, "Automation agent result"); - AutomationDatabaseUtility.updateNodeSuccess(runId, waitingNodeId, utcNow(), 0, + AutomationRunStore.updateNodeSuccess(runId, waitingNodeId, utcNow(), 0, (String) waitingNode.get(AutomationConstants.NODE_FIELD_OUTPUT_VAR), output, AutomationRuntimeUtils.generatePreview(output), null, agentRunId); - AutomationDatabaseUtility.resolveWait(runId, waitId, currentUserId()); + AutomationRunStore.resolveWait(runId, waitId, currentUserId()); executionInsightLease = openExecutionInsight(projectId, runId); executionInsight = executionInsightLease.insight(); @@ -1244,17 +1294,17 @@ Map resumeWaitingRun(String runId, String projectId) { if (translator == null) { throw new IllegalStateException("Python runtime is not available for this insight."); } - int completedNodes = (int) AutomationDatabaseUtility.getNodeOutputsForRun(runId).stream() + int completedNodes = (int) AutomationRunStore.getNodeOutputsForRun(runId).stream() .filter(row -> AutomationConstants.NODE_STATUS_SUCCESS.equals(row.get(AutomationConstants.STATUS)) || AutomationConstants.NODE_STATUS_SKIPPED.equals(row.get(AutomationConstants.STATUS))) .count(); - AutomationPythonRunRegistry.register(runId, translator, executionInsight, streamJobId, completedNodes); + AutomationRunRegistry.register(runId, translator, executionInsight, streamJobId, completedNodes); Map scope = reconstructScope(runId, executionInsight.getUser(), AutomationRuntime.startNodeId(definition)); String resumeNodeId = stringValue(wait.get(AutomationConstants.RESUME_NODE_ID)); if (resumeNodeId != null) { continuation = executeInControlOrder(executionInsight, projectId, runId, definition, runNodes, - AutomationDatabaseUtility.getRunNodeSources(runId), scope, traceRoomIds(runId), resumeNodeId); + AutomationRunStore.getRunNodeSources(runId), scope, traceRoomIds(runId), resumeNodeId); } else { continuation.put("scope", scope); } @@ -1266,7 +1316,7 @@ Map resumeWaitingRun(String runId, String projectId) { finishFailedRun(runId, projectId, e); continuation = Map.of("error", safeMessage(e)); } finally { - AutomationPythonRunRegistry.unregister(runId); + AutomationRunRegistry.unregister(runId); if (executionInsightLease == null || executionInsightLease.cleanupOnCompletion()) { cleanupExecutionInsight(executionInsight); } @@ -1287,11 +1337,11 @@ private Map finishTerminalAgentWait(String runId, String project } String output = AutomationRuntimeUtils.toBoundedRuntimeJson(agent.get("finalText"), AutomationConstants.NODE_OUTPUT_MAX_BYTES, "Automation agent terminal result"); - AutomationDatabaseUtility.updateNodeFailedWithResult(runId, waitingNodeId, utcNow(), 0, null, output, + AutomationRunStore.updateNodeFailedWithResult(runId, waitingNodeId, utcNow(), 0, null, output, AutomationRuntimeUtils.generatePreview(output), agentRunId, message); - AutomationDatabaseUtility.resolveWait(runId, waitId, currentUserId()); - AutomationDatabaseUtility.skipPendingNodes(runId, "Skipped because the agent run did not complete"); - AutomationDatabaseUtility.completeRun(runId, projectId, + AutomationRunStore.resolveWait(runId, waitId, currentUserId()); + AutomationRunStore.skipPendingNodes(runId, "Skipped because the agent run did not complete"); + AutomationRunStore.completeRun(runId, projectId, "CANCELLED".equalsIgnoreCase(agentStatus) ? AutomationConstants.STATUS_CANCELLED : AutomationConstants.STATUS_FAILED, waitingNodeId, message); @@ -1319,8 +1369,8 @@ private Map finishTerminalAgentWait(String runId, String project */ private static Map reconstructScope(String runId, prerna.auth.User user, String startNodeId) { Map scope = AutomationRuntimeUtils.buildInitialScope(runId, user); - scope.putAll(AutomationDatabaseUtility.getRunInputs(runId)); - for (Map row : AutomationDatabaseUtility.getNodeOutputsForRun(runId)) { + scope.putAll(AutomationRunStore.getRunInputs(runId)); + for (Map row : AutomationRunStore.getNodeOutputsForRun(runId)) { if (!AutomationConstants.NODE_STATUS_SUCCESS.equals(row.get(AutomationConstants.STATUS))) { continue; } @@ -1349,7 +1399,7 @@ private static Map reconstructScope(String runId, prerna.auth.Us */ private static Map traceRoomIds(String runId) { Map roomIds = new LinkedHashMap<>(); - for (Map row : AutomationDatabaseUtility.getNodeOutputsForRun(runId)) { + for (Map row : AutomationRunStore.getNodeOutputsForRun(runId)) { String roomId = stringValue(row.get(AutomationConstants.ROOM_ID)); if (roomId != null) { roomIds.put(String.valueOf(row.get(AutomationConstants.NODE_ID)), roomId); @@ -1682,10 +1732,10 @@ private static Map nodeResult(String nodeId, String status, Obje * settled is recorded as successful. */ private void finishRun(String runId, String projectId) { - List> outputs = AutomationDatabaseUtility.getNodeOutputsForRun(runId); - if (AutomationPythonRunRegistry.isCancellationRequested(runId)) { - AutomationDatabaseUtility.skipPendingNodes(runId, "Run cancelled by user"); - AutomationDatabaseUtility.completeRun(runId, projectId, AutomationConstants.STATUS_CANCELLED, null, + List> outputs = AutomationRunStore.getNodeOutputsForRun(runId); + if (AutomationRunRegistry.isCancellationRequested(runId)) { + AutomationRunStore.skipPendingNodes(runId, "Run cancelled by user"); + AutomationRunStore.completeRun(runId, projectId, AutomationConstants.STATUS_CANCELLED, null, "Run cancelled by user"); return; } @@ -1695,8 +1745,8 @@ private void finishRun(String runId, String projectId) { .findFirst().orElse(null); if (failed != null) { String nodeId = (String) failed.get(AutomationConstants.NODE_ID); - AutomationDatabaseUtility.skipPendingNodes(runId, "Skipped because an earlier node failed"); - AutomationDatabaseUtility.completeRun(runId, projectId, AutomationConstants.STATUS_FAILED, nodeId, + AutomationRunStore.skipPendingNodes(runId, "Skipped because an earlier node failed"); + AutomationRunStore.completeRun(runId, projectId, AutomationConstants.STATUS_FAILED, nodeId, (String) failed.get(AutomationConstants.ERROR_MESSAGE)); return; } @@ -1708,8 +1758,8 @@ private void finishRun(String runId, String projectId) { if (incomplete != null) { String nodeId = (String) incomplete.get(AutomationConstants.NODE_ID); String message = "Python source did not return a structured result for node " + nodeId + "."; - AutomationDatabaseUtility.skipPendingNodes(runId, message); - AutomationDatabaseUtility.completeRun(runId, projectId, AutomationConstants.STATUS_FAILED, nodeId, message); + AutomationRunStore.skipPendingNodes(runId, message); + AutomationRunStore.completeRun(runId, projectId, AutomationConstants.STATUS_FAILED, nodeId, message); return; } @@ -1717,8 +1767,8 @@ private void finishRun(String runId, String projectId) { .filter(output -> AutomationConstants.NODE_STATUS_SUCCESS.equals(output.get(AutomationConstants.STATUS)) || AutomationConstants.NODE_STATUS_SKIPPED.equals(output.get(AutomationConstants.STATUS))) .count(); - AutomationDatabaseUtility.updateHeartbeat(runId, completed); - AutomationDatabaseUtility.completeRun(runId, projectId, AutomationConstants.STATUS_SUCCESS, null, null); + AutomationRunStore.updateHeartbeat(runId, completed); + AutomationRunStore.completeRun(runId, projectId, AutomationConstants.STATUS_SUCCESS, null, null); } /** @@ -1727,17 +1777,17 @@ private void finishRun(String runId, String projectId) { * what interrupted the run. */ private void finishFailedRun(String runId, String projectId, Exception error) { - if (AutomationPythonRunRegistry.isCancellationRequested(runId)) { - AutomationDatabaseUtility.skipPendingNodes(runId, "Run cancelled by user"); - AutomationDatabaseUtility.completeRun(runId, projectId, AutomationConstants.STATUS_CANCELLED, null, + if (AutomationRunRegistry.isCancellationRequested(runId)) { + AutomationRunStore.skipPendingNodes(runId, "Run cancelled by user"); + AutomationRunStore.completeRun(runId, projectId, AutomationConstants.STATUS_CANCELLED, null, "Run cancelled by user"); return; } - String failedNodeId = AutomationDatabaseUtility.getNodeOutputsForRun(runId).stream() + String failedNodeId = AutomationRunStore.getNodeOutputsForRun(runId).stream() .filter(output -> AutomationConstants.NODE_STATUS_FAILED.equals(output.get(AutomationConstants.STATUS))) .map(output -> (String) output.get(AutomationConstants.NODE_ID)).findFirst().orElse(null); - AutomationDatabaseUtility.skipPendingNodes(runId, "Python runtime failed before this node executed"); - AutomationDatabaseUtility.completeRun(runId, projectId, AutomationConstants.STATUS_FAILED, failedNodeId, + AutomationRunStore.skipPendingNodes(runId, "Python runtime failed before this node executed"); + AutomationRunStore.completeRun(runId, projectId, AutomationConstants.STATUS_FAILED, failedNodeId, safeMessage(error)); } @@ -1746,14 +1796,14 @@ private void finishFailedRun(String runId, String projectId, Exception error) { * scope and globals, and persists the human-readable summary alongside it. */ private Map buildResult(String runId, String projectId, Map pythonResult) { - Map detail = AutomationDatabaseUtility.getRunDetail(runId); + Map detail = AutomationRunStore.getRunDetail(runId); if (detail == null) { detail = new LinkedHashMap<>(); detail.put(AutomationConstants.RUN_ID, runId); detail.put(AutomationConstants.PROJECT_ID, projectId); } - List> nodeResults = AutomationDatabaseUtility - .buildNodeResults(AutomationDatabaseUtility.getNodeOutputsForRun(runId)); + List> nodeResults = AutomationRunStore + .buildNodeResults(AutomationRunStore.getNodeOutputsForRun(runId)); detail.put(AutomationConstants.RESULT_NODE_RESULTS, nodeResults); detail.put("scope", normalizeScope(pythonResult.get("scope"))); detail.put(AutomationConstants.RESULT_GLOBALS, @@ -1771,7 +1821,7 @@ private Map buildResult(String runId, String projectId, Map RUNS = new ConcurrentHashMap<>(); + private static final Logger classLogger = LogManager.getLogger(AutomationRunRegistry.class); + private static final ConcurrentHashMap RUNS = new ConcurrentHashMap<>(); private static final ScheduledExecutorService HEARTBEAT_SCHEDULER = Executors .newSingleThreadScheduledExecutor(r -> { Thread thread = new Thread(r, "automation-python-heartbeat"); @@ -61,7 +62,7 @@ final class AutomationPythonRunRegistry { return thread; }); - private AutomationPythonRunRegistry() { + private AutomationRunRegistry() { } static void register(String runId, PyTranslator translator, Insight insight, String jobId) { @@ -69,13 +70,13 @@ static void register(String runId, PyTranslator translator, Insight insight, Str } static void register(String runId, PyTranslator translator, Insight insight, String jobId, int completedNodes) { - ActivePythonRun active = new ActivePythonRun(translator, insight.getInsightId(), jobId, completedNodes); + ActiveRun active = new ActiveRun(translator, insight.getInsightId(), jobId, completedNodes); if (RUNS.putIfAbsent(runId, active) != null) { throw new IllegalStateException("Python automation run is already registered: " + runId); } active.heartbeat = HEARTBEAT_SCHEDULER.scheduleAtFixedRate(() -> { try { - AutomationDatabaseUtility.touchHeartbeat(runId); + AutomationRunStore.touchHeartbeat(runId); } catch (Exception e) { classLogger.warn("Heartbeat update failed for Python automation run {}: {}", runId, e.getMessage()); } @@ -84,19 +85,19 @@ static void register(String runId, PyTranslator translator, Insight insight, Str } static void unregister(String runId) { - ActivePythonRun active = RUNS.remove(runId); + ActiveRun active = RUNS.remove(runId); if (active != null && active.heartbeat != null) { active.heartbeat.cancel(false); } } static boolean isCancellationRequested(String runId) { - ActivePythonRun active = RUNS.get(runId); - return (active != null && active.cancelled.get()) || AutomationDatabaseUtility.isCancelRequested(runId); + ActiveRun active = RUNS.get(runId); + return (active != null && active.cancelled.get()) || AutomationRunStore.isCancelRequested(runId); } static boolean requestCancellation(String runId) { - ActivePythonRun active = RUNS.get(runId); + ActiveRun active = RUNS.get(runId); if (active == null) { return false; } @@ -112,13 +113,13 @@ static boolean requestCancellation(String runId) { } static void nodeCompleted(String runId) { - ActivePythonRun active = RUNS.get(runId); + ActiveRun active = RUNS.get(runId); if (active != null) { - AutomationDatabaseUtility.updateHeartbeat(runId, active.completedNodes.incrementAndGet()); + AutomationRunStore.updateHeartbeat(runId, active.completedNodes.incrementAndGet()); } } - private static final class ActivePythonRun { + private static final class ActiveRun { private final PyTranslator translator; private final String insightId; private final String jobId; @@ -126,7 +127,7 @@ private static final class ActivePythonRun { private final AtomicInteger completedNodes; private ScheduledFuture heartbeat; - private ActivePythonRun(PyTranslator translator, String insightId, String jobId, int completedNodes) { + private ActiveRun(PyTranslator translator, String insightId, String jobId, int completedNodes) { this.translator = translator; this.insightId = insightId; this.jobId = jobId; diff --git a/src/prerna/reactor/automation/AutomationDatabaseUtility.java b/src/prerna/reactor/automation/run/AutomationRunStore.java similarity index 98% rename from src/prerna/reactor/automation/AutomationDatabaseUtility.java rename to src/prerna/reactor/automation/run/AutomationRunStore.java index 3348ff4a14..86856a96aa 100644 --- a/src/prerna/reactor/automation/AutomationDatabaseUtility.java +++ b/src/prerna/reactor/automation/run/AutomationRunStore.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import static prerna.reactor.automation.AutomationConstants.AGENT_RUN_ID; import static prerna.reactor.automation.AutomationConstants.AUTOMATION_ID; @@ -117,6 +117,9 @@ import prerna.query.querystruct.filters.SimpleQueryFilter; import prerna.query.querystruct.selectors.QueryColumnOrderBySelector; import prerna.query.querystruct.selectors.QueryColumnSelector; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; +import prerna.reactor.automation.definition.AutomationDefinitionService; import prerna.reactor.automation.utils.AutomationRuntimeUtils; import prerna.reactor.scheduler.SchedulerOwlCreator; import prerna.sablecc2.om.PixelDataType; @@ -136,9 +139,9 @@ * The scheduler OWL owns the logical schema; startup initialization creates the * corresponding physical tables in the scheduler database. */ -public final class AutomationDatabaseUtility { +public final class AutomationRunStore { - private static final Logger classLogger = LogManager.getLogger(AutomationDatabaseUtility.class); + private static final Logger classLogger = LogManager.getLogger(AutomationRunStore.class); /** Prefix identifying the automation tables inside the scheduler OWL schema. */ private static final String AUTOMATION_TABLE_PREFIX = "AUTOMATION_"; @@ -161,7 +164,7 @@ public final class AutomationDatabaseUtility { private static final String TABLE_NODE_OUTPUTS = TABLE_AUTOMATION_NODE_OUTPUTS; private static final String TABLE_RUN_WAITS = TABLE_AUTOMATION_RUN_WAITS; - private AutomationDatabaseUtility() { + private AutomationRunStore() { } // AUTOMATION_RUNS @@ -1274,9 +1277,10 @@ public static boolean hasAgentRunTrace(String projectId, String automationRunId, * to callers. * *

- * Each entry contains: nodeId, nodeLabel, status, durationMs, outputPreview - * (falls back from outputValue when blank), outputValue, errorMessage, and an - * optional trace map. + * Each entry contains: nodeId, nodeLabel, status, durationMs, a bounded display + * preview, an optional small outputValue, errorMessage, and an optional trace + * map. Large values remain server-owned and are not copied into run-detail + * responses. * * @param nodeOutputs ordered rows from {@link #getNodeOutputsForRun(String)} * @return mutable list of node result maps (empty when {@code nodeOutputs} is @@ -1293,12 +1297,16 @@ public static List> buildNodeResults(List trace = new LinkedHashMap<>(); putIfPresent(trace, AutomationConstants.TRACE_AUTOMATION_RUN_ID, output.get(AutomationConstants.RUN_ID)); diff --git a/src/prerna/reactor/automation/CancelAutomationRunReactor.java b/src/prerna/reactor/automation/run/CancelAutomationRunReactor.java similarity index 90% rename from src/prerna/reactor/automation/CancelAutomationRunReactor.java rename to src/prerna/reactor/automation/run/CancelAutomationRunReactor.java index fb0834a49e..80e23bf9e6 100644 --- a/src/prerna/reactor/automation/CancelAutomationRunReactor.java +++ b/src/prerna/reactor/automation/run/CancelAutomationRunReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import java.util.HashMap; import java.util.Map; @@ -35,6 +35,9 @@ import prerna.reactor.AbstractReactor; import prerna.reactor.agent.run.AgentRunService; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.agent.AutomationAgentRunAccess; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; import prerna.sablecc2.om.ReactorKeysEnum; @@ -53,7 +56,7 @@ *

* Sets a cluster-safe cancellation flag * ({@code AUTOMATION_RUNS.CANCEL_REQUESTED}, via - * {@link AutomationDatabaseUtility#setCancelRequested(String)}) that the + * {@link AutomationRunStore#setCancelRequested(String)}) that the * executing pod polls regardless of which pod owns the Python run. The same-pod * fast path also interrupts the matching Python socket job, allowing native * Python and blocking bridge calls to stop promptly. A running run's @@ -90,14 +93,14 @@ public NounMetadata execute() { throw new IllegalArgumentException("Must provide the run id to cancel"); } - projectId = AutomationProjectUtils.getEditableAutomationProject(this.insight.getUser(), projectId) + projectId = AutomationProjectService.getEditableAutomationProject(this.insight.getUser(), projectId) .getProjectId(); // Validate the run exists and belongs to this project. Scoping by // PROJECT_ID prevents a user with edit access to their own project from // cancelling // a run that belongs to a project they were never granted access to. - Map runDetail = AutomationDatabaseUtility.getRunDetail(runId); + Map runDetail = AutomationRunStore.getRunDetail(runId); if (runDetail == null || !projectId.equals(runDetail.get(AutomationConstants.PROJECT_ID))) { throw new IllegalArgumentException("Run not found: " + runId); } @@ -116,8 +119,8 @@ public NounMetadata execute() { // The executing pod owns the terminal status transition; stale recovery handles // a run // whose owner disappeared before observing the request. - AutomationDatabaseUtility.setCancelRequested(runId); - boolean signalledLocally = AutomationPythonRunRegistry.requestCancellation(runId); + AutomationRunStore.setCancelRequested(runId); + boolean signalledLocally = AutomationRunRegistry.requestCancellation(runId); classLogger.info("Cancel requested for automation run {}: signalledLocally={}", runId, signalledLocally); @@ -133,7 +136,7 @@ public NounMetadata execute() { * waiting run. */ private NounMetadata cancelWaitingAgentRun(String projectId, String runId) { - Map wait = AutomationDatabaseUtility.getActiveWait(runId); + Map wait = AutomationRunStore.getActiveWait(runId); if (wait == null) { throw new IllegalStateException("Waiting Automation run has no active input boundary: " + runId); } @@ -144,8 +147,8 @@ private NounMetadata cancelWaitingAgentRun(String projectId, String runId) { // Record the request before stopping the child agent. A resume that wins the // wait claim while the child settles observes this flag between nodes, so the // remaining graph stops either way. - AutomationDatabaseUtility.setCancelRequested(runId); - boolean signalledLocally = AutomationPythonRunRegistry.requestCancellation(runId); + AutomationRunStore.setCancelRequested(runId); + boolean signalledLocally = AutomationRunRegistry.requestCancellation(runId); AgentRunService.get().stopForAutomation(agentRunId, this.insight); Map run = new AutomationRunExecutionService(this.insight, null).cancelWaitingRun(runId, diff --git a/src/prerna/reactor/automation/GetAutomationRunReactor.java b/src/prerna/reactor/automation/run/GetAutomationRunReactor.java similarity index 75% rename from src/prerna/reactor/automation/GetAutomationRunReactor.java rename to src/prerna/reactor/automation/run/GetAutomationRunReactor.java index 5e854866c6..db40895bff 100644 --- a/src/prerna/reactor/automation/GetAutomationRunReactor.java +++ b/src/prerna/reactor/automation/run/GetAutomationRunReactor.java @@ -25,14 +25,17 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; +import prerna.om.Insight; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; import prerna.sablecc2.om.ReactorKeysEnum; @@ -49,6 +52,8 @@ */ public class GetAutomationRunReactor extends AbstractReactor { + private static final String OUTPUT_FRAME_KEY = "OUTPUT_FRAME"; + // Not standardized in ReactorKeysEnum — matches the local-key convention used // by prerna.reactor.agent (e.g. GetAgentRunReactor.RUN_ID_KEY). private static final String RUN_ID_KEY = "runId"; @@ -71,10 +76,10 @@ public NounMetadata execute() { throw new IllegalArgumentException("Must provide a run id"); } - projectId = AutomationProjectUtils.getViewableAutomationProject(this.insight.getUser(), projectId) + projectId = AutomationProjectService.getViewableAutomationProject(this.insight.getUser(), projectId) .getProjectId(); - Map runDetail = AutomationDatabaseUtility.getRunDetail(runId); + Map runDetail = AutomationRunStore.getRunDetail(runId); // Scope by PROJECT_ID so a user with view access to one project cannot read // another project's run detail/node outputs by guessing or reusing a runId. if (runDetail == null || !projectId.equals(runDetail.get(AutomationConstants.PROJECT_ID))) { @@ -84,15 +89,25 @@ public NounMetadata execute() { return new NounMetadata(notFound, PixelDataType.MAP, PixelOperationType.OPERATION); } - List> nodeOutputs = AutomationDatabaseUtility.getNodeOutputsForRun(runId); - List> nodeResults = AutomationDatabaseUtility.buildNodeResults(nodeOutputs); + List> nodeOutputs = AutomationRunStore.getNodeOutputsForRun(runId); + List> nodeResults = AutomationRunStore.buildNodeResults(nodeOutputs); runDetail.put(AutomationConstants.RESULT_NODE_RESULTS, nodeResults); - String executionInsightId = AutomationRunExecutionService.getAvailableExecutionInsightId(runId); - if (executionInsightId != null && !executionInsightId.isBlank()) { - runDetail.put(AutomationConstants.RESULT_EXECUTION_INSIGHT_ID, executionInsightId); + Insight executionInsight = AutomationRunExecutionService.getAvailableExecutionInsight(runId); + if (executionInsight != null && projectId.equals(executionInsight.getProjectId())) { + runDetail.put(AutomationConstants.RESULT_EXECUTION_INSIGHT_ID, executionInsight.getInsightId()); + for (int index = 0; index < nodeOutputs.size(); index++) { + Object outputVariable = nodeOutputs.get(index).get(AutomationConstants.OUTPUT_VAR_NAME); + if (!(outputVariable instanceof String name)) { + continue; + } + NounMetadata frame = executionInsight.getVarStore().get(name); + if (frame != null && frame.getNounType() == PixelDataType.FRAME) { + nodeResults.get(index).put(OUTPUT_FRAME_KEY, processNounMetadata(frame)); + } + } } - Map wait = AutomationDatabaseUtility.getActiveWait(runId); + Map wait = AutomationRunStore.getActiveWait(runId); if (wait != null) { runDetail.put("wait", wait); } diff --git a/src/prerna/reactor/automation/ListAutomationRunsReactor.java b/src/prerna/reactor/automation/run/ListAutomationRunsReactor.java similarity index 90% rename from src/prerna/reactor/automation/ListAutomationRunsReactor.java rename to src/prerna/reactor/automation/run/ListAutomationRunsReactor.java index 29a049fc64..53a67d90d2 100644 --- a/src/prerna/reactor/automation/ListAutomationRunsReactor.java +++ b/src/prerna/reactor/automation/run/ListAutomationRunsReactor.java @@ -25,13 +25,15 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import java.util.ArrayList; import java.util.List; import java.util.Map; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; import prerna.sablecc2.om.ReactorKeysEnum; @@ -69,10 +71,10 @@ public NounMetadata execute() { } limit = Math.min(limit, MAXIMUM_LIMIT); - projectId = AutomationProjectUtils.getViewableAutomationProject(this.insight.getUser(), projectId) + projectId = AutomationProjectService.getViewableAutomationProject(this.insight.getUser(), projectId) .getProjectId(); - List> runs = AutomationDatabaseUtility.getRunsForProject(projectId, limit); + List> runs = AutomationRunStore.getRunsForProject(projectId, limit); if (runs == null) { runs = new ArrayList<>(); } diff --git a/src/prerna/reactor/automation/ResumeAutomationRunReactor.java b/src/prerna/reactor/automation/run/ResumeAutomationRunReactor.java similarity index 88% rename from src/prerna/reactor/automation/ResumeAutomationRunReactor.java rename to src/prerna/reactor/automation/run/ResumeAutomationRunReactor.java index a10d6ed2a5..5daa51b996 100644 --- a/src/prerna/reactor/automation/ResumeAutomationRunReactor.java +++ b/src/prerna/reactor/automation/run/ResumeAutomationRunReactor.java @@ -25,12 +25,15 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import java.util.Map; import prerna.project.api.IProject; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.agent.AutomationAgentRunAccess; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; import prerna.sablecc2.om.ReactorKeysEnum; @@ -65,15 +68,15 @@ public NounMetadata execute() { throw new IllegalArgumentException("Must provide a run id"); } - IProject project = AutomationProjectUtils.getEditableAutomationProject(this.insight.getUser(), + IProject project = AutomationProjectService.getEditableAutomationProject(this.insight.getUser(), requestedProjectId); String projectId = project.getProjectId(); - Map run = AutomationDatabaseUtility.getRunDetail(runId); + Map run = AutomationRunStore.getRunDetail(runId); if (run == null || !projectId.equals(run.get(AutomationConstants.PROJECT_ID))) { throw new IllegalArgumentException("Automation run not found: " + runId); } - Map wait = AutomationDatabaseUtility.getActiveWait(runId); + Map wait = AutomationRunStore.getActiveWait(runId); if (wait != null) { AutomationAgentRunAccess.authorizeEdit(this.insight, projectId, runId, String.valueOf(wait.get(AutomationConstants.NODE_ID)), diff --git a/src/prerna/reactor/automation/TriggerAutomationReactor.java b/src/prerna/reactor/automation/run/TriggerAutomationReactor.java similarity index 91% rename from src/prerna/reactor/automation/TriggerAutomationReactor.java rename to src/prerna/reactor/automation/run/TriggerAutomationReactor.java index d9226a47ff..283638e9a6 100644 --- a/src/prerna/reactor/automation/TriggerAutomationReactor.java +++ b/src/prerna/reactor/automation/run/TriggerAutomationReactor.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import java.util.LinkedHashMap; import java.util.List; @@ -34,6 +34,11 @@ import prerna.om.ThreadStore; import prerna.reactor.AbstractReactor; +import prerna.reactor.automation.AutomationConstants; +import prerna.reactor.automation.AutomationRuntime; +import prerna.reactor.automation.definition.AutomationDefinitionService; +import prerna.reactor.automation.definition.AutomationDefinitionValidator; +import prerna.reactor.automation.project.AutomationProjectService; import prerna.sablecc2.om.PixelDataType; import prerna.sablecc2.om.PixelOperationType; import prerna.sablecc2.om.ReactorKeysEnum; @@ -63,7 +68,7 @@ public NounMetadata execute() { AutomationDefinitionService.DefinitionFiles files = AutomationDefinitionService.load(projectId); AutomationDefinitionValidator.ValidatedDefinition definition = AutomationDefinitionValidator.parseAndValidate(files.definition()); - AutomationProjectUtils.validateDefinitionReferences(definition, this.insight.getUser()); + AutomationProjectService.validateDefinitionReferences(definition, this.insight.getUser()); List> runNodes = AutomationRuntime.nodesForRun(definition); @SuppressWarnings("unchecked") Map inputs = this.getMap(AutomationConstants.AUTOMATION_INPUTS_KEY); @@ -88,14 +93,14 @@ private void initializeRun(String runId, String projectId, Map inputs, List> runNodes, Map traceRoomIds, Map nodeSources) { - AutomationDatabaseUtility.initializeRun(runId, projectId, + AutomationRunStore.initializeRun(runId, projectId, AutomationConstants.DEFAULT_AUTOMATION_ID, AutomationConstants.PYTHON_DOC_CURRENT_VERSION, definition.hash(), definition.snapshot(), inputs, getTriggerType(), getUserId(), runNodes, traceRoomIds, nodeSources); } private String getProjectId() { - return AutomationProjectUtils.getEditableAutomationProject(this.insight.getUser(), + return AutomationProjectService.getEditableAutomationProject(this.insight.getUser(), this.keyValue.get(ReactorKeysEnum.PROJECT.getKey())).getProjectId(); } diff --git a/src/prerna/reactor/automation/utils/AGENTS.md b/src/prerna/reactor/automation/utils/AGENTS.md new file mode 100644 index 0000000000..f8283433d6 --- /dev/null +++ b/src/prerna/reactor/automation/utils/AGENTS.md @@ -0,0 +1,20 @@ +# Automation Runtime Utilities Package + +This package is limited to shared runtime-boundary transformations. Follow the +parent `../AGENTS.md` for the complete Automation contract. + +## Invariants + +- Keep JSON serialization, bounded scope construction, and output-preview + transformations deterministic and side-effect free. +- Reject unsupported or oversized values with clear errors. +- Do not place graph ownership, persistence, permissions, execution lifecycle, + or project coordination here. +- Prefer a concrete owner in `definition`, `run`, `project`, or `agent` over a + new generic helper. + +## Verification + +Keep serialization and boundary tests under +`test/prerna/reactor/automation/utils`, including size, type, date/time, and +error-path coverage. diff --git a/src/prerna/util/SMSSWebWatcher.java b/src/prerna/util/SMSSWebWatcher.java index 21f98ef50c..576c0dd7a8 100644 --- a/src/prerna/util/SMSSWebWatcher.java +++ b/src/prerna/util/SMSSWebWatcher.java @@ -48,7 +48,7 @@ import prerna.collaboration.CollaborationDbUtils; import prerna.notifications.NotificationDbUtils; import prerna.prompt.PromptUtils; -import prerna.reactor.automation.AutomationDatabaseUtility; +import prerna.reactor.automation.run.AutomationRunStore; import prerna.reactor.scheduler.SchedulerDatabaseUtility; import prerna.theme.AbstractThemeUtils; import prerna.usertracking.UserTrackingUtils; @@ -216,8 +216,8 @@ public void init() { SchedulerDatabaseUtility.startServer(); // Automation tables live in the scheduler DB, so only initialize them // after the scheduler DB has started successfully. - AutomationDatabaseUtility.initialize(); - AutomationDatabaseUtility.markStaleRunsInterrupted(); + AutomationRunStore.initialize(); + AutomationRunStore.markStaleRunsInterrupted(); } catch (Exception e) { classLogger.error("Failed to load and start the scheduler database", e); } diff --git a/test/prerna/reactor/automation/AutomationRuntimeUnitTests.java b/test/prerna/reactor/automation/AutomationRuntimeUnitTests.java index 478064e49f..d278b9b9f4 100644 --- a/test/prerna/reactor/automation/AutomationRuntimeUnitTests.java +++ b/test/prerna/reactor/automation/AutomationRuntimeUnitTests.java @@ -38,6 +38,8 @@ import org.junit.jupiter.api.Test; +import prerna.reactor.automation.definition.AutomationDefinitionService; + /** * Covers how the trigger node contributes to a run. The trigger is the only node * whose Python lives inside the definition rather than in its own file, so the diff --git a/test/prerna/reactor/automation/AutomationDefinitionValidatorUnitTests.java b/test/prerna/reactor/automation/definition/AutomationDefinitionValidatorUnitTests.java similarity index 90% rename from test/prerna/reactor/automation/AutomationDefinitionValidatorUnitTests.java rename to test/prerna/reactor/automation/definition/AutomationDefinitionValidatorUnitTests.java index 3895000a14..e26e876605 100644 --- a/test/prerna/reactor/automation/AutomationDefinitionValidatorUnitTests.java +++ b/test/prerna/reactor/automation/definition/AutomationDefinitionValidatorUnitTests.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -37,6 +37,7 @@ import org.junit.jupiter.api.Test; +import prerna.reactor.automation.AutomationConstants; import prerna.reactor.automation.utils.AutomationRuntimeUtils; /** @@ -125,6 +126,12 @@ private static Map databaseConfig(String query) { return config; } + private static Map databaseQueryConfig() { + Map config = databaseConfig("SELECT * FROM CLAIMS ORDER BY ID"); + config.put(AutomationConstants.CONFIG_LIMIT, 100); + return config; + } + private static Map jevConfig() { Map config = new LinkedHashMap<>(); config.put(AutomationConstants.CONFIG_ENGINE_ID, "typesafe-engine"); @@ -181,6 +188,23 @@ void acceptsEveryDatabaseWriteNode() { } } + @Test + void acceptsDatabaseQueryPagingOptions() { + AutomationDefinitionValidator.ValidatedDefinition validated = AutomationDefinitionValidator + .parseAndValidateForAuthoring(definition(Map.of(), + workNode(AutomationConstants.NODE_DATABASE_QUERY, databaseQueryConfig()))); + assertEquals(2, validated.nodes().size()); + } + + @Test + void rejectsInvalidDatabaseQueryLimit() { + Map invalidLimit = databaseQueryConfig(); + invalidLimit.put(AutomationConstants.CONFIG_LIMIT, AutomationConstants.DB_QUERY_MAX_LIMIT + 1); + assertThrows(IllegalArgumentException.class, + () -> AutomationDefinitionValidator.parseAndValidateForAuthoring(definition(Map.of(), + workNode(AutomationConstants.NODE_DATABASE_QUERY, invalidLimit)))); + } + @Test void acceptsAJevDecisionDraft() { Map node = workNode(AutomationConstants.NODE_CONTROL_JEV, jevConfig()); @@ -322,6 +346,18 @@ void detachedNodesAreDraftOnly() { assertThrows(IllegalArgumentException.class, () -> AutomationDefinitionValidator.parseAndValidate(json)); } + @Test + void visionAcceptsNestedScopeReferencesForMediaUrls() { + Map config = new LinkedHashMap<>(); + config.put(AutomationConstants.CONFIG_ENGINE_ID, "model-1"); + config.put("prompt", "Describe every downloaded image."); + config.put("image", ""); + config.put("urls", "${download.files}"); + + AutomationDefinitionValidator.parseAndValidateForAuthoring( + definition(Map.of(), workNode(AutomationConstants.NODE_MODEL_VISION, config))); + } + @Test void rejectsAnUnknownFormatVersion() { String json = document(Map.of(), null, AutomationConstants.PYTHON_DOC_CURRENT_VERSION + 1, true); diff --git a/test/prerna/reactor/automation/AutomationNodeCatalogUnitTests.java b/test/prerna/reactor/automation/definition/AutomationNodeCatalogUnitTests.java similarity index 74% rename from test/prerna/reactor/automation/AutomationNodeCatalogUnitTests.java rename to test/prerna/reactor/automation/definition/AutomationNodeCatalogUnitTests.java index 7e13c9e8bb..9829bd573a 100644 --- a/test/prerna/reactor/automation/AutomationNodeCatalogUnitTests.java +++ b/test/prerna/reactor/automation/definition/AutomationNodeCatalogUnitTests.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; @@ -38,6 +38,8 @@ import org.junit.jupiter.api.Test; +import prerna.reactor.automation.AutomationConstants; + /** * Covers the node catalog the canvas builds its palette from. The catalog is the * only description of a node the client ever sees, so a node missing from it is @@ -78,6 +80,19 @@ void everyDatabaseOperationIsOffered() { } } + @Test + void databaseQueryAdvertisesEditableLimit() { + AutomationNodeDefinition definition = find(AutomationConstants.NODE_DATABASE_QUERY); + assertNotNull(definition); + Map> fieldsByKey = definition.configFields().stream() + .collect(java.util.stream.Collectors.toMap(AutomationNodeDefinition.ConfigField::key, + AutomationNodeDefinition.ConfigField::toMap)); + assertEquals(AutomationConstants.DEFAULT_DB_QUERY_LIMIT, + definition.defaultConfig().get(AutomationConstants.CONFIG_LIMIT)); + assertEquals(AutomationConstants.DB_QUERY_MAX_LIMIT, + fieldsByKey.get(AutomationConstants.CONFIG_LIMIT).get("maximum")); + } + @Test void storageDownloadAdvertisesItsStructuredResult() { AutomationNodeDefinition definition = find(AutomationConstants.NODE_STORAGE_DOWNLOAD); @@ -91,6 +106,25 @@ void storageDownloadAdvertisesItsStructuredResult() { assertEquals(false, fieldsByKey.get("filePath").get("required")); } + @Test + void optionalEngineParametersAreBackendOwned() { + Map> uploadFields = find(AutomationConstants.NODE_STORAGE_UPLOAD).configFields() + .stream().collect(java.util.stream.Collectors.toMap(AutomationNodeDefinition.ConfigField::key, + AutomationNodeDefinition.ConfigField::toMap)); + Map> searchFields = find(AutomationConstants.NODE_VECTOR_SEARCH).configFields() + .stream().collect(java.util.stream.Collectors.toMap(AutomationNodeDefinition.ConfigField::key, + AutomationNodeDefinition.ConfigField::toMap)); + Map> visionFields = find(AutomationConstants.NODE_MODEL_VISION).configFields() + .stream().collect(java.util.stream.Collectors.toMap(AutomationNodeDefinition.ConfigField::key, + AutomationNodeDefinition.ConfigField::toMap)); + + assertEquals("json", uploadFields.get("metadata").get("type")); + assertEquals("json", searchFields.get("filters").get("type")); + assertEquals("json", searchFields.get("paramValues").get("type")); + assertEquals("string[]", visionFields.get("urls").get("type")); + assertEquals("json", visionFields.get("paramValues").get("type")); + } + @Test void jevDecisionAdvertisesTypeSafeRoutingConfiguration() { AutomationNodeDefinition definition = find(AutomationConstants.NODE_CONTROL_JEV); diff --git a/test/prerna/reactor/automation/AutomationScopeCatalogUnitTests.java b/test/prerna/reactor/automation/definition/AutomationScopeCatalogUnitTests.java similarity index 98% rename from test/prerna/reactor/automation/AutomationScopeCatalogUnitTests.java rename to test/prerna/reactor/automation/definition/AutomationScopeCatalogUnitTests.java index 9e982bda59..b58c3ec1a2 100644 --- a/test/prerna/reactor/automation/AutomationScopeCatalogUnitTests.java +++ b/test/prerna/reactor/automation/definition/AutomationScopeCatalogUnitTests.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -35,6 +35,7 @@ import org.junit.jupiter.api.Test; +import prerna.reactor.automation.AutomationConstants; import prerna.reactor.automation.utils.AutomationRuntimeUtils; /** Covers the server-owned scope metadata consumed by editor variable helpers. */ diff --git a/test/prerna/reactor/automation/AutomationSourceRendererUnitTests.java b/test/prerna/reactor/automation/definition/AutomationSourceRendererUnitTests.java similarity index 86% rename from test/prerna/reactor/automation/AutomationSourceRendererUnitTests.java rename to test/prerna/reactor/automation/definition/AutomationSourceRendererUnitTests.java index 44b259b4c7..58875b9d65 100644 --- a/test/prerna/reactor/automation/AutomationSourceRendererUnitTests.java +++ b/test/prerna/reactor/automation/definition/AutomationSourceRendererUnitTests.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; @@ -37,6 +37,8 @@ import org.junit.jupiter.api.Test; +import prerna.reactor.automation.AutomationConstants; + /** * Covers the Python each generated node renders. The renderer is the only place * a node's behavior is defined, so a wrong SDK call here is a silent data bug at @@ -65,6 +67,7 @@ void readsGoThroughExecQuery() { String source = AutomationSourceRenderer.renderNode( node(AutomationConstants.NODE_DATABASE_QUERY, databaseConfig())); assertTrue(source.contains("def run(scope):"), "every node source defines the run entry point"); + assertTrue(source.contains("_pixel_value(\"limit\", int(scope.resolve(LIMIT)))")); assertFalse(source.contains("insertData")); assertFalse(source.contains("removeData")); } @@ -131,10 +134,32 @@ void visionAcceptsOneOrManyMediaPathsWithoutNestingThem() { String source = AutomationSourceRenderer.renderNode(node(AutomationConstants.NODE_MODEL_VISION, config)); assertTrue(source.contains("value if isinstance(value, list) else [value]")); - assertTrue(source.contains("media=_automation_media(scope.resolve(MEDIA))")); + assertTrue(source.contains("media=media or None")); assertFalse(source.contains("image=[scope.resolve(IMAGE)]")); } + @Test + void optionalEngineParametersReachTheirSdkCalls() { + Map vectorConfig = new LinkedHashMap<>(); + vectorConfig.put("engineId", "vector-1"); + vectorConfig.put("value", "claims"); + vectorConfig.put("filters", Map.of("category", "report")); + vectorConfig.put("paramValues", Map.of("threshold", 0.7)); + String vectorSource = AutomationSourceRenderer + .renderNode(node(AutomationConstants.NODE_VECTOR_SEARCH, vectorConfig)); + assertTrue(vectorSource.contains("filters=scope.resolve_config(FILTERS)")); + assertTrue(vectorSource.contains("param_dict=scope.resolve_config(PARAMETERS_JSON)")); + + Map uploadConfig = new LinkedHashMap<>(); + uploadConfig.put("engineId", "storage-1"); + uploadConfig.put("path", "archive"); + uploadConfig.put("destination", "report.pdf"); + uploadConfig.put("metadata", Map.of("case", "123")); + String uploadSource = AutomationSourceRenderer + .renderNode(node(AutomationConstants.NODE_STORAGE_UPLOAD, uploadConfig)); + assertTrue(uploadSource.contains("metadata=scope.resolve_config(METADATA)")); + } + /** * The canvas seeds its trigger editor with an identical template, and the save * path compares the persisted source against this string to decide whether the diff --git a/test/prerna/reactor/automation/RemoveAutomationStepUnitTests.java b/test/prerna/reactor/automation/definition/RemoveAutomationStepUnitTests.java similarity index 96% rename from test/prerna/reactor/automation/RemoveAutomationStepUnitTests.java rename to test/prerna/reactor/automation/definition/RemoveAutomationStepUnitTests.java index 7843f11a83..fe50317cca 100644 --- a/test/prerna/reactor/automation/RemoveAutomationStepUnitTests.java +++ b/test/prerna/reactor/automation/definition/RemoveAutomationStepUnitTests.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.definition; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -35,6 +35,8 @@ import org.junit.jupiter.api.Test; +import prerna.reactor.automation.AutomationConstants; + /** Covers downstream selection for the destructive branch-removal option. */ public class RemoveAutomationStepUnitTests { diff --git a/test/prerna/reactor/automation/AutomationRunExecutionServiceUnitTests.java b/test/prerna/reactor/automation/run/AutomationRunExecutionServiceUnitTests.java similarity index 98% rename from test/prerna/reactor/automation/AutomationRunExecutionServiceUnitTests.java rename to test/prerna/reactor/automation/run/AutomationRunExecutionServiceUnitTests.java index 3606e3f64e..6f27f4d5c4 100644 --- a/test/prerna/reactor/automation/AutomationRunExecutionServiceUnitTests.java +++ b/test/prerna/reactor/automation/run/AutomationRunExecutionServiceUnitTests.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNull; @@ -39,6 +39,7 @@ import prerna.engine.impl.model.responses.TypeSafeModelEngineResponse; import prerna.om.Insight; import prerna.om.InsightStore; +import prerna.reactor.automation.AutomationConstants; /** Covers deterministic execution-service mappings and run Insight lookup. */ public class AutomationRunExecutionServiceUnitTests { diff --git a/test/prerna/reactor/automation/AutomationDatabaseUtilityUnitTests.java b/test/prerna/reactor/automation/run/AutomationRunStoreUnitTests.java similarity index 81% rename from test/prerna/reactor/automation/AutomationDatabaseUtilityUnitTests.java rename to test/prerna/reactor/automation/run/AutomationRunStoreUnitTests.java index 1c74981e42..daeda99f3e 100644 --- a/test/prerna/reactor/automation/AutomationDatabaseUtilityUnitTests.java +++ b/test/prerna/reactor/automation/run/AutomationRunStoreUnitTests.java @@ -25,7 +25,7 @@ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. *******************************************************************************/ -package prerna.reactor.automation; +package prerna.reactor.automation.run; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; @@ -48,14 +48,14 @@ import prerna.util.JdbcTestDatabase; import prerna.util.SystemEngineRegistry; -class AutomationDatabaseUtilityUnitTests { +class AutomationRunStoreUnitTests { @Test void missingSchedulerDatabaseKeepsInitializationStateError() { try (var registry = mockStatic(SystemEngineRegistry.class, CALLS_REAL_METHODS)) { registry.when(SystemEngineRegistry::isSchedulerDbLoaded).thenReturn(false); IllegalStateException error = assertThrows(IllegalStateException.class, - () -> AutomationDatabaseUtility.initializeRun("run", "project", "automation", 1, "hash", "{}", + () -> AutomationRunStore.initializeRun("run", "project", "automation", 1, "hash", "{}", Map.of(), "manual", "user", List.of(), Map.of(), Map.of())); assertEquals("System database 'scheduler' is required for initializing automation run history", error.getMessage()); @@ -71,7 +71,7 @@ private void schema(JdbcTestDatabase db) throws Exception { } private void initialize(String runId, Map sources) { - AutomationDatabaseUtility.initializeRun(runId, "project", "automation", 1, "hash", "{}", + AutomationRunStore.initializeRun(runId, "project", "automation", 1, "hash", "{}", Map.of("value", ""), "manual", "user", List.of(), Map.of(), sources); } @@ -84,15 +84,15 @@ void initializationAndClaimPreserveSnapshotAndSingleClaim() throws Exception { db.connection.rollback(); assertEquals(1, db.count("AUTOMATION_RUNS")); assertEquals(" return 1\n", db.value("SELECT SOURCE_CODE FROM AUTOMATION_RUN_NODE_SOURCES")); - assertTrue(AutomationDatabaseUtility.claimRun("run")); - assertFalse(AutomationDatabaseUtility.claimRun("run")); + assertTrue(AutomationRunStore.claimRun("run")); + assertFalse(AutomationRunStore.claimRun("run")); assertEquals("RUNNING", db.value("SELECT STATUS FROM AUTOMATION_RUNS")); - AutomationDatabaseUtility.setCancelRequested("run"); + AutomationRunStore.setCancelRequested("run"); assertEquals(true, db.value("SELECT CANCEL_REQUESTED FROM AUTOMATION_RUNS")); - assertTrue(AutomationDatabaseUtility.updateHeartbeat("run", 2)); - assertTrue(AutomationDatabaseUtility.touchHeartbeat("run")); - assertTrue(AutomationDatabaseUtility.updateRunSummary("run", " result ")); - AutomationDatabaseUtility.completeRun("run", "project", "SUCCESS", null, null); + assertTrue(AutomationRunStore.updateHeartbeat("run", 2)); + assertTrue(AutomationRunStore.touchHeartbeat("run")); + assertTrue(AutomationRunStore.updateRunSummary("run", " result ")); + AutomationRunStore.completeRun("run", "project", "SUCCESS", null, null); assertEquals("SUCCESS", db.value("SELECT STATUS FROM AUTOMATION_RUNS")); assertEquals(" result ", db.value("SELECT RESULT_SUMMARY FROM AUTOMATION_RUNS")); assertFalse(db.connection.isClosed()); @@ -119,7 +119,7 @@ void zeroRowCancellationStillFailsAndRollsBack() throws Exception { schema(db); db.manual(); var error = assertThrows(IllegalStateException.class, - () -> AutomationDatabaseUtility.setCancelRequested("absent")); + () -> AutomationRunStore.setCancelRequested("absent")); assertEquals("Unable to persist the automation cancellation request.", error.getMessage()); verify(db.connection).rollback(); verify(db.connection, never()).commit(); @@ -133,36 +133,36 @@ void nodeResultsAndWaitContinuationRetainStatePredicatesAndTrace() throws Except schema(db); db.execute( "CREATE TABLE AUTOMATION_RUN_WAITS (WAIT_ID VARCHAR, RUN_ID VARCHAR, NODE_ID VARCHAR, WAIT_TYPE VARCHAR, AGENT_RUN_ID VARCHAR, ROOM_ID VARCHAR, RESUME_NODE_ID VARCHAR, STATUS VARCHAR, CREATED_BY VARCHAR, STARTED_AT TIMESTAMP, EXPIRES_AT TIMESTAMP, RESOLVED_AT TIMESTAMP, RESOLVED_BY VARCHAR)"); - AutomationDatabaseUtility.initializeRun("run", "project", "automation", 1, "hash", "{}", Map.of(), "manual", + AutomationRunStore.initializeRun("run", "project", "automation", 1, "hash", "{}", Map.of(), "manual", "user", List.of(Map.of("id", "node", "label", "Node"), Map.of("id", "skipped", "label", "Skipped")), Map.of(), Map.of()); - assertTrue(AutomationDatabaseUtility.claimRun("run")); - AutomationDatabaseUtility.markNodeRunning("run", "node"); - AutomationDatabaseUtility.updateNodeAgentRunTrace("run", "node", "agent"); - String waitId = AutomationDatabaseUtility.persistAgentWait("run", "project", "node", "output", " {} ", + assertTrue(AutomationRunStore.claimRun("run")); + AutomationRunStore.markNodeRunning("run", "node"); + AutomationRunStore.updateNodeAgentRunTrace("run", "node", "agent"); + String waitId = AutomationRunStore.persistAgentWait("run", "project", "node", "output", " {} ", "preview", "agent", "room", null, "user", java.time.Instant.now().plusSeconds(600), 4L); assertEquals("WAITING_FOR_INPUT", db.value("SELECT STATUS FROM AUTOMATION_RUNS")); var waiting = new java.util.HashMap( Map.of("WAIT_ID", waitId, "NODE_ID", "node", "STATUS", "PENDING")); queries.when(() -> prerna.util.QueryExecutionUtility.flushRsToMap(eq(db.engine), any(prerna.query.querystruct.SelectQueryStruct.class))).thenReturn(List.of(waiting)); - assertNotNull(AutomationDatabaseUtility.claimWaitingRun("run", "project")); - assertNull(AutomationDatabaseUtility.claimWaitingRun("run", "project")); - AutomationDatabaseUtility.resolveWait("run", waitId, "resolver"); + assertNotNull(AutomationRunStore.claimWaitingRun("run", "project")); + assertNull(AutomationRunStore.claimWaitingRun("run", "project")); + AutomationRunStore.resolveWait("run", waitId, "resolver"); assertEquals("RESOLVED", db.value("SELECT STATUS FROM AUTOMATION_RUN_WAITS")); var started = java.sql.Timestamp.from(java.time.Instant.now()); - AutomationDatabaseUtility.updateNodeSuccess("run", "node", started, 5L, "output", " {} ", "preview", null, + AutomationRunStore.updateNodeSuccess("run", "node", started, 5L, "output", " {} ", "preview", null, "agent"); assertEquals(" {} ", db.value("SELECT OUTPUT_VALUE FROM AUTOMATION_NODE_OUTPUTS WHERE NODE_ID='node'")); - AutomationDatabaseUtility.updateNodeFailed("run", "node", started, 6L, "failure"); - AutomationDatabaseUtility.updateNodeFailedWithResult("run", "node", started, 7L, "output", null, null, null, + AutomationRunStore.updateNodeFailed("run", "node", started, 6L, "failure"); + AutomationRunStore.updateNodeFailedWithResult("run", "node", started, 7L, "output", null, null, null, "failed result"); assertEquals("agent", db.value("SELECT AGENT_RUN_ID FROM AUTOMATION_NODE_OUTPUTS WHERE NODE_ID='node'")); assertNull(db.value("SELECT OUTPUT_VALUE FROM AUTOMATION_NODE_OUTPUTS WHERE NODE_ID='node'")); - AutomationDatabaseUtility.skipPendingNodes("run", "not selected"); + AutomationRunStore.skipPendingNodes("run", "not selected"); assertEquals("SKIPPED", db.value("SELECT STATUS FROM AUTOMATION_NODE_OUTPUTS WHERE NODE_ID='skipped'")); assertThrows(IllegalStateException.class, - () -> AutomationDatabaseUtility.updateNodeAgentRunTrace("run", "missing", "agent")); + () -> AutomationRunStore.updateNodeAgentRunTrace("run", "missing", "agent")); } } @@ -184,7 +184,7 @@ void staleRecoveryLeavesFreshHeartbeatsRunning() throws Exception { any(prerna.query.querystruct.SelectQueryStruct.class))) .thenReturn(List.of(Map.of("RUN_ID", "stale", "STATUS", "SUBMITTED", "LAST_HEARTBEAT", stale), Map.of("RUN_ID", "fresh", "STATUS", "SUBMITTED", "LAST_HEARTBEAT", fresh))); - AutomationDatabaseUtility.markStaleRunsInterrupted(); + AutomationRunStore.markStaleRunsInterrupted(); assertEquals("INTERRUPTED", db.value("SELECT STATUS FROM AUTOMATION_RUNS WHERE RUN_ID='stale'")); assertEquals("SUBMITTED", db.value("SELECT STATUS FROM AUTOMATION_RUNS WHERE RUN_ID='fresh'")); }