Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 29 additions & 2 deletions py/semoss_automation_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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(
Expand Down Expand Up @@ -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)
Expand Down
31 changes: 28 additions & 3 deletions src/prerna/reactor/automation/AGENTS.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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.
41 changes: 23 additions & 18 deletions src/prerna/reactor/automation/AutomationRuntime.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<Map<String, Object>> nodesForRun(AutomationDefinitionValidator.ValidatedDefinition definition) {
public static List<Map<String, Object>> nodesForRun(AutomationDefinitionValidator.ValidatedDefinition definition) {
List<Map<String, Object>> nodes = new ArrayList<>();
for (Map<String, Object> original : controlOrderedNodes(definition)) {
Map<String, Object> node = new LinkedHashMap<>(original);
Expand Down Expand Up @@ -86,7 +87,8 @@ private static String defaultOutputVariable(Map<String, Object> node) {
* initialization. Runtime traversal still selects only one condition path and
* remains Java-owned.
*/
static List<Map<String, Object>> controlOrderedNodes(AutomationDefinitionValidator.ValidatedDefinition definition) {
static List<Map<String, Object>> controlOrderedNodes(
AutomationDefinitionValidator.ValidatedDefinition definition) {
Map<String, Map<String, Object>> nodes = new LinkedHashMap<>();
Map<String, Integer> incoming = new LinkedHashMap<>();
for (Map<String, Object> node : definition.nodes()) {
Expand Down Expand Up @@ -128,7 +130,7 @@ static List<Map<String, Object>> controlOrderedNodes(AutomationDefinitionValidat
return ordered;
}

static String startNodeId(AutomationDefinitionValidator.ValidatedDefinition definition) {
public static String startNodeId(AutomationDefinitionValidator.ValidatedDefinition definition) {
for (Map<String, Object> node : definition.nodes()) {
if (AutomationConstants.NODE_START.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) {
return (String) node.get(AutomationConstants.NODE_FIELD_ID);
Expand All @@ -137,7 +139,7 @@ static String startNodeId(AutomationDefinitionValidator.ValidatedDefinition defi
throw new IllegalArgumentException("Automation definition has no trigger.start node.");
}

static Map<String, Map<String, String>> controlTargets(
public static Map<String, Map<String, String>> controlTargets(
AutomationDefinitionValidator.ValidatedDefinition definition) {
Map<String, Map<String, String>> targets = new LinkedHashMap<>();
for (Map<String, Object> edge : definition.edges()) {
Expand All @@ -155,43 +157,46 @@ static Map<String, Map<String, String>> controlTargets(
/**
* Runs one node module with the workflow scope supplied by the Java scheduler.
*/
static String buildNodeInvocationScript(String source, Map<String, Object> scope) {
return buildPythonInvocation("execute_node", source, scope);
public static String buildNodeInvocationScript(String source, Map<String, Object> scope, String outputVariable) {
return buildPythonInvocation("execute_node", source, scope, outputVariable);
}

/**
* Executes trigger Python in an isolated module and returns its non-private,
* JSON-compatible globals. A trigger may also return a map from
* {@code run(scope)} to define computed globals.
*/
static String buildTriggerInvocationScript(String source, Map<String, Object> scope) {
return buildPythonInvocation("execute_trigger", source, scope);
public static String buildTriggerInvocationScript(String source, Map<String, Object> scope) {
return buildPythonInvocation("execute_trigger", source, scope, null);
}

private static String buildPythonInvocation(String function, String source, Map<String, Object> scope) {
private static String buildPythonInvocation(String function, String source, Map<String, Object> 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<String, Object> declaredGlobals(AutomationDefinitionValidator.ValidatedDefinition definition) {
public static Map<String, Object> declaredGlobals(AutomationDefinitionValidator.ValidatedDefinition definition) {
for (Map<String, Object> node : definition.nodes()) {
if (AutomationConstants.NODE_START.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) {
return triggerGlobalDefaults(node);
Expand All @@ -204,7 +209,7 @@ static Map<String, Object> declaredGlobals(AutomationDefinitionValidator.Validat
* Returns the trigger declarations held in
* {@code trigger.start.config.globals}.
*/
static List<Map<String, Object>> triggerGlobalDefinitions(
public static List<Map<String, Object>> triggerGlobalDefinitions(
AutomationDefinitionValidator.ValidatedDefinition definition) {
for (Map<String, Object> node : definition.nodes()) {
if (AutomationConstants.NODE_START.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) {
Expand Down Expand Up @@ -235,7 +240,7 @@ private static List<Map<String, Object>> triggerGlobalDefinitions(Map<String, Ob
* {@code trigger.start.config.pythonSource}.
*/
@SuppressWarnings("unchecked")
static String triggerSource(Map<String, Object> node) {
public static String triggerSource(Map<String, Object> node) {
Object rawConfig = node.get(AutomationConstants.NODE_FIELD_CONFIG);
if (rawConfig instanceof Map<?, ?> raw) {
return sourceValue(((Map<String, Object>) raw).get(AutomationConstants.CONFIG_PYTHON_SOURCE));
Expand All @@ -247,7 +252,7 @@ static String triggerSource(Map<String, Object> node) {
* Returns the declared default for each trigger global the run should seed into
* scope.
*/
static Map<String, Object> triggerGlobalDefaults(Map<String, Object> node) {
public static Map<String, Object> triggerGlobalDefaults(Map<String, Object> node) {
Map<String, Object> globals = new LinkedHashMap<>();
for (Map<String, Object> global : triggerGlobalDefinitions(node)) {
if (global.get("name") instanceof String name
Expand All @@ -264,7 +269,7 @@ static Map<String, Object> triggerGlobalDefaults(Map<String, Object> 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<String, Object> node) {
public static String configuredAgentWorkspaceId(Map<String, Object> node) {
if (!AutomationConstants.NODE_AGENT_RUN.equals(node.get(AutomationConstants.NODE_FIELD_TYPE))) {
return null;
}
Expand All @@ -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 {
Expand Down
19 changes: 19 additions & 0 deletions src/prerna/reactor/automation/agent/AGENTS.md
Original file line number Diff line number Diff line change
@@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -39,38 +41,37 @@
* 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() {
}

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;
}
}

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.");
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
30 changes: 30 additions & 0 deletions src/prerna/reactor/automation/definition/AGENTS.md
Original file line number Diff line number Diff line change
@@ -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.
Loading
Loading