mirror of
https://github.com/Skyvern-AI/skyvern.git
synced 2026-10-02 19:57:59 +00:00
Reuse upstream block outputs from the origin run in Copilot repair tests (SKY-16495) (#8737)
This commit is contained in:
parent
2d090daf1b
commit
577ece32b6
14 changed files with 1812 additions and 69 deletions
|
|
@ -1302,6 +1302,8 @@ class CopilotContext(AgentContext):
|
|||
dispatched_run_ids_this_turn: set[str] = field(default_factory=set)
|
||||
# The files the last recorded worker run published, keyed by that run's id.
|
||||
delivered_output_files: tuple[str, list[DeliveredOutputFile]] | None = None
|
||||
# Producers a dispatched run was seeded with and never ran; their values belong to another run.
|
||||
seeded_only_labels_by_run_id: dict[str, frozenset[str]] = field(default_factory=dict)
|
||||
# The browser session the last run actually executed in. On the fresh-session replay path this
|
||||
# is not ctx.browser_session_id, which stays pointed at the debug/scout browser.
|
||||
last_run_blocks_browser_session_id: str | None = None
|
||||
|
|
|
|||
|
|
@ -50,6 +50,7 @@ def _trust_snapshot(ctx: AgentContext) -> dict[str, Any]:
|
|||
"verified_prefix_block_end_session_id": ctx.verified_prefix_block_end_session_id,
|
||||
"verified_prefix_terminal_label": ctx.verified_prefix_terminal_label,
|
||||
"frontier_resume_session_id": ctx.frontier_resume_session_id,
|
||||
"repair_origin_outputs_run_id": ctx.repair_origin_outputs_run_id,
|
||||
"frontier_start_provenance": ctx.frontier_start_provenance,
|
||||
"last_full_workflow_test_ok": ctx.last_full_workflow_test_ok,
|
||||
"last_requested_block_labels": list(ctx.last_requested_block_labels or []),
|
||||
|
|
|
|||
|
|
@ -19,17 +19,29 @@ logs, though a test run's own output can echo one like any run input.
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from collections.abc import Iterable
|
||||
from dataclasses import dataclass, field, replace
|
||||
from datetime import UTC, datetime
|
||||
from enum import StrEnum
|
||||
from typing import Protocol
|
||||
from typing import Any, Protocol
|
||||
|
||||
import structlog
|
||||
|
||||
from skyvern.constants import SCRUBBED_VALUE
|
||||
from skyvern.exceptions import WorkflowRunNotFound
|
||||
from skyvern.forge import app
|
||||
from skyvern.forge.sdk.api.files import is_uploaded_file_id
|
||||
from skyvern.forge.sdk.schemas.workflow_runs import WorkflowRunBlock
|
||||
from skyvern.forge.sdk.workflow.models.parameter import WorkflowParameter
|
||||
from skyvern.forge.sdk.workflow.models.workflow import WorkflowRunParameter, WorkflowRunStatus
|
||||
from skyvern.forge.sdk.workflow.models.workflow import (
|
||||
Workflow,
|
||||
WorkflowDefinition,
|
||||
WorkflowRun,
|
||||
WorkflowRunOutputParameter,
|
||||
WorkflowRunParameter,
|
||||
WorkflowRunStatus,
|
||||
)
|
||||
from skyvern.schemas.proxy_location import ResolvedProxyLocationInput, runtime_proxy_location
|
||||
from skyvern.services.uploaded_file_service import resolve_file_reference
|
||||
|
||||
LOG = structlog.get_logger()
|
||||
|
|
@ -55,6 +67,139 @@ _REFUSAL_SENTENCES: dict[RepairOriginRefusal, str] = {
|
|||
}
|
||||
|
||||
|
||||
class OriginOutputRefusal(StrEnum):
|
||||
FOREIGN_OR_MISMATCHED_ORIGIN = "foreign_or_mismatched_origin"
|
||||
ORIGIN_UNSETTLED = "origin_unsettled"
|
||||
ORDER_UNPROVABLE = "order_unprovable"
|
||||
UPSTREAM_ABSENT = "upstream_absent"
|
||||
UPSTREAM_FAILED = "upstream_failed"
|
||||
CHANGED_PRODUCER = "changed_producer"
|
||||
CHANGED_INPUT = "changed_input"
|
||||
CHANGED_EXECUTION_SETTINGS = "changed_execution_settings"
|
||||
OUTPUT_UNAVAILABLE = "output_unavailable"
|
||||
|
||||
|
||||
OriginInputValue = bool | int | float | str | dict | list
|
||||
|
||||
|
||||
_BINDING_REFUSAL_TO_OUTPUT_REFUSAL: dict[RepairOriginRefusal, OriginOutputRefusal] = {
|
||||
RepairOriginRefusal.RUN_NOT_FOUND: OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN,
|
||||
RepairOriginRefusal.FOREIGN_ORGANIZATION: OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN,
|
||||
RepairOriginRefusal.WORKFLOW_MISMATCH: OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN,
|
||||
}
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class OriginBlockOutput:
|
||||
status: str | None
|
||||
has_value: bool
|
||||
created_at: datetime
|
||||
value: dict | list | str | None = field(default=None, repr=False)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class OriginExecutionSettings:
|
||||
"""The workflow-level settings a run executes under: the run's own override where it set one, else its version's."""
|
||||
|
||||
proxy_location: ResolvedProxyLocationInput
|
||||
browser_profile_id: str | None
|
||||
browser_profile_key: str | None
|
||||
model: dict[str, Any] | None
|
||||
extra_http_headers: dict[str, str] | None = field(repr=False)
|
||||
|
||||
@classmethod
|
||||
def of(cls, workflow: Workflow, run: WorkflowRun | None = None) -> OriginExecutionSettings:
|
||||
run_proxy = run.proxy_location if run is not None else None
|
||||
run_profile_id = run.browser_profile_id if run is not None else None
|
||||
run_headers = run.extra_http_headers if run is not None else None
|
||||
# Saving through Copilot turns absent headers into {}, which sends the same request.
|
||||
headers = (run_headers if run_headers is not None else workflow.extra_http_headers) or None
|
||||
return cls(
|
||||
proxy_location=runtime_proxy_location(run_proxy if run_proxy is not None else workflow.proxy_location),
|
||||
browser_profile_id=run_profile_id if run_profile_id is not None else workflow.browser_profile_id,
|
||||
browser_profile_key=workflow.browser_profile_key,
|
||||
model=workflow.model,
|
||||
extra_http_headers=headers,
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class OriginOutputSnapshot:
|
||||
"""The origin run's latest top-level row per label, beside the definition version that run executed."""
|
||||
|
||||
definition: WorkflowDefinition = field(repr=False)
|
||||
outputs: dict[str, OriginBlockOutput] = field(repr=False)
|
||||
# None when the run's settings were not read, which no test run's settings can be proven equal to.
|
||||
settings: OriginExecutionSettings | None = None
|
||||
# Every input value the run recorded, unfiltered: an output is only reusable with the inputs it was computed from.
|
||||
input_values: dict[str, OriginInputValue] = field(default_factory=dict, repr=False)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class OriginOutputRefusalDetail:
|
||||
reason: OriginOutputRefusal
|
||||
block_label: str
|
||||
origin_workflow_run_id: str | None
|
||||
parameter_key: str | None = None
|
||||
changed_label: str | None = None
|
||||
changed_settings: tuple[str, ...] = ()
|
||||
|
||||
@property
|
||||
def output_key(self) -> str:
|
||||
return f"{self.block_label}_output"
|
||||
|
||||
def as_payload(self) -> dict[str, str | list[str]]:
|
||||
payload: dict[str, str | list[str]] = {
|
||||
"reason": self.reason.value,
|
||||
"block_label": self.block_label,
|
||||
"output_key": self.output_key,
|
||||
}
|
||||
optional = {
|
||||
"origin_workflow_run_id": self.origin_workflow_run_id,
|
||||
"parameter_key": self.parameter_key,
|
||||
"changed_label": self.changed_label,
|
||||
}
|
||||
payload.update({key: value for key, value in optional.items() if value is not None})
|
||||
if self.changed_settings:
|
||||
payload["changed_settings"] = list(self.changed_settings)
|
||||
return payload
|
||||
|
||||
|
||||
def origin_block_outputs_from_rows(
|
||||
definition: WorkflowDefinition,
|
||||
run_blocks: Iterable[WorkflowRunBlock],
|
||||
output_parameter_rows: Iterable[WorkflowRunOutputParameter],
|
||||
input_values: dict[str, OriginInputValue] | None = None,
|
||||
settings: OriginExecutionSettings | None = None,
|
||||
) -> OriginOutputSnapshot:
|
||||
"""Values each block the way verified recording does: its registered output parameter first, even
|
||||
when that value is None, else the row's own output."""
|
||||
output_parameter_ids = {block.label: block.output_parameter.output_parameter_id for block in definition.blocks}
|
||||
registered_by_parameter_id = {
|
||||
row.output_parameter_id: row for row in sorted(output_parameter_rows, key=lambda row: row.created_at)
|
||||
}
|
||||
latest_rows: dict[str, WorkflowRunBlock] = {}
|
||||
for run_block in sorted(run_blocks, key=lambda run_block: run_block.created_at):
|
||||
label = run_block.label
|
||||
if run_block.parent_workflow_run_block_id is None and label is not None and label in output_parameter_ids:
|
||||
latest_rows[label] = run_block
|
||||
outputs: dict[str, OriginBlockOutput] = {}
|
||||
for label, run_block in latest_rows.items():
|
||||
registered = registered_by_parameter_id.get(output_parameter_ids[label])
|
||||
row = OriginBlockOutput(status=run_block.status, has_value=False, created_at=run_block.created_at)
|
||||
if registered is not None:
|
||||
row = replace(row, has_value=True, value=registered.value)
|
||||
elif run_block.output is not None:
|
||||
row = replace(row, has_value=True, value=run_block.output)
|
||||
outputs[label] = row
|
||||
return OriginOutputSnapshot(
|
||||
definition=definition,
|
||||
outputs=outputs,
|
||||
settings=settings,
|
||||
input_values=dict(input_values or {}),
|
||||
)
|
||||
|
||||
|
||||
class RepairTurnContext(Protocol):
|
||||
"""The part of a turn's context this binding reads and writes.
|
||||
|
||||
|
|
@ -73,6 +218,8 @@ class RepairTurnContext(Protocol):
|
|||
last_run_binding_unavailable_reason: str | None
|
||||
repair_origin_input_values: tuple[tuple[WorkflowParameter, WorkflowRunParameter], ...]
|
||||
repair_origin_is_copilot_run: bool
|
||||
repair_origin_outputs: OriginOutputSnapshot | OriginOutputRefusal | None
|
||||
repair_origin_outputs_run_id: str | None
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
|
|
@ -82,6 +229,8 @@ class RepairOriginBinding:
|
|||
refusal: RepairOriginRefusal | None
|
||||
status: WorkflowRunStatus | None = None
|
||||
copilot_run: bool = False
|
||||
debug_run: bool = False
|
||||
run: WorkflowRun | None = field(default=None, repr=False)
|
||||
|
||||
@property
|
||||
def usable(self) -> bool:
|
||||
|
|
@ -95,10 +244,21 @@ class RepairOriginBinding:
|
|||
|
||||
|
||||
def _refused(
|
||||
reason: RepairOriginRefusal, status: WorkflowRunStatus | None = None, *, copilot_run: bool = False
|
||||
reason: RepairOriginRefusal,
|
||||
status: WorkflowRunStatus | None = None,
|
||||
*,
|
||||
copilot_run: bool = False,
|
||||
debug_run: bool = False,
|
||||
run: WorkflowRun | None = None,
|
||||
) -> RepairOriginBinding:
|
||||
return RepairOriginBinding(
|
||||
workflow_run_id=None, browser_session_id=None, refusal=reason, status=status, copilot_run=copilot_run
|
||||
workflow_run_id=None,
|
||||
browser_session_id=None,
|
||||
refusal=reason,
|
||||
status=status,
|
||||
copilot_run=copilot_run,
|
||||
debug_run=debug_run,
|
||||
run=run,
|
||||
)
|
||||
|
||||
|
||||
|
|
@ -127,8 +287,15 @@ async def resolve_repair_origin_binding(
|
|||
if run.workflow_permanent_id != workflow_permanent_id:
|
||||
return _refused(RepairOriginRefusal.WORKFLOW_MISMATCH)
|
||||
copilot_run = run.copilot_session_id is not None
|
||||
debug_run = run.is_debug_session
|
||||
if not run.browser_session_id:
|
||||
return _refused(RepairOriginRefusal.NO_RECORDED_BROWSER, status=run.status, copilot_run=copilot_run)
|
||||
return _refused(
|
||||
RepairOriginRefusal.NO_RECORDED_BROWSER,
|
||||
status=run.status,
|
||||
copilot_run=copilot_run,
|
||||
debug_run=debug_run,
|
||||
run=run,
|
||||
)
|
||||
|
||||
return RepairOriginBinding(
|
||||
workflow_run_id=run.workflow_run_id,
|
||||
|
|
@ -136,6 +303,8 @@ async def resolve_repair_origin_binding(
|
|||
refusal=None,
|
||||
status=run.status,
|
||||
copilot_run=copilot_run,
|
||||
debug_run=debug_run,
|
||||
run=run,
|
||||
)
|
||||
|
||||
|
||||
|
|
@ -154,6 +323,69 @@ async def _file_id_is_reusable(run_parameter: WorkflowRunParameter, organization
|
|||
return await resolve_file_reference(file_id=file_id, organization_id=organization_id) is not None
|
||||
|
||||
|
||||
def _as_utc(value: datetime) -> datetime:
|
||||
return value if value.tzinfo is not None else value.replace(tzinfo=UTC)
|
||||
|
||||
|
||||
async def _load_origin_outputs(
|
||||
binding: RepairOriginBinding,
|
||||
*,
|
||||
workflow_run_id: str,
|
||||
organization_id: str,
|
||||
run_parameters: list[tuple[WorkflowParameter, WorkflowRunParameter]] | None,
|
||||
) -> OriginOutputSnapshot | OriginOutputRefusal:
|
||||
if binding.refusal is not None and binding.refusal is not RepairOriginRefusal.NO_RECORDED_BROWSER:
|
||||
return _BINDING_REFUSAL_TO_OUTPUT_REFUSAL.get(binding.refusal, OriginOutputRefusal.OUTPUT_UNAVAILABLE)
|
||||
if not binding.finished:
|
||||
return OriginOutputRefusal.ORIGIN_UNSETTLED
|
||||
if binding.run is None or run_parameters is None:
|
||||
return OriginOutputRefusal.OUTPUT_UNAVAILABLE
|
||||
# The retention scrubber nulls output values in place and marks the run's failure reason and inputs,
|
||||
# so a scrubbed run's null outputs are not the values its blocks produced.
|
||||
if binding.run.failure_reason == SCRUBBED_VALUE or any(
|
||||
run_parameter.value == SCRUBBED_VALUE for _, run_parameter in run_parameters
|
||||
):
|
||||
return OriginOutputRefusal.OUTPUT_UNAVAILABLE
|
||||
try:
|
||||
origin_workflow = await app.DATABASE.workflows.get_workflow(
|
||||
workflow_id=binding.run.workflow_id, organization_id=organization_id
|
||||
)
|
||||
# Auto-accept overwrites a version in place, so a row changed after the run started may no
|
||||
# longer be the definition the run executed.
|
||||
if origin_workflow is None or _as_utc(origin_workflow.modified_at) > _as_utc(binding.run.created_at):
|
||||
return OriginOutputRefusal.OUTPUT_UNAVAILABLE
|
||||
# A test of selected blocks never loads a cached script, so an output a script produced is not
|
||||
# what the test's agent run of the same block would produce. `script_run` records that a script
|
||||
# actually loaded, including runs a rollout upgraded to code without declaring run_with.
|
||||
if binding.run.script_run is not None:
|
||||
return OriginOutputRefusal.CHANGED_EXECUTION_SETTINGS
|
||||
# A child run's LLM blocks also render every ancestor's workflow_system_prompt, which a
|
||||
# standalone repair test never inherits.
|
||||
if binding.run.parent_workflow_run_id is not None:
|
||||
return OriginOutputRefusal.CHANGED_EXECUTION_SETTINGS
|
||||
run_blocks = await app.DATABASE.observer.get_workflow_run_blocks(
|
||||
workflow_run_id=workflow_run_id, organization_id=organization_id
|
||||
)
|
||||
output_parameter_rows = await app.DATABASE.workflow_runs.get_workflow_run_output_parameters(
|
||||
workflow_run_id=workflow_run_id
|
||||
)
|
||||
return origin_block_outputs_from_rows(
|
||||
origin_workflow.workflow_definition,
|
||||
run_blocks,
|
||||
output_parameter_rows,
|
||||
input_values={parameter.key: run_parameter.value for parameter, run_parameter in run_parameters},
|
||||
settings=OriginExecutionSettings.of(origin_workflow, binding.run),
|
||||
)
|
||||
except Exception as exc:
|
||||
# No exc_info: a row that fails to parse is quoted in its own exception message.
|
||||
LOG.warning(
|
||||
"copilot_repair_origin_outputs_unavailable",
|
||||
workflow_run_id=workflow_run_id,
|
||||
error_type=type(exc).__name__,
|
||||
)
|
||||
return OriginOutputRefusal.OUTPUT_UNAVAILABLE
|
||||
|
||||
|
||||
async def seed_repair_origin_run(ctx: RepairTurnContext, *, workflow_run_id: str | None) -> RepairOriginBinding:
|
||||
"""Seed the turn's last-run identity from the run it was opened about, before it acts.
|
||||
|
||||
|
|
@ -178,7 +410,9 @@ async def seed_repair_origin_run(ctx: RepairTurnContext, *, workflow_run_id: str
|
|||
_REFUSAL_SENTENCES.get(binding.refusal) if binding.refusal is not None else None
|
||||
)
|
||||
origin_input_values: tuple[tuple[WorkflowParameter, WorkflowRunParameter], ...] = ()
|
||||
if workflow_run_id and binding.refusal in (None, RepairOriginRefusal.NO_RECORDED_BROWSER):
|
||||
owned_run = bool(workflow_run_id) and binding.refusal in (None, RepairOriginRefusal.NO_RECORDED_BROWSER)
|
||||
loaded: list[tuple[WorkflowParameter, WorkflowRunParameter]] | None = None
|
||||
if workflow_run_id and owned_run:
|
||||
try:
|
||||
loaded = await app.DATABASE.workflow_runs.get_workflow_run_parameters(workflow_run_id=workflow_run_id)
|
||||
origin_input_values = tuple(
|
||||
|
|
@ -193,6 +427,21 @@ async def seed_repair_origin_run(ctx: RepairTurnContext, *, workflow_run_id: str
|
|||
)
|
||||
ctx.repair_origin_input_values = origin_input_values
|
||||
ctx.repair_origin_is_copilot_run = binding.copilot_run
|
||||
# A Copilot test run or a debugger block run may itself have been seeded, so neither is ever an
|
||||
# output origin and the turn plans exactly as with no origin.
|
||||
seeded_run = binding.copilot_run or binding.debug_run
|
||||
ctx.repair_origin_outputs_run_id = workflow_run_id if owned_run and not seeded_run else None
|
||||
ctx.repair_origin_outputs = (
|
||||
await _load_origin_outputs(
|
||||
binding,
|
||||
workflow_run_id=workflow_run_id,
|
||||
organization_id=ctx.organization_id,
|
||||
run_parameters=loaded,
|
||||
)
|
||||
if workflow_run_id and not seeded_run
|
||||
else None
|
||||
)
|
||||
origin_outputs = ctx.repair_origin_outputs
|
||||
LOG.info(
|
||||
"copilot_repair_origin_binding",
|
||||
requested_workflow_run_id=workflow_run_id,
|
||||
|
|
@ -201,5 +450,7 @@ async def seed_repair_origin_run(ctx: RepairTurnContext, *, workflow_run_id: str
|
|||
workflow_run_id=binding.workflow_run_id,
|
||||
origin_input_keys=[parameter.key for parameter, _ in origin_input_values],
|
||||
origin_is_copilot_run=binding.copilot_run,
|
||||
origin_output_labels=sorted(origin_outputs.outputs) if isinstance(origin_outputs, OriginOutputSnapshot) else [],
|
||||
origin_output_refusal=origin_outputs if isinstance(origin_outputs, OriginOutputRefusal) else None,
|
||||
)
|
||||
return binding
|
||||
|
|
|
|||
|
|
@ -88,6 +88,11 @@ if TYPE_CHECKING:
|
|||
from skyvern.forge.sdk.copilot.completion_verification import CompletionVerificationResult
|
||||
from skyvern.forge.sdk.copilot.context import CodeAuthoringRepairContext, SignedOutPageObservation
|
||||
from skyvern.forge.sdk.copilot.mcp_adapter import SkyvernOverlayMCPServer
|
||||
from skyvern.forge.sdk.copilot.repair_origin_run import (
|
||||
OriginOutputRefusal,
|
||||
OriginOutputRefusalDetail,
|
||||
OriginOutputSnapshot,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.request_policy import RequestPolicy
|
||||
from skyvern.forge.sdk.copilot.result_evidence import ScoutObservationContract
|
||||
from skyvern.forge.sdk.copilot.run_outcome import RecordedRunOutcome
|
||||
|
|
@ -697,6 +702,12 @@ class AgentContext:
|
|||
default=(), repr=False
|
||||
)
|
||||
repair_origin_is_copilot_run: bool = False
|
||||
# Block outputs the same run recorded, or why they cannot be used; loaded only when a run was named.
|
||||
repair_origin_outputs: OriginOutputSnapshot | OriginOutputRefusal | None = field(default=None, repr=False)
|
||||
repair_origin_outputs_run_id: str | None = None
|
||||
# Set by the planner, consumed and cleared by the run it planned, like frontier_resume_session_id.
|
||||
frontier_origin_reused_labels: list[str] = field(default_factory=list)
|
||||
frontier_origin_output_refusal: OriginOutputRefusalDetail | None = None
|
||||
# Source page of an in-flight scout action, captured before it may navigate away.
|
||||
pending_scout_source_url: str | None = None
|
||||
# The withheld page's URL, read before a navigation meant to leave it, keyed by the browser the
|
||||
|
|
|
|||
|
|
@ -84,6 +84,7 @@ from skyvern.forge.sdk.copilot.workflow_yaml import (
|
|||
stored_block_code,
|
||||
stored_workflow_yaml,
|
||||
)
|
||||
from skyvern.forge.sdk.workflow.models.workflow import WorkflowDefinition
|
||||
from skyvern.utils.yaml_loader import dump_workflow_yaml
|
||||
|
||||
from ._shared import _COMPOSITION_STRIPPED_HTML_MAX_CHARS as _COMPOSITION_STRIPPED_HTML_MAX_CHARS
|
||||
|
|
@ -1575,15 +1576,13 @@ async def _run_updated_workflow_blocks(
|
|||
tool_name: str,
|
||||
arguments: dict[str, Any],
|
||||
update_result: dict[str, Any],
|
||||
prior_definition: object | None,
|
||||
prior_definition: WorkflowDefinition | None,
|
||||
block_labels: list[str],
|
||||
parameters: dict[str, Any],
|
||||
handler_start: float,
|
||||
) -> str:
|
||||
"""Run a just-persisted definition through the shared frontier and debug-evidence seam."""
|
||||
new_definition = None
|
||||
if copilot_ctx.last_workflow is not None:
|
||||
new_definition = getattr(copilot_ctx.last_workflow, "workflow_definition", None)
|
||||
new_definition = copilot_ctx.last_workflow.workflow_definition if copilot_ctx.last_workflow is not None else None
|
||||
|
||||
labels_to_execute, block_outputs_to_seed, frontier_start_label, start_provenance = _plan_frontier(
|
||||
copilot_ctx,
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ import json
|
|||
import re
|
||||
import textwrap
|
||||
from collections import Counter
|
||||
from dataclasses import asdict
|
||||
from typing import Any, Protocol, runtime_checkable
|
||||
from urllib.parse import urlparse
|
||||
|
||||
|
|
@ -18,7 +19,9 @@ except ImportError: # pragma: no cover — bs4 is a transitive dep but discover
|
|||
BeautifulSoup = None # type: ignore[assignment, misc]
|
||||
from jinja2 import TemplateError, meta
|
||||
from jinja2.sandbox import SandboxedEnvironment
|
||||
from pydantic import JsonValue
|
||||
|
||||
from skyvern.constants import SCRUBBED_VALUE
|
||||
from skyvern.forge import app
|
||||
from skyvern.forge.sdk.copilot.block_goal_wrapping import wrap_workflow_block_goals
|
||||
from skyvern.forge.sdk.copilot.build_test_outcome import RecordedBuildTestOutcome
|
||||
|
|
@ -36,20 +39,29 @@ from skyvern.forge.sdk.copilot.frontier_provenance_dump import (
|
|||
write_packet,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.output_utils import INTERNAL_VALIDATION_FAILURE_PREFIX
|
||||
from skyvern.forge.sdk.copilot.repair_origin_run import (
|
||||
OriginBlockOutput,
|
||||
OriginExecutionSettings,
|
||||
OriginInputValue,
|
||||
OriginOutputRefusal,
|
||||
OriginOutputRefusalDetail,
|
||||
OriginOutputSnapshot,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.runtime import (
|
||||
AgentContext,
|
||||
FrontierStartProvenance,
|
||||
resolve_browser_state_for_context,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.turn_origin import TurnOrigin
|
||||
from skyvern.forge.sdk.copilot.workflow_yaml import _process_workflow_yaml
|
||||
from skyvern.forge.sdk.copilot.workflow_yaml import _process_workflow_yaml, copilot_round_trip_definition
|
||||
from skyvern.forge.sdk.workflow.models.block import BlockTypeVar, get_all_blocks
|
||||
from skyvern.forge.sdk.workflow.models.parameter import (
|
||||
RESERVED_PARAMETER_KEYS,
|
||||
Parameter,
|
||||
WorkflowParameter,
|
||||
)
|
||||
from skyvern.forge.sdk.workflow.models.workflow import Workflow
|
||||
from skyvern.schemas.workflows import BlockType
|
||||
from skyvern.forge.sdk.workflow.models.workflow import Workflow, WorkflowDefinition
|
||||
from skyvern.schemas.workflows import BlockStatus, BlockType
|
||||
|
||||
from ._shared import (
|
||||
_block_type_name,
|
||||
|
|
@ -763,6 +775,19 @@ def _jinja_roots(serialized: str) -> set[str]:
|
|||
return roots
|
||||
|
||||
|
||||
def _root_parameter_key(root: str, parameter_keys: set[str]) -> str | None:
|
||||
if root in parameter_keys:
|
||||
return root
|
||||
return next(
|
||||
(
|
||||
root[: -len(suffix)]
|
||||
for suffix in _CREDENTIAL_REAL_VALUE_SUFFIXES
|
||||
if root.endswith(suffix) and root[: -len(suffix)] in parameter_keys
|
||||
),
|
||||
None,
|
||||
)
|
||||
|
||||
|
||||
def _classify_frontier_jinja_refs(
|
||||
frontier_labels: list[str],
|
||||
new_definition: object | None,
|
||||
|
|
@ -790,12 +815,7 @@ def _classify_frontier_jinja_refs(
|
|||
continue
|
||||
if root in known_labels:
|
||||
block_form_refs.add(root)
|
||||
if root in known_roots:
|
||||
continue
|
||||
if any(
|
||||
root.endswith(suffix) and root[: -len(suffix)] in parameter_keys
|
||||
for suffix in _CREDENTIAL_REAL_VALUE_SUFFIXES
|
||||
):
|
||||
if root in known_roots or _root_parameter_key(root, parameter_keys) is not None:
|
||||
continue
|
||||
unknown_roots.add(root)
|
||||
|
||||
|
|
@ -820,14 +840,11 @@ def _unknown_jinja_roots(
|
|||
return unknown_roots
|
||||
|
||||
|
||||
async def _get_prior_workflow_definition(ctx: AgentContext) -> object | None:
|
||||
async def _get_prior_workflow_definition(ctx: AgentContext) -> WorkflowDefinition | None:
|
||||
"""Hybrid: prefer ctx.last_workflow, fall back to DB fetch on cold start."""
|
||||
last_workflow = getattr(ctx, "last_workflow", None)
|
||||
if last_workflow is not None:
|
||||
definition = getattr(last_workflow, "workflow_definition", None)
|
||||
if definition is not None:
|
||||
return definition
|
||||
last_yaml = getattr(ctx, "last_workflow_yaml", None)
|
||||
if ctx.last_workflow is not None:
|
||||
return ctx.last_workflow.workflow_definition
|
||||
last_yaml = ctx.last_workflow_yaml
|
||||
if last_yaml:
|
||||
try:
|
||||
workflow = await _process_workflow_yaml(
|
||||
|
|
@ -1076,15 +1093,16 @@ def _anchored_plan(
|
|||
def _plan_frontier(
|
||||
ctx: AgentContext,
|
||||
requested_labels: list[str],
|
||||
old_definition: object | None,
|
||||
new_definition: object | None,
|
||||
old_definition: WorkflowDefinition | None,
|
||||
new_definition: WorkflowDefinition | None,
|
||||
runtime_page_url: str | None = None,
|
||||
) -> FrontierPlan:
|
||||
if frontier_dump_root() is None:
|
||||
return _plan_frontier_uncaptured(ctx, requested_labels, old_definition, new_definition, runtime_page_url)
|
||||
before = trust_snapshot(ctx)
|
||||
plan = _plan_frontier_uncaptured(ctx, requested_labels, old_definition, new_definition, runtime_page_url)
|
||||
labels_to_execute, _seed, frontier_start_label, start_provenance = plan
|
||||
labels_to_execute, seed, frontier_start_label, start_provenance = plan
|
||||
refusal = ctx.frontier_origin_output_refusal
|
||||
write_packet(
|
||||
"frontier",
|
||||
{
|
||||
|
|
@ -1097,6 +1115,9 @@ def _plan_frontier(
|
|||
"labels_to_execute": list(labels_to_execute),
|
||||
"frontier_start_label": frontier_start_label,
|
||||
"frontier_start_provenance": start_provenance,
|
||||
"seed_labels": sorted(seed),
|
||||
"origin_reused_labels": list(ctx.frontier_origin_reused_labels),
|
||||
"origin_output_refusal": refusal.as_payload() if refusal is not None else None,
|
||||
},
|
||||
},
|
||||
)
|
||||
|
|
@ -1106,8 +1127,300 @@ def _plan_frontier(
|
|||
def _plan_frontier_uncaptured(
|
||||
ctx: AgentContext,
|
||||
requested_labels: list[str],
|
||||
old_definition: object | None,
|
||||
new_definition: object | None,
|
||||
old_definition: WorkflowDefinition | None,
|
||||
new_definition: WorkflowDefinition | None,
|
||||
runtime_page_url: str | None = None,
|
||||
) -> FrontierPlan:
|
||||
ctx.frontier_origin_reused_labels = []
|
||||
ctx.frontier_origin_output_refusal = None
|
||||
plan = _plan_frontier_base(ctx, requested_labels, old_definition, new_definition, runtime_page_url)
|
||||
return _fill_seed_from_origin(ctx, plan, new_definition)
|
||||
|
||||
|
||||
def _without_per_version_parameter_ids(value: JsonValue) -> JsonValue:
|
||||
# Each persisted version mints its own parameter rows, so an unchanged block's nested parameters
|
||||
# differ between versions in their ids, owning workflow_id and timestamps.
|
||||
if isinstance(value, dict):
|
||||
volatile = _PARAMETER_FINGERPRINT_VOLATILE_KEYS | {"workflow_id"} if "parameter_type" in value else frozenset()
|
||||
return {key: _without_per_version_parameter_ids(item) for key, item in value.items() if key not in volatile}
|
||||
if isinstance(value, list):
|
||||
return [_without_per_version_parameter_ids(item) for item in value]
|
||||
return value
|
||||
|
||||
|
||||
# Routing and display fields: storage order already proves the producer runs first, and neither
|
||||
# changes what the producer outputs.
|
||||
_PRODUCER_SHAPE_IGNORED_KEYS = frozenset({"next_block_label", "title"})
|
||||
|
||||
|
||||
def _producer_shape(config: dict[str, JsonValue]) -> JsonValue:
|
||||
# The export schema shapes the output only while export is on.
|
||||
ignored = (
|
||||
_PRODUCER_SHAPE_IGNORED_KEYS
|
||||
if config.get("export_enabled")
|
||||
else _PRODUCER_SHAPE_IGNORED_KEYS | {"export_data_schema"}
|
||||
)
|
||||
return _without_per_version_parameter_ids({key: value for key, value in config.items() if key not in ignored})
|
||||
|
||||
|
||||
def _cross_version_block_shapes(definition: WorkflowDefinition, workflow_id: str) -> dict[str, JsonValue] | None:
|
||||
# Both versions go through Copilot's own save path first, so a version authored elsewhere (the API,
|
||||
# say) is compared by what Copilot would persist rather than by how its author serialized it.
|
||||
try:
|
||||
normalized = copilot_round_trip_definition(definition, workflow_id=workflow_id)
|
||||
except Exception as exc:
|
||||
LOG.info("copilot_frontier_origin_producer_shape_unavailable", error_type=type(exc).__name__)
|
||||
return None
|
||||
return {
|
||||
label: _producer_shape(_canonical_block_config(block)) for label, block in _blocks_by_label(normalized).items()
|
||||
}
|
||||
|
||||
|
||||
def _upstream_prefix(producer: str, definition: WorkflowDefinition) -> list[str]:
|
||||
"""The producer and every block stored before it: what its output was computed after."""
|
||||
labels = _executable_workflow_block_labels(definition)
|
||||
return labels[: labels.index(producer) + 1] if producer in labels else []
|
||||
|
||||
|
||||
def _eligible_origin_output(
|
||||
ctx: AgentContext,
|
||||
producer: str,
|
||||
first_executed_label: str,
|
||||
new_definition: WorkflowDefinition,
|
||||
shapes: tuple[dict[str, JsonValue] | None, dict[str, JsonValue] | None],
|
||||
) -> OriginBlockOutput | OriginOutputRefusalDetail:
|
||||
def refused(reason: OriginOutputRefusal, changed_label: str | None = None) -> OriginOutputRefusalDetail:
|
||||
return OriginOutputRefusalDetail(
|
||||
reason=reason,
|
||||
block_label=producer,
|
||||
origin_workflow_run_id=ctx.repair_origin_outputs_run_id,
|
||||
changed_label=changed_label if changed_label != producer else None,
|
||||
)
|
||||
|
||||
snapshot = ctx.repair_origin_outputs
|
||||
if not isinstance(snapshot, OriginOutputSnapshot):
|
||||
return refused(snapshot or OriginOutputRefusal.OUTPUT_UNAVAILABLE)
|
||||
workflow_labels = _executable_workflow_block_labels(new_definition)
|
||||
if (
|
||||
producer not in workflow_labels
|
||||
or first_executed_label not in workflow_labels
|
||||
or workflow_labels.index(producer) >= workflow_labels.index(first_executed_label)
|
||||
or not _storage_order_is_traversal_order(new_definition, first_executed_label)
|
||||
or not _storage_order_is_traversal_order(snapshot.definition, producer)
|
||||
):
|
||||
return refused(OriginOutputRefusal.ORDER_UNPROVABLE)
|
||||
recorded = snapshot.outputs.get(producer)
|
||||
if recorded is None:
|
||||
return refused(OriginOutputRefusal.UPSTREAM_ABSENT)
|
||||
if recorded.status != BlockStatus.completed:
|
||||
return refused(OriginOutputRefusal.UPSTREAM_FAILED)
|
||||
origin_shapes, candidate_shapes = shapes
|
||||
if origin_shapes is None or candidate_shapes is None:
|
||||
return refused(OriginOutputRefusal.OUTPUT_UNAVAILABLE)
|
||||
# Every LLM-driven block renders the workflow-level prompt, so a change there makes the output stale too.
|
||||
if (snapshot.definition.workflow_system_prompt or None) != (new_definition.workflow_system_prompt or None):
|
||||
return refused(OriginOutputRefusal.CHANGED_PRODUCER)
|
||||
# Like verified-state invalidation, a change anywhere upstream makes the producer's output stale.
|
||||
prefix = _upstream_prefix(producer, new_definition)
|
||||
origin_labels = _executable_workflow_block_labels(snapshot.definition)
|
||||
for index, label in enumerate(prefix):
|
||||
origin_shape = origin_shapes.get(label)
|
||||
if (
|
||||
index >= len(origin_labels)
|
||||
or origin_labels[index] != label
|
||||
or origin_shape is None
|
||||
or origin_shape != candidate_shapes.get(label)
|
||||
):
|
||||
return refused(OriginOutputRefusal.CHANGED_PRODUCER, changed_label=label)
|
||||
# The output was computed after the prefix only if the origin itself ran each prefix block first.
|
||||
for label in prefix[:-1]:
|
||||
predecessor = snapshot.outputs.get(label)
|
||||
if predecessor is None:
|
||||
return refused(OriginOutputRefusal.UPSTREAM_ABSENT, changed_label=label)
|
||||
if predecessor.status != BlockStatus.completed:
|
||||
return refused(OriginOutputRefusal.UPSTREAM_FAILED, changed_label=label)
|
||||
if predecessor.created_at > recorded.created_at:
|
||||
return refused(OriginOutputRefusal.ORDER_UNPROVABLE, changed_label=label)
|
||||
# The retention scrubber nulls output values in place, and a completed run with no inputs keeps no
|
||||
# scrub marker, so a stored null cannot be told apart from a scrubbed value.
|
||||
if not recorded.has_value or recorded.value is None:
|
||||
return refused(OriginOutputRefusal.OUTPUT_UNAVAILABLE)
|
||||
return recorded
|
||||
|
||||
|
||||
def _fill_seed_from_origin(
|
||||
ctx: AgentContext, plan: FrontierPlan, new_definition: WorkflowDefinition | None
|
||||
) -> FrontierPlan:
|
||||
"""Seed each producer the executed blocks reference but do not run: this turn's verified output
|
||||
first, else the output the origin run recorded; otherwise the run dispatches unseeded and reports
|
||||
why. Only a turn whose bound run loaded origin outputs is filled, and never a resumed plan."""
|
||||
labels_to_execute, seed, frontier_start_label, start_provenance = plan
|
||||
origin_run_id = ctx.repair_origin_outputs_run_id
|
||||
if (
|
||||
ctx.repair_origin_outputs is None
|
||||
or new_definition is None
|
||||
or ctx.frontier_resume_session_id
|
||||
or not labels_to_execute
|
||||
):
|
||||
return plan
|
||||
missing = _referenced_output_labels(labels_to_execute, new_definition) - set(labels_to_execute) - set(seed)
|
||||
if not missing:
|
||||
return plan
|
||||
verified_outputs = ctx.verified_block_outputs or {}
|
||||
shapes: tuple[dict[str, JsonValue] | None, dict[str, JsonValue] | None] = (None, None)
|
||||
if isinstance(ctx.repair_origin_outputs, OriginOutputSnapshot) and not missing.issubset(verified_outputs):
|
||||
shapes = (
|
||||
_cross_version_block_shapes(ctx.repair_origin_outputs.definition, ctx.workflow_id),
|
||||
_cross_version_block_shapes(new_definition, ctx.workflow_id),
|
||||
)
|
||||
filled = dict(seed)
|
||||
reused: list[str] = []
|
||||
for producer in [label for label in _workflow_definition_block_labels(new_definition) if label in missing]:
|
||||
if producer in verified_outputs:
|
||||
filled[producer] = verified_outputs[producer]
|
||||
continue
|
||||
eligible = _eligible_origin_output(ctx, producer, labels_to_execute[0], new_definition, shapes)
|
||||
if isinstance(eligible, OriginOutputRefusalDetail):
|
||||
ctx.frontier_origin_output_refusal = logged_origin_refusal(eligible)
|
||||
return plan
|
||||
filled[producer] = eligible.value
|
||||
reused.append(producer)
|
||||
ctx.frontier_origin_reused_labels = reused
|
||||
if reused:
|
||||
LOG.info("copilot_frontier_origin_outputs_reused", block_labels=reused, origin_workflow_run_id=origin_run_id)
|
||||
return labels_to_execute, filled, frontier_start_label, start_provenance
|
||||
|
||||
|
||||
def _producer_input_keys(producer: str, definition: WorkflowDefinition) -> set[str]:
|
||||
"""Run inputs the producer reads: workflow parameters it binds or its templates name."""
|
||||
block = _blocks_by_label(definition).get(producer)
|
||||
if block is None:
|
||||
return set()
|
||||
input_keys = {parameter.key for parameter in definition.parameters if isinstance(parameter, WorkflowParameter)}
|
||||
referenced: set[str] = set()
|
||||
stack: list[JsonValue] = [_canonical_block_config(block)]
|
||||
while stack:
|
||||
item = stack.pop()
|
||||
if isinstance(item, dict):
|
||||
key = item.get("key")
|
||||
if isinstance(key, str) and key in input_keys and "parameter_type" in item:
|
||||
referenced.add(key)
|
||||
stack.extend(item.values())
|
||||
elif isinstance(item, list):
|
||||
stack.extend(item)
|
||||
for serialized in _serialized_frontier_block_configs([producer], definition):
|
||||
referenced.update(
|
||||
key for root in _jinja_roots(serialized) if (key := _root_parameter_key(root, input_keys)) is not None
|
||||
)
|
||||
return referenced
|
||||
|
||||
|
||||
def _workflow_prompt_input_keys(definition: WorkflowDefinition) -> set[str]:
|
||||
input_keys = {parameter.key for parameter in definition.parameters if isinstance(parameter, WorkflowParameter)}
|
||||
roots = _jinja_roots(definition.workflow_system_prompt) if definition.workflow_system_prompt else set()
|
||||
return {key for root in roots if (key := _root_parameter_key(root, input_keys)) is not None}
|
||||
|
||||
|
||||
def origin_input_refusal(
|
||||
ctx: AgentContext,
|
||||
reused_labels: list[str],
|
||||
definition: WorkflowDefinition,
|
||||
resolved_inputs: dict[str, Any],
|
||||
) -> OriginOutputRefusalDetail | None:
|
||||
"""Refuse a reused origin output when the producer or a block before it would read a different input
|
||||
in this test run than in the origin run, or when the origin input cannot be proven."""
|
||||
snapshot = ctx.repair_origin_outputs
|
||||
origin_run_id = ctx.repair_origin_outputs_run_id
|
||||
origin_parameters = (
|
||||
{p.key: p for p in snapshot.definition.parameters if isinstance(p, WorkflowParameter)}
|
||||
if isinstance(snapshot, OriginOutputSnapshot)
|
||||
else {}
|
||||
)
|
||||
for producer in reused_labels:
|
||||
upstream_keys = set().union(
|
||||
*(_producer_input_keys(label, definition) for label in _upstream_prefix(producer, definition)),
|
||||
_workflow_prompt_input_keys(definition),
|
||||
)
|
||||
for key in sorted(upstream_keys):
|
||||
origin_parameter = origin_parameters.get(key)
|
||||
if not isinstance(snapshot, OriginOutputSnapshot) or origin_parameter is None:
|
||||
proven = False
|
||||
elif key in snapshot.input_values:
|
||||
origin_value = snapshot.input_values[key]
|
||||
proven = origin_value != SCRUBBED_VALUE and _same_input(origin_value, resolved_inputs.get(key))
|
||||
else:
|
||||
proven = _same_input(origin_parameter.default_value, resolved_inputs.get(key))
|
||||
if not proven:
|
||||
return logged_origin_refusal(
|
||||
OriginOutputRefusalDetail(
|
||||
reason=OriginOutputRefusal.CHANGED_INPUT,
|
||||
block_label=producer,
|
||||
origin_workflow_run_id=origin_run_id,
|
||||
parameter_key=key,
|
||||
)
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
def origin_definition_refusal(
|
||||
ctx: AgentContext, reused_labels: list[str], first_executed_label: str, definition: WorkflowDefinition
|
||||
) -> OriginOutputRefusalDetail | None:
|
||||
"""Re-prove each reused output against the definition that will be dispatched, which the planner
|
||||
did not necessarily read."""
|
||||
if not reused_labels:
|
||||
return None
|
||||
snapshot = ctx.repair_origin_outputs
|
||||
shapes = (
|
||||
_cross_version_block_shapes(snapshot.definition, ctx.workflow_id)
|
||||
if isinstance(snapshot, OriginOutputSnapshot)
|
||||
else None,
|
||||
_cross_version_block_shapes(definition, ctx.workflow_id),
|
||||
)
|
||||
for producer in reused_labels:
|
||||
eligible = _eligible_origin_output(ctx, producer, first_executed_label, definition, shapes)
|
||||
if isinstance(eligible, OriginOutputRefusalDetail):
|
||||
return logged_origin_refusal(eligible)
|
||||
return None
|
||||
|
||||
|
||||
def origin_settings_refusal(
|
||||
ctx: AgentContext, reused_labels: list[str], workflow: Workflow
|
||||
) -> OriginOutputRefusalDetail | None:
|
||||
"""Refuse reused origin outputs when the test would run under other workflow-level execution settings
|
||||
than the origin run did, or when the origin's cannot be read; it names the settings, never their values."""
|
||||
if not reused_labels:
|
||||
return None
|
||||
snapshot = ctx.repair_origin_outputs
|
||||
origin_settings = snapshot.settings if isinstance(snapshot, OriginOutputSnapshot) else None
|
||||
origin = asdict(origin_settings) if origin_settings is not None else {}
|
||||
candidate = asdict(OriginExecutionSettings.of(workflow))
|
||||
changed = tuple(name for name, value in candidate.items() if name not in origin or origin[name] != value)
|
||||
if not changed:
|
||||
return None
|
||||
return logged_origin_refusal(
|
||||
OriginOutputRefusalDetail(
|
||||
reason=OriginOutputRefusal.CHANGED_EXECUTION_SETTINGS,
|
||||
block_label=reused_labels[0],
|
||||
origin_workflow_run_id=ctx.repair_origin_outputs_run_id,
|
||||
changed_settings=changed,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def logged_origin_refusal(detail: OriginOutputRefusalDetail) -> OriginOutputRefusalDetail:
|
||||
LOG.info("copilot_frontier_origin_output_refused", **detail.as_payload())
|
||||
return detail
|
||||
|
||||
|
||||
def _same_input(origin_value: OriginInputValue | None, test_value: OriginInputValue | None) -> bool:
|
||||
return json.dumps(origin_value, sort_keys=True, default=str) == json.dumps(test_value, sort_keys=True, default=str)
|
||||
|
||||
|
||||
def _plan_frontier_base(
|
||||
ctx: AgentContext,
|
||||
requested_labels: list[str],
|
||||
old_definition: WorkflowDefinition | None,
|
||||
new_definition: WorkflowDefinition | None,
|
||||
runtime_page_url: str | None = None,
|
||||
) -> FrontierPlan:
|
||||
"""Plan the frontier execution.
|
||||
|
|
|
|||
|
|
@ -148,6 +148,7 @@ from skyvern.forge.sdk.copilot.output_utils import (
|
|||
screened_recorded_url,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.reached_download_target import generated_file_artifact_ids
|
||||
from skyvern.forge.sdk.copilot.repair_origin_run import OriginOutputRefusal, OriginOutputRefusalDetail
|
||||
from skyvern.forge.sdk.copilot.review_gate import workflow_block_fingerprints
|
||||
from skyvern.forge.sdk.copilot.run_outcome import (
|
||||
TERMINAL_CHALLENGE_RUN_OUTCOME_REASON_CODE,
|
||||
|
|
@ -292,6 +293,10 @@ from .frontier import (
|
|||
_workflow_with_runtime_block_goal_context,
|
||||
_workflow_with_runtime_frontier_anchor,
|
||||
_workflow_with_runtime_frontier_starter_url_seed,
|
||||
logged_origin_refusal,
|
||||
origin_definition_refusal,
|
||||
origin_input_refusal,
|
||||
origin_settings_refusal,
|
||||
)
|
||||
from .guardrails import (
|
||||
_authority_tool_error,
|
||||
|
|
@ -1663,6 +1668,9 @@ class _RunExecution:
|
|||
unbound_keys: list[str]
|
||||
explicit_blank: bool
|
||||
reused_origin_input_keys: list[str] = dataclass_field(default_factory=list)
|
||||
reused_origin_output_labels: list[str] = dataclass_field(default_factory=list)
|
||||
origin_workflow_run_id: str | None = None
|
||||
origin_output_not_reused: OriginOutputRefusalDetail | None = None
|
||||
proposal_owner_turn_id: str | None = None
|
||||
proposal_revision: int | None = None
|
||||
outcome: RecordedRunOutcome | None = None
|
||||
|
|
@ -1798,6 +1806,7 @@ async def _attach_registered_output_parameter_values(
|
|||
workflow: Workflow | None,
|
||||
data: dict[str, Any],
|
||||
persisted_output_parameters: list[Any] | None = None,
|
||||
excluded_block_labels: frozenset[str] = frozenset(),
|
||||
) -> dict[str, Any]:
|
||||
try:
|
||||
registered_rows = await app.DATABASE.workflow_runs.get_workflow_run_output_parameters(
|
||||
|
|
@ -1838,6 +1847,8 @@ async def _attach_registered_output_parameter_values(
|
|||
block_info["output_parameter_key"] = output_parameter_key
|
||||
if output_parameter_key and not block_info.get("block_label"):
|
||||
block_info.update(index_by_key.get(output_parameter_key, {}))
|
||||
if block_info.get("block_label") in excluded_block_labels:
|
||||
continue
|
||||
value = getattr(row, "value", None)
|
||||
item = {
|
||||
"workflow_run_id": workflow_run_id,
|
||||
|
|
@ -3191,6 +3202,25 @@ async def _halt_turn_if_superseded(
|
|||
stash_build_test_superseded_halt(ctx, workflow_run_id=workflow_run_id)
|
||||
|
||||
|
||||
def _attach_reused_origin_outputs(data: dict[str, Any], execution: _RunExecution) -> None:
|
||||
if execution.reused_origin_output_labels:
|
||||
data["reused_origin_output_labels"] = list(execution.reused_origin_output_labels)
|
||||
data["origin_workflow_run_id"] = execution.origin_workflow_run_id
|
||||
if execution.origin_output_not_reused is not None:
|
||||
data["origin_output_not_reused"] = execution.origin_output_not_reused.as_payload()
|
||||
|
||||
|
||||
def _stop_reusing_origin_outputs(
|
||||
execution: _RunExecution, block_outputs_to_seed: dict[str, Any], refusal: OriginOutputRefusalDetail
|
||||
) -> dict[str, Any]:
|
||||
"""Drop the origin seeds so the test runs as it would with no origin, and keep why as a fact."""
|
||||
reused = set(execution.reused_origin_output_labels)
|
||||
execution.reused_origin_output_labels = []
|
||||
execution.origin_workflow_run_id = None
|
||||
execution.origin_output_not_reused = refusal
|
||||
return {label: value for label, value in block_outputs_to_seed.items() if label not in reused}
|
||||
|
||||
|
||||
async def _run_blocks_and_collect_debug(
|
||||
params: dict[str, Any],
|
||||
ctx: CopilotContext,
|
||||
|
|
@ -3221,10 +3251,14 @@ async def _run_blocks_and_collect_debug(
|
|||
start_provenance: FrontierStartProvenance = (
|
||||
"initial" if explicit_blank else ctx.frontier_start_provenance or "unanchored"
|
||||
)
|
||||
origin_output_refusal = None if explicit_blank else ctx.frontier_origin_output_refusal
|
||||
reused_origin_output_labels = [] if explicit_blank else list(ctx.frontier_origin_reused_labels)
|
||||
if not explicit_blank:
|
||||
ctx.frontier_resume_session_id = None
|
||||
ctx.frontier_requires_own_browser = False
|
||||
ctx.frontier_start_provenance = None
|
||||
ctx.frontier_origin_output_refusal = None
|
||||
ctx.frontier_origin_reused_labels = []
|
||||
|
||||
block_labels = params["block_labels"]
|
||||
if not block_labels:
|
||||
|
|
@ -3270,6 +3304,8 @@ async def _run_blocks_and_collect_debug(
|
|||
proposal_revision=ctx.proposal_revision if snapshot.provenance == "staged" else None,
|
||||
unbound_keys=[],
|
||||
explicit_blank=explicit_blank,
|
||||
reused_origin_output_labels=reused_origin_output_labels,
|
||||
origin_workflow_run_id=ctx.repair_origin_outputs_run_id if reused_origin_output_labels else None,
|
||||
)
|
||||
|
||||
if explicit_blank and snapshot.provenance == "canonical" and execution.source_is_current(ctx):
|
||||
|
|
@ -3281,6 +3317,35 @@ async def _run_blocks_and_collect_debug(
|
|||
if not workflow.get_output_parameter(label):
|
||||
return {"ok": False, "error": f"Block label not found in saved workflow: {label!r}"}
|
||||
|
||||
# Resolved once: the origin-input check must compare exactly what gets dispatched.
|
||||
origin_checked_resolution: tuple[dict[str, Any], list[str], list[str]] | None = None
|
||||
origin_checked_parameter_keys: set[str] = set()
|
||||
if reused_origin_output_labels:
|
||||
origin_output_refusal = origin_definition_refusal(
|
||||
ctx, reused_origin_output_labels, labels_to_execute[0], workflow.workflow_definition
|
||||
) or origin_settings_refusal(ctx, reused_origin_output_labels, workflow)
|
||||
if origin_output_refusal is None:
|
||||
origin_checked_parameter_keys = {parameter.key for parameter in snapshot.workflow_parameters}
|
||||
origin_checked_resolution = _resolve_run_data_and_unbound_keys(
|
||||
snapshot.workflow_parameters,
|
||||
params.get("parameters") or {},
|
||||
ephemeral_input_values=(
|
||||
_ephemeral_input_values_by_parameter_key(execution.metadata, ctx.scout_trajectory)
|
||||
if use_ephemeral_inputs
|
||||
else {}
|
||||
),
|
||||
origin_parameters=ctx.repair_origin_input_values,
|
||||
origin_is_copilot_run=ctx.repair_origin_is_copilot_run,
|
||||
)
|
||||
origin_output_refusal = origin_input_refusal(
|
||||
ctx, reused_origin_output_labels, workflow.workflow_definition, origin_checked_resolution[0]
|
||||
)
|
||||
if origin_output_refusal is not None:
|
||||
block_outputs_to_seed = _stop_reusing_origin_outputs(
|
||||
execution, block_outputs_to_seed, origin_output_refusal
|
||||
)
|
||||
execution.origin_output_not_reused = origin_output_refusal
|
||||
|
||||
workflow_definition = workflow.workflow_definition
|
||||
finally_block_label = (
|
||||
workflow_definition.get("finally_block_label")
|
||||
|
|
@ -3471,12 +3536,15 @@ async def _run_blocks_and_collect_debug(
|
|||
if use_ephemeral_inputs
|
||||
else {}
|
||||
)
|
||||
data, unbound_required_parameter_keys, reused_origin_input_keys = _resolve_run_data_and_unbound_keys(
|
||||
list(snapshot.workflow_parameters),
|
||||
user_params,
|
||||
ephemeral_input_values=ephemeral_input_values,
|
||||
origin_parameters=ctx.repair_origin_input_values,
|
||||
origin_is_copilot_run=ctx.repair_origin_is_copilot_run,
|
||||
data, unbound_required_parameter_keys, reused_origin_input_keys = (
|
||||
origin_checked_resolution
|
||||
or _resolve_run_data_and_unbound_keys(
|
||||
list(snapshot.workflow_parameters),
|
||||
user_params,
|
||||
ephemeral_input_values=ephemeral_input_values,
|
||||
origin_parameters=ctx.repair_origin_input_values,
|
||||
origin_is_copilot_run=ctx.repair_origin_is_copilot_run,
|
||||
)
|
||||
)
|
||||
browser_seed: BuildTestBrowserSeed | None = None
|
||||
browser_seed_source: BrowserSeedSource | None = None
|
||||
|
|
@ -3671,6 +3739,25 @@ async def _run_blocks_and_collect_debug(
|
|||
|
||||
all_workflow_params = list(snapshot.workflow_parameters)
|
||||
all_output_params = list(snapshot.output_parameters)
|
||||
# The check above resolved every key the dispatched run reads only if persistence added none and dropped none.
|
||||
unchecked_parameter_keys = (
|
||||
sorted(origin_checked_parameter_keys ^ {parameter.key for parameter in all_workflow_params})
|
||||
if origin_checked_resolution is not None
|
||||
else []
|
||||
)
|
||||
if unchecked_parameter_keys and execution.reused_origin_output_labels:
|
||||
block_outputs_to_seed = _stop_reusing_origin_outputs(
|
||||
execution,
|
||||
block_outputs_to_seed,
|
||||
logged_origin_refusal(
|
||||
OriginOutputRefusalDetail(
|
||||
reason=OriginOutputRefusal.CHANGED_INPUT,
|
||||
block_label=execution.reused_origin_output_labels[0],
|
||||
origin_workflow_run_id=ctx.repair_origin_outputs_run_id,
|
||||
parameter_key=unchecked_parameter_keys[0],
|
||||
)
|
||||
),
|
||||
)
|
||||
|
||||
ctx.unbound_required_parameter_keys = unbound_required_parameter_keys
|
||||
execution.reused_origin_input_keys = reused_origin_input_keys
|
||||
|
|
@ -3801,6 +3888,10 @@ async def _run_blocks_and_collect_debug(
|
|||
"error": "The pending Copilot proposal changed before the test started; reload and try again.",
|
||||
}
|
||||
ctx.dispatched_run_ids_this_turn.add(workflow_run.workflow_run_id)
|
||||
# A seeded producer's value belongs to the run that produced it, never to this one.
|
||||
seeded_only_labels = frozenset(block_outputs_to_seed) - frozenset(labels_that_may_execute)
|
||||
if seeded_only_labels:
|
||||
ctx.seeded_only_labels_by_run_id[workflow_run.workflow_run_id] = seeded_only_labels
|
||||
# The browser session the run attaches overrides the run row's proxy in the browser layer,
|
||||
# so the proxy this run acts through is the session's whenever it declares one.
|
||||
session_made_hop, run_session_proxy_location = await browser_session_hop_proxy(ctx, run_session_id)
|
||||
|
|
@ -4158,6 +4249,7 @@ async def _run_blocks_and_collect_debug(
|
|||
result["data"]["user_facing_summary"] = user_facing_summary
|
||||
if execution.reused_origin_input_keys:
|
||||
result["data"]["reused_origin_input_keys"] = list(execution.reused_origin_input_keys)
|
||||
_attach_reused_origin_outputs(result["data"], execution)
|
||||
if run_cancelled_by_watchdog:
|
||||
result[_INTERNAL_RUN_CANCELLED_BY_WATCHDOG_KEY] = True
|
||||
failed_result = _newest_failed_result(result["data"]["blocks"])
|
||||
|
|
@ -4314,6 +4406,7 @@ async def _run_blocks_and_collect_debug(
|
|||
result_data["runtime_frontier_starter_url_seeded"] = True
|
||||
if execution.reused_origin_input_keys:
|
||||
result_data["reused_origin_input_keys"] = list(execution.reused_origin_input_keys)
|
||||
_attach_reused_origin_outputs(result_data, execution)
|
||||
if not run_ok and run and getattr(run, "failure_reason", None):
|
||||
result_data["failure_reason"] = redact_totp_runtime_values(run.failure_reason)
|
||||
if not run_ok and run and getattr(run, "failure_category", None):
|
||||
|
|
@ -4337,6 +4430,7 @@ async def _run_blocks_and_collect_debug(
|
|||
workflow=output_identity_workflow,
|
||||
data=result_data,
|
||||
persisted_output_parameters=all_output_params,
|
||||
excluded_block_labels=seeded_only_labels,
|
||||
)
|
||||
# The run's terminal facts are known here, so the outcome is committed before any
|
||||
# browser-dependent probe: a probe the turn deadline cancels must not erase what the run did.
|
||||
|
|
@ -4421,6 +4515,12 @@ async def _run_blocks_and_collect_debug(
|
|||
# trust blocks that individually succeeded inside it.
|
||||
if run_fully_completed and execution.source_is_current(ctx):
|
||||
for label, output in block_outputs_by_label.items():
|
||||
# Seeded as-is into the next run's `<label>_output`, so the block's own registered value is
|
||||
# kept rather than the {output_key: value} map the result packet shows.
|
||||
output_parameter = snapshot.workflow.get_output_parameter(label)
|
||||
registered = registered_outputs_by_label.get(label, {})
|
||||
if output_parameter is not None and output_parameter.key in registered:
|
||||
output = registered[output_parameter.key]
|
||||
ctx.verified_block_outputs[label] = output
|
||||
# Rebuilt from this run's rows alone: the position was forgotten at dispatch, and the
|
||||
# browser these pages describe is the one this run used.
|
||||
|
|
@ -4724,6 +4824,7 @@ async def _get_run_results(
|
|||
if run_workflow is not None
|
||||
else None
|
||||
),
|
||||
excluded_block_labels=ctx.seeded_only_labels_by_run_id.get(workflow_run_id, frozenset()),
|
||||
)
|
||||
|
||||
cold_artifact_requires_redaction_context = _workflow_requires_terminal_artifact_redaction(run_workflow)
|
||||
|
|
|
|||
|
|
@ -29,7 +29,7 @@ from skyvern.forge.sdk.copilot.workflow_block_traversal import (
|
|||
workflow_link_node_mappings,
|
||||
)
|
||||
from skyvern.forge.sdk.workflow.models.parameter import ParameterType
|
||||
from skyvern.forge.sdk.workflow.models.workflow import Workflow
|
||||
from skyvern.forge.sdk.workflow.models.workflow import Workflow, WorkflowDefinition
|
||||
from skyvern.forge.sdk.workflow.private_settings import (
|
||||
resolve_cdp_connect_headers,
|
||||
)
|
||||
|
|
@ -282,11 +282,8 @@ def _strip_runtime_block_fields(block: dict[str, Any]) -> dict[str, Any]:
|
|||
return cleaned
|
||||
|
||||
|
||||
def workflow_to_copilot_yaml(workflow: Workflow) -> str:
|
||||
workflow_data = workflow.model_dump(mode="json", exclude_none=True)
|
||||
strip_private_workflow_settings(workflow_data)
|
||||
workflow_definition = deepcopy(workflow_data.get("workflow_definition") or {})
|
||||
|
||||
def _copilot_yaml_workflow_definition(persisted_definition: dict[str, Any]) -> dict[str, Any]:
|
||||
workflow_definition = deepcopy(persisted_definition)
|
||||
parameters = workflow_definition.get("parameters")
|
||||
if isinstance(parameters, list):
|
||||
workflow_definition["parameters"] = [
|
||||
|
|
@ -300,6 +297,13 @@ def workflow_to_copilot_yaml(workflow: Workflow) -> str:
|
|||
workflow_definition["blocks"] = [
|
||||
_strip_runtime_block_fields(block) if isinstance(block, dict) else block for block in blocks
|
||||
]
|
||||
return workflow_definition
|
||||
|
||||
|
||||
def workflow_to_copilot_yaml(workflow: Workflow) -> str:
|
||||
workflow_data = workflow.model_dump(mode="json", exclude_none=True)
|
||||
strip_private_workflow_settings(workflow_data)
|
||||
workflow_definition = _copilot_yaml_workflow_definition(workflow_data.get("workflow_definition") or {})
|
||||
|
||||
request_data = {
|
||||
key: workflow_data[key]
|
||||
|
|
@ -829,6 +833,34 @@ def redact_credentials_in_workflow_yaml(
|
|||
return workflow_yaml
|
||||
|
||||
|
||||
def _copilot_definition_from_yaml(
|
||||
workflow_yaml: str, workflow_id: str
|
||||
) -> tuple[WorkflowCreateYAMLRequest, WorkflowDefinition]:
|
||||
# Single seam every copilot YAML->Workflow conversion passes through, so code
|
||||
# blocks get their plain-view steps regardless of which path produced the YAML
|
||||
# (the update_workflow tool derives them upstream; the inline REPLACE_WORKFLOW
|
||||
# fallbacks would otherwise surface "No steps yet").
|
||||
workflow_yaml = derive_code_block_steps_in_yaml(workflow_yaml)
|
||||
# Same reasoning one field over: a block whose code names a declared parameter but omits it
|
||||
# from parameter_keys gets no value at runtime and dies on NameError mid-login.
|
||||
workflow_yaml = bind_referenced_parameters_in_yaml(workflow_yaml)
|
||||
workflow_yaml_request = _normalize_copilot_yaml(workflow_yaml)
|
||||
return workflow_yaml_request, convert_workflow_definition(
|
||||
workflow_definition_yaml=workflow_yaml_request.workflow_definition,
|
||||
workflow_id=workflow_id,
|
||||
)
|
||||
|
||||
|
||||
def copilot_round_trip_definition(definition: WorkflowDefinition, *, workflow_id: str) -> WorkflowDefinition:
|
||||
"""The definition as Copilot's own save path would persist it, whoever authored it."""
|
||||
workflow_definition = _copilot_yaml_workflow_definition(definition.model_dump(mode="json", exclude_none=True))
|
||||
# The save path stitches next_block_label after validating, so a version it persisted can carry
|
||||
# version 1 beside explicit routing; letting validation derive the version again round-trips it.
|
||||
workflow_definition.pop("version", None)
|
||||
workflow_yaml = yaml.safe_dump({"title": "", "workflow_definition": workflow_definition}, sort_keys=False)
|
||||
return _copilot_definition_from_yaml(workflow_yaml, workflow_id)[1]
|
||||
|
||||
|
||||
async def _process_workflow_yaml(
|
||||
workflow_id: str,
|
||||
workflow_permanent_id: str,
|
||||
|
|
@ -839,19 +871,7 @@ async def _process_workflow_yaml(
|
|||
private_workflow_settings: dict[str, Any] | None = None,
|
||||
prefer_live_title: bool = False,
|
||||
) -> Workflow:
|
||||
# Single seam every copilot YAML->Workflow conversion passes through, so code
|
||||
# blocks get their plain-view steps regardless of which path produced the YAML
|
||||
# (the update_workflow tool derives them upstream; the inline REPLACE_WORKFLOW
|
||||
# fallbacks would otherwise surface "No steps yet").
|
||||
workflow_yaml = derive_code_block_steps_in_yaml(workflow_yaml)
|
||||
# Same reasoning one field over: a block whose code names a declared parameter but omits it
|
||||
# from parameter_keys gets no value at runtime and dies on NameError mid-login.
|
||||
workflow_yaml = bind_referenced_parameters_in_yaml(workflow_yaml)
|
||||
workflow_yaml_request = _normalize_copilot_yaml(workflow_yaml)
|
||||
updated_workflow_definition = convert_workflow_definition(
|
||||
workflow_definition_yaml=workflow_yaml_request.workflow_definition,
|
||||
workflow_id=workflow_id,
|
||||
)
|
||||
workflow_yaml_request, updated_workflow_definition = _copilot_definition_from_yaml(workflow_yaml, workflow_id)
|
||||
|
||||
enable_self_healing = workflow_yaml_request.enable_self_healing
|
||||
if enable_self_healing is None:
|
||||
|
|
|
|||
|
|
@ -6643,7 +6643,7 @@ class WorkflowService:
|
|||
organization_id=organization_id,
|
||||
browser_session_id=browser_session_id,
|
||||
block_labels=block_labels,
|
||||
block_outputs=block_outputs,
|
||||
block_output_labels=list(block_outputs or ()),
|
||||
)
|
||||
workflow_run = await self.get_workflow_run(workflow_run_id=workflow_run_id, organization_id=organization_id)
|
||||
|
||||
|
|
@ -8136,7 +8136,7 @@ class WorkflowService:
|
|||
workflow_run_id=workflow_run_id,
|
||||
block_cnt=len(blocks),
|
||||
block_labels=block_labels,
|
||||
block_outputs=block_outputs,
|
||||
block_output_labels=list(block_outputs or ()),
|
||||
)
|
||||
|
||||
else:
|
||||
|
|
|
|||
|
|
@ -57,7 +57,13 @@ from skyvern.forge.sdk.schemas.organizations import Organization
|
|||
from skyvern.forge.sdk.schemas.workflow_copilot import WorkflowCopilotChatRequest, WorkflowCopilotTitleUpdate
|
||||
from skyvern.forge.sdk.schemas.workflow_runs import WorkflowRunBlock
|
||||
from skyvern.forge.sdk.workflow.models.parameter import OutputParameter, WorkflowParameter, WorkflowParameterType
|
||||
from skyvern.forge.sdk.workflow.models.workflow import WorkflowRunParameter, WorkflowRunStatus
|
||||
from skyvern.forge.sdk.workflow.models.workflow import (
|
||||
Workflow,
|
||||
WorkflowRun,
|
||||
WorkflowRunOutputParameter,
|
||||
WorkflowRunParameter,
|
||||
WorkflowRunStatus,
|
||||
)
|
||||
from skyvern.schemas.proxy_location import ProxyLocationInput
|
||||
from skyvern.schemas.runs import ProxyLocation
|
||||
from skyvern.schemas.workflows import BlockType
|
||||
|
|
@ -328,6 +334,7 @@ def install_get_run_results_harness(
|
|||
last_run_blocks_workflow_run_id=carried_run_id,
|
||||
proposal_workflow_run_id=None,
|
||||
dispatched_run_ids_this_turn=set(),
|
||||
seeded_only_labels_by_run_id={},
|
||||
)
|
||||
|
||||
|
||||
|
|
@ -859,6 +866,126 @@ def origin_run_input(
|
|||
)
|
||||
|
||||
|
||||
INERT_APPROVAL_WORKFLOW_YAML = """
|
||||
title: inert approval
|
||||
workflow_definition:
|
||||
parameters:
|
||||
- parameter_type: workflow
|
||||
workflow_parameter_type: string
|
||||
key: request_id
|
||||
blocks:
|
||||
- block_type: extraction
|
||||
label: approval
|
||||
url: https://example.test/approval
|
||||
data_extraction_goal: "Extract whether request {{ request_id }} is authorized."
|
||||
parameter_keys:
|
||||
- request_id
|
||||
- block_type: extraction
|
||||
label: source_status
|
||||
data_extraction_goal: "Report the source status for authorization {{ approval.output.authorized }}."
|
||||
"""
|
||||
|
||||
REPAIRED_APPROVAL_WORKFLOW_YAML = INERT_APPROVAL_WORKFLOW_YAML.replace(
|
||||
"Report the source status for authorization {{ approval.output.authorized }}.",
|
||||
"Report the repaired source status for authorization {{ approval.output.authorized }} "
|
||||
"and {{ approval_output.extracted_information.authorized }}.",
|
||||
)
|
||||
|
||||
ORIGIN_RUN_ID = "wr_origin"
|
||||
|
||||
|
||||
async def inert_approval_workflow(workflow_yaml: str, *, workflow_id: str) -> Workflow:
|
||||
"""A converted version mints its own parameter ids, as each persisted version does."""
|
||||
return await process_workflow_yaml(
|
||||
settings_fallback_yaml="enable_self_healing: false",
|
||||
workflow_id=workflow_id,
|
||||
workflow_permanent_id="wfp-1",
|
||||
organization_id="org-1",
|
||||
workflow_yaml=workflow_yaml,
|
||||
)
|
||||
|
||||
|
||||
ORIGIN_OUTPUT_SENTINEL = "origin-output-sentinel-7f3a"
|
||||
|
||||
|
||||
def origin_block_rows(
|
||||
workflow: Workflow,
|
||||
label: str,
|
||||
*,
|
||||
status: str = "completed",
|
||||
registered: bool = True,
|
||||
value: dict | list | str | None = None,
|
||||
row_output: dict | list | str | None = None,
|
||||
minute: int = 0,
|
||||
parent: str | None = None,
|
||||
) -> tuple[list[WorkflowRunBlock], list[WorkflowRunOutputParameter]]:
|
||||
now = datetime(2026, 9, 1, 12, minute, tzinfo=UTC)
|
||||
run_block = WorkflowRunBlock(
|
||||
workflow_run_block_id=f"wrb_origin_{label}_{minute}",
|
||||
workflow_run_id=ORIGIN_RUN_ID,
|
||||
organization_id="org-1",
|
||||
parent_workflow_run_block_id=parent,
|
||||
block_type="extraction",
|
||||
label=label,
|
||||
status=status,
|
||||
output=row_output,
|
||||
created_at=now,
|
||||
modified_at=now,
|
||||
)
|
||||
output_parameter = workflow.get_output_parameter(label)
|
||||
if not registered or output_parameter is None:
|
||||
return [run_block], []
|
||||
return [run_block], [
|
||||
WorkflowRunOutputParameter(
|
||||
workflow_run_id=ORIGIN_RUN_ID,
|
||||
output_parameter_id=output_parameter.output_parameter_id,
|
||||
value=value,
|
||||
created_at=now,
|
||||
)
|
||||
]
|
||||
|
||||
|
||||
def merge_origin_rows(
|
||||
*rows: tuple[list[WorkflowRunBlock], list[WorkflowRunOutputParameter]],
|
||||
) -> tuple[list[WorkflowRunBlock], list[WorkflowRunOutputParameter]]:
|
||||
return [block for blocks, _ in rows for block in blocks], [row for _, registered in rows for row in registered]
|
||||
|
||||
|
||||
ORIGIN_RUN_CREATED_AT = datetime.now(UTC) + timedelta(days=1)
|
||||
|
||||
|
||||
def origin_run_row(**overrides: object) -> WorkflowRun:
|
||||
fields: dict[str, object] = {
|
||||
"workflow_run_id": ORIGIN_RUN_ID,
|
||||
"workflow_id": "w_origin",
|
||||
"organization_id": "org-1",
|
||||
"workflow_permanent_id": "wfp-1",
|
||||
"browser_session_id": "pbs_origin",
|
||||
"status": WorkflowRunStatus.completed,
|
||||
# Later than any workflow version a test builds, as a real run starts after its version was saved.
|
||||
"created_at": ORIGIN_RUN_CREATED_AT,
|
||||
"modified_at": ORIGIN_RUN_CREATED_AT,
|
||||
}
|
||||
return WorkflowRun.model_validate({**fields, **overrides})
|
||||
|
||||
|
||||
def install_origin_run(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
*,
|
||||
origin_workflow: Workflow,
|
||||
rows: tuple[list[WorkflowRunBlock], list[WorkflowRunOutputParameter]],
|
||||
**run_overrides: object,
|
||||
) -> None:
|
||||
run = origin_run_row(**{"workflow_id": origin_workflow.workflow_id, **run_overrides})
|
||||
monkeypatch.setattr(forge_app.WORKFLOW_SERVICE, "get_workflow_run", AsyncMock(return_value=run))
|
||||
monkeypatch.setattr(forge_app.DATABASE.workflow_runs, "get_workflow_run_parameters", AsyncMock(return_value=[]))
|
||||
monkeypatch.setattr(forge_app.DATABASE.workflows, "get_workflow", AsyncMock(return_value=origin_workflow))
|
||||
monkeypatch.setattr(forge_app.DATABASE.observer, "get_workflow_run_blocks", AsyncMock(return_value=rows[0]))
|
||||
monkeypatch.setattr(
|
||||
forge_app.DATABASE.workflow_runs, "get_workflow_run_output_parameters", AsyncMock(return_value=rows[1])
|
||||
)
|
||||
|
||||
|
||||
def make_copilot_ctx(**overrides: object) -> CopilotContext:
|
||||
defaults: dict[str, object] = dict(
|
||||
organization_id="org-1",
|
||||
|
|
|
|||
|
|
@ -5,6 +5,8 @@ from __future__ import annotations
|
|||
import asyncio
|
||||
import copy
|
||||
import json
|
||||
from collections.abc import Callable
|
||||
from dataclasses import replace
|
||||
from itertools import pairwise
|
||||
from pathlib import Path
|
||||
from typing import Any, cast
|
||||
|
|
@ -12,6 +14,7 @@ from unittest.mock import AsyncMock
|
|||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
import yaml
|
||||
from agents.items import ModelResponse
|
||||
from agents.models.interface import Model
|
||||
from agents.run_config import RunConfig
|
||||
|
|
@ -19,6 +22,8 @@ from agents.usage import Usage
|
|||
from jinja2.sandbox import SandboxedEnvironment
|
||||
from openai.types.responses import Response, ResponseCompletedEvent, ResponseOutputMessage, ResponseOutputText
|
||||
|
||||
from skyvern.constants import SCRUBBED_VALUE
|
||||
from skyvern.forge import app
|
||||
from skyvern.forge.sdk.copilot import agent as agent_module
|
||||
from skyvern.forge.sdk.copilot import tools
|
||||
from skyvern.forge.sdk.copilot.agent import _verified_workflow_or_none
|
||||
|
|
@ -33,6 +38,11 @@ from skyvern.forge.sdk.copilot.output_utils import (
|
|||
sanitize_tool_result_for_llm,
|
||||
summarize_tool_result,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.repair_origin_run import (
|
||||
OriginOutputRefusal,
|
||||
OriginOutputSnapshot,
|
||||
seed_repair_origin_run,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.request_policy import RequestPolicy
|
||||
from skyvern.forge.sdk.copilot.review_gate import workflow_block_fingerprints
|
||||
from skyvern.forge.sdk.copilot.run_outcome import RecordedRunOutcome
|
||||
|
|
@ -58,7 +68,27 @@ from skyvern.forge.sdk.copilot.tools.run_execution import (
|
|||
terminal_ready_for_latch,
|
||||
)
|
||||
from skyvern.forge.sdk.schemas.workflow_copilot import WorkflowCopilotChatRequest
|
||||
from skyvern.forge.sdk.schemas.workflow_runs import WorkflowRunBlock
|
||||
from skyvern.forge.sdk.workflow.models.parameter import RESERVED_PARAMETER_KEYS
|
||||
from skyvern.forge.sdk.workflow.models.workflow import (
|
||||
Workflow,
|
||||
WorkflowDefinition,
|
||||
WorkflowRunOutputParameter,
|
||||
WorkflowRunStatus,
|
||||
)
|
||||
from skyvern.forge.sdk.workflow.workflow_definition_converter import convert_workflow_definition
|
||||
from skyvern.schemas.workflows import WorkflowCreateYAMLRequest
|
||||
from tests.unit.copilot_test_helpers import (
|
||||
INERT_APPROVAL_WORKFLOW_YAML,
|
||||
ORIGIN_OUTPUT_SENTINEL,
|
||||
ORIGIN_RUN_ID,
|
||||
REPAIRED_APPROVAL_WORKFLOW_YAML,
|
||||
inert_approval_workflow,
|
||||
install_origin_run,
|
||||
make_copilot_ctx,
|
||||
merge_origin_rows,
|
||||
origin_block_rows,
|
||||
)
|
||||
|
||||
|
||||
class _FakeBlock:
|
||||
|
|
@ -4880,3 +4910,617 @@ async def test_test_end_to_end_will_not_touch_the_browser_after_a_raw_secret(
|
|||
result = await run_workflow_end_to_end(ctx, "workflow: yaml")
|
||||
|
||||
assert result["ok"] is False
|
||||
|
||||
|
||||
_APPROVAL_VALUE = {"extracted_information": {"authorized": True, "note": ORIGIN_OUTPUT_SENTINEL}}
|
||||
|
||||
_SOURCE_STATUS_FIRST_YAML = """
|
||||
title: inert approval
|
||||
workflow_definition:
|
||||
parameters:
|
||||
- parameter_type: workflow
|
||||
workflow_parameter_type: string
|
||||
key: request_id
|
||||
blocks:
|
||||
- block_type: extraction
|
||||
label: source_status
|
||||
data_extraction_goal: "Report the repaired source status for authorization {{ approval.output.authorized }}."
|
||||
- block_type: extraction
|
||||
label: approval
|
||||
url: https://example.test/approval
|
||||
data_extraction_goal: "Extract whether request {{ request_id }} is authorized."
|
||||
parameter_keys:
|
||||
- request_id
|
||||
"""
|
||||
|
||||
_SIGN_IN_SOURCE_STATUS_YAML = REPAIRED_APPROVAL_WORKFLOW_YAML.replace(
|
||||
" - block_type: extraction\n label: source_status\n",
|
||||
" - block_type: login\n label: source_status\n url: https://example.test/sign-in\n"
|
||||
" navigation_goal: Sign in.\n",
|
||||
).replace('data_extraction_goal: "Report the repaired', 'complete_criterion: "Report the repaired')
|
||||
|
||||
|
||||
OriginRows = tuple[list[WorkflowRunBlock], list[WorkflowRunOutputParameter]]
|
||||
|
||||
|
||||
async def _origin_turn(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
*,
|
||||
origin_yaml: str = INERT_APPROVAL_WORKFLOW_YAML,
|
||||
rows: str | Callable[[Workflow], OriginRows] = "completed",
|
||||
requested: str | None = ORIGIN_RUN_ID,
|
||||
origin: Workflow | None = None,
|
||||
**run_overrides: object,
|
||||
) -> CopilotContext:
|
||||
monkeypatch.setattr(app.WORKFLOW_SERVICE, "get_workflow_by_permanent_id", AsyncMock(return_value=None))
|
||||
origin = origin or await inert_approval_workflow(origin_yaml, workflow_id="w_origin")
|
||||
origin_rows = (
|
||||
rows(origin)
|
||||
if callable(rows)
|
||||
else {
|
||||
"completed": origin_block_rows(origin, "approval", value=_APPROVAL_VALUE),
|
||||
"explicit_null": origin_block_rows(origin, "approval", value=None),
|
||||
"absent": ([], []),
|
||||
"output_row_only": ([], origin_block_rows(origin, "approval", value=_APPROVAL_VALUE)[1]),
|
||||
"failed": origin_block_rows(origin, "approval", status="failed", value=_APPROVAL_VALUE),
|
||||
"unregistered": origin_block_rows(origin, "approval", registered=False),
|
||||
}[rows]
|
||||
)
|
||||
install_origin_run(monkeypatch, origin_workflow=origin, rows=origin_rows, **run_overrides)
|
||||
ctx = make_copilot_ctx()
|
||||
await seed_repair_origin_run(ctx, workflow_run_id=requested)
|
||||
return ctx
|
||||
|
||||
|
||||
async def _definitions(
|
||||
candidate_yaml: str = REPAIRED_APPROVAL_WORKFLOW_YAML,
|
||||
) -> tuple[WorkflowDefinition, WorkflowDefinition]:
|
||||
old = await inert_approval_workflow(INERT_APPROVAL_WORKFLOW_YAML, workflow_id="w_origin")
|
||||
new = await inert_approval_workflow(candidate_yaml, workflow_id="w_candidate")
|
||||
return old.workflow_definition, new.workflow_definition
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_edited_downstream_block_is_seeded_with_the_origin_output_and_nothing_else(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch)
|
||||
old, new = await _definitions()
|
||||
|
||||
labels, seed, start, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert (labels, start) == (["source_status"], "source_status")
|
||||
assert seed == {"approval": _APPROVAL_VALUE}
|
||||
assert ctx.frontier_origin_reused_labels == ["approval"]
|
||||
assert ctx.frontier_origin_output_refusal is None
|
||||
assert ctx.verified_block_outputs == {}
|
||||
assert ctx.verified_prefix_labels == []
|
||||
assert ctx.composition_verified_labels == []
|
||||
assert ctx.verified_prefix_block_end_urls == {}
|
||||
assert ctx.frontier_resume_session_id is None
|
||||
assert ORIGIN_OUTPUT_SENTINEL not in repr(ctx)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("candidate_yaml", "own_browser"),
|
||||
[(INERT_APPROVAL_WORKFLOW_YAML, False), (_SIGN_IN_SOURCE_STATUS_YAML, True)],
|
||||
ids=["run_blocks_unchanged_definition", "planner_gives_the_start_its_own_browser"],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_every_planner_branch_that_drops_the_seed_is_refilled_from_the_origin(
|
||||
monkeypatch: pytest.MonkeyPatch, candidate_yaml: str, own_browser: bool
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch)
|
||||
old, new = await _definitions(candidate_yaml)
|
||||
if candidate_yaml == INERT_APPROVAL_WORKFLOW_YAML:
|
||||
old = new
|
||||
|
||||
labels, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert frontier_module._plan_frontier_base(make_copilot_ctx(), ["source_status"], old, new)[1] == {}
|
||||
assert ctx.frontier_requires_own_browser is own_browser
|
||||
assert labels == ["source_status"]
|
||||
assert seed == {"approval": _APPROVAL_VALUE}
|
||||
assert ctx.frontier_origin_reused_labels == ["approval"]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize("rows", ["explicit_null", "unregistered"])
|
||||
async def test_a_stored_null_or_missing_output_is_never_reused(monkeypatch: pytest.MonkeyPatch, rows: str) -> None:
|
||||
ctx = await _origin_turn(monkeypatch, rows=rows)
|
||||
old, new = await _definitions()
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert "approval" not in seed
|
||||
assert ctx.frontier_origin_output_refusal is not None
|
||||
assert ctx.frontier_origin_output_refusal.reason is OriginOutputRefusal.OUTPUT_UNAVAILABLE
|
||||
|
||||
|
||||
def _with_workflow_prompt(workflow_yaml: str, prompt: str) -> str:
|
||||
return workflow_yaml.replace(
|
||||
"workflow_definition:\n", f'workflow_definition:\n workflow_system_prompt: "{prompt}"\n', 1
|
||||
)
|
||||
|
||||
|
||||
def _with_approval_export(workflow_yaml: str, schema_type: str) -> str:
|
||||
goal = ' data_extraction_goal: "Extract whether request {{ request_id }} is authorized."\n'
|
||||
return workflow_yaml.replace(
|
||||
goal, f"{goal} export_enabled: true\n export_data_schema:\n type: {schema_type}\n"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_producer_that_cannot_be_normalized_is_unavailable_not_changed(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch)
|
||||
old, new = await _definitions()
|
||||
|
||||
def _unparseable(*_args: object, **_kwargs: object) -> WorkflowDefinition:
|
||||
raise ValueError("unparseable")
|
||||
|
||||
monkeypatch.setattr(frontier_module, "copilot_round_trip_definition", _unparseable)
|
||||
_plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert ctx.frontier_origin_output_refusal is not None
|
||||
assert ctx.frontier_origin_output_refusal.reason is OriginOutputRefusal.OUTPUT_UNAVAILABLE
|
||||
|
||||
|
||||
_EXTRACTION_APPROVAL_BLOCK = """ - block_type: extraction
|
||||
label: approval
|
||||
url: https://example.test/approval
|
||||
data_extraction_goal: "Extract whether request {{ request_id }} is authorized."
|
||||
parameter_keys:
|
||||
- request_id
|
||||
"""
|
||||
_UNBOUND_CODE_APPROVAL_BLOCK = """ - block_type: code
|
||||
label: approval
|
||||
code: "result = {'authorized': bool(request_id)}"
|
||||
"""
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"approval_block",
|
||||
[_EXTRACTION_APPROVAL_BLOCK, _UNBOUND_CODE_APPROVAL_BLOCK],
|
||||
ids=["unrouted_extraction", "code_without_parameter_keys"],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_origin_version_saved_outside_copilot_is_compared_by_what_copilot_would_save(
|
||||
monkeypatch: pytest.MonkeyPatch, approval_block: str
|
||||
) -> None:
|
||||
monkeypatch.setattr(app.WORKFLOW_SERVICE, "get_workflow_by_permanent_id", AsyncMock(return_value=None))
|
||||
origin_yaml = INERT_APPROVAL_WORKFLOW_YAML.replace(_EXTRACTION_APPROVAL_BLOCK, approval_block)
|
||||
copilot_saved = await inert_approval_workflow(origin_yaml, workflow_id="w_origin")
|
||||
raw = yaml.safe_load(origin_yaml)
|
||||
for block in raw["workflow_definition"]["blocks"]:
|
||||
block["title"] = ""
|
||||
api_definition = convert_workflow_definition(
|
||||
workflow_definition_yaml=WorkflowCreateYAMLRequest.model_validate(raw).workflow_definition,
|
||||
workflow_id="w_origin",
|
||||
)
|
||||
assert api_definition.blocks[0] != copilot_saved.workflow_definition.blocks[0]
|
||||
ctx = await _origin_turn(
|
||||
monkeypatch, origin=copilot_saved.model_copy(update={"workflow_definition": api_definition})
|
||||
)
|
||||
old = copilot_saved.workflow_definition
|
||||
new = (
|
||||
await inert_approval_workflow(
|
||||
REPAIRED_APPROVAL_WORKFLOW_YAML.replace(_EXTRACTION_APPROVAL_BLOCK, approval_block),
|
||||
workflow_id="w_candidate",
|
||||
)
|
||||
).workflow_definition
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert seed == {"approval": _APPROVAL_VALUE}
|
||||
assert ctx.frontier_origin_output_refusal is None
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("candidate_yaml", "requested"),
|
||||
[
|
||||
(
|
||||
REPAIRED_APPROVAL_WORKFLOW_YAML.replace(
|
||||
" - block_type: extraction\n label: source_status\n",
|
||||
" - block_type: extraction\n label: review\n data_extraction_goal: Review the request.\n"
|
||||
" - block_type: extraction\n label: source_status\n",
|
||||
),
|
||||
"source_status",
|
||||
),
|
||||
(REPAIRED_APPROVAL_WORKFLOW_YAML.replace("label: source_status", "label: source_check"), "source_check"),
|
||||
(
|
||||
REPAIRED_APPROVAL_WORKFLOW_YAML.replace(
|
||||
" label: approval\n", " label: approval\n title: Approval check\n"
|
||||
),
|
||||
"source_status",
|
||||
),
|
||||
],
|
||||
ids=["block_inserted_after_producer", "producer_successor_renamed", "producer_display_title_edited"],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_unchanged_producer_is_reused_when_only_its_successor_changes(
|
||||
monkeypatch: pytest.MonkeyPatch, candidate_yaml: str, requested: str
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch)
|
||||
old, new = await _definitions(candidate_yaml)
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, [requested], old, new)
|
||||
|
||||
assert ctx.frontier_origin_output_refusal is None
|
||||
assert seed == {"approval": _APPROVAL_VALUE}
|
||||
assert ctx.frontier_origin_reused_labels == ["approval"]
|
||||
|
||||
|
||||
def _with_intake_before_approval(workflow_yaml: str, intake_goal: str = "Read the intake for {{ region }}.") -> str:
|
||||
return workflow_yaml.replace(
|
||||
" key: request_id\n blocks:\n",
|
||||
" key: request_id\n - parameter_type: workflow\n workflow_parameter_type: string\n"
|
||||
" key: region\n blocks:\n - block_type: extraction\n label: intake\n"
|
||||
f' url: https://example.test/intake\n data_extraction_goal: "{intake_goal}"\n'
|
||||
" parameter_keys:\n - region\n",
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_changed_block_before_the_producer_makes_its_origin_output_stale(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch, origin_yaml=_with_intake_before_approval(INERT_APPROVAL_WORKFLOW_YAML))
|
||||
old, new = await _definitions(
|
||||
_with_intake_before_approval(REPAIRED_APPROVAL_WORKFLOW_YAML, "Read the archived intake for {{ region }}.")
|
||||
)
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert "approval" not in seed
|
||||
refusal = ctx.frontier_origin_output_refusal
|
||||
assert refusal is not None
|
||||
assert refusal.as_payload() == {
|
||||
"reason": "changed_producer",
|
||||
"block_label": "approval",
|
||||
"output_key": "approval_output",
|
||||
"origin_workflow_run_id": ORIGIN_RUN_ID,
|
||||
"changed_label": "intake",
|
||||
}
|
||||
|
||||
|
||||
def _with_intake_after_approval(workflow_yaml: str) -> str:
|
||||
moved = _with_intake_before_approval(workflow_yaml)
|
||||
intake = moved[
|
||||
moved.index(" - block_type: extraction\n label: intake\n") : moved.index(
|
||||
" - block_type: extraction\n label: approval\n"
|
||||
)
|
||||
]
|
||||
moved = moved.replace(intake, "", 1)
|
||||
return moved.replace(
|
||||
" - block_type: extraction\n label: source_status\n",
|
||||
intake + " - block_type: extraction\n label: source_status\n",
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"candidate_yaml",
|
||||
[REPAIRED_APPROVAL_WORKFLOW_YAML, _with_intake_after_approval(REPAIRED_APPROVAL_WORKFLOW_YAML)],
|
||||
ids=["block_before_producer_removed", "block_before_producer_moved_after_it"],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_producer_whose_predecessors_differ_from_the_origins_is_refused(
|
||||
monkeypatch: pytest.MonkeyPatch, candidate_yaml: str
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch, origin_yaml=_with_intake_before_approval(INERT_APPROVAL_WORKFLOW_YAML))
|
||||
old, new = await _definitions(candidate_yaml)
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert "approval" not in seed
|
||||
assert ctx.frontier_origin_output_refusal is not None
|
||||
assert ctx.frontier_origin_output_refusal.reason is OriginOutputRefusal.CHANGED_PRODUCER
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_the_dispatch_recheck_reports_the_reason_that_failed(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
ctx = await _origin_turn(monkeypatch, rows="unregistered")
|
||||
_, new = await _definitions()
|
||||
|
||||
refusal = frontier_module.origin_definition_refusal(ctx, ["approval"], "source_status", new)
|
||||
|
||||
assert refusal is not None
|
||||
assert refusal.reason is OriginOutputRefusal.OUTPUT_UNAVAILABLE
|
||||
|
||||
|
||||
@pytest.mark.parametrize(("tone", "refused"), [(None, False), ("formal", True)])
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_input_the_workflow_prompt_reads_must_match_the_origin(
|
||||
monkeypatch: pytest.MonkeyPatch, tone: str | None, refused: bool
|
||||
) -> None:
|
||||
def with_tone(workflow_yaml: str) -> str:
|
||||
return _with_workflow_prompt(workflow_yaml, "Answer in a {{ tone }} tone.").replace(
|
||||
" key: request_id\n",
|
||||
" key: request_id\n - parameter_type: workflow\n workflow_parameter_type: string\n"
|
||||
" key: tone\n",
|
||||
)
|
||||
|
||||
ctx = await _origin_turn(monkeypatch, origin_yaml=with_tone(INERT_APPROVAL_WORKFLOW_YAML))
|
||||
_, new = await _definitions(with_tone(REPAIRED_APPROVAL_WORKFLOW_YAML))
|
||||
|
||||
refusal = frontier_module.origin_input_refusal(ctx, ["approval"], new, {"tone": tone})
|
||||
|
||||
assert (refusal.parameter_key if refusal is not None else None) == ("tone" if refused else None)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(("region", "refused"), [(None, False), ("west", True)])
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_input_read_only_by_a_block_before_the_producer_must_match_the_origin(
|
||||
monkeypatch: pytest.MonkeyPatch, region: str | None, refused: bool
|
||||
) -> None:
|
||||
workflow_yaml = _with_intake_before_approval(REPAIRED_APPROVAL_WORKFLOW_YAML)
|
||||
ctx = await _origin_turn(monkeypatch, origin_yaml=_with_intake_before_approval(INERT_APPROVAL_WORKFLOW_YAML))
|
||||
_, new = await _definitions(workflow_yaml)
|
||||
|
||||
refusal = frontier_module.origin_input_refusal(ctx, ["approval"], new, {"region": region})
|
||||
|
||||
assert (refusal is not None) is refused
|
||||
if refusal is not None:
|
||||
assert (refusal.reason, refusal.block_label, refusal.parameter_key) == (
|
||||
OriginOutputRefusal.CHANGED_INPUT,
|
||||
"approval",
|
||||
"region",
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_same_turn_verified_output_wins_over_the_origin_output(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
ctx = await _origin_turn(monkeypatch)
|
||||
verified = {"extracted_information": {"authorized": False}}
|
||||
ctx.verified_block_outputs = {"approval": verified}
|
||||
old, new = await _definitions()
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert seed == {"approval": verified}
|
||||
assert ctx.frontier_origin_reused_labels == []
|
||||
assert ctx.verified_block_outputs == {"approval": verified}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_producer_the_request_names_runs_instead_of_being_seeded(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
ctx = await _origin_turn(monkeypatch, rows="failed")
|
||||
old, new = await _definitions()
|
||||
|
||||
labels, seed, _, _ = _plan_frontier(ctx, ["approval", "source_status"], old, new)
|
||||
|
||||
assert labels == ["approval", "source_status"]
|
||||
assert seed == {}
|
||||
assert ctx.frontier_origin_output_refusal is None
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"turn",
|
||||
[{"requested": None}, {"debug_session_id": "ds_debugger"}],
|
||||
ids=["no_run_named", "debugger_block_run_is_never_an_origin"],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_turn_opened_about_no_run_plans_exactly_as_before(
|
||||
monkeypatch: pytest.MonkeyPatch, turn: dict[str, str | None]
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch, **turn)
|
||||
old, new = await _definitions()
|
||||
|
||||
plan = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert plan[:3] == (["source_status"], {}, "source_status")
|
||||
assert ctx.frontier_origin_output_refusal is None
|
||||
assert ctx.frontier_origin_reused_labels == []
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_resumed_plan_is_never_given_origin_outputs(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
ctx = await _origin_turn(monkeypatch)
|
||||
_, new = await _definitions()
|
||||
ctx.frontier_resume_session_id = "pbs_prefix"
|
||||
plan: frontier_module.FrontierPlan = (["source_status"], {}, "source_status", "resumed")
|
||||
|
||||
assert frontier_module._fill_seed_from_origin(ctx, plan, new) == plan
|
||||
assert ctx.frontier_origin_reused_labels == []
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("rows", "run_overrides", "origin_yaml", "candidate_yaml", "expected"),
|
||||
[
|
||||
("completed", {"organization_id": "org-other"}, None, None, OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN),
|
||||
(
|
||||
"completed",
|
||||
{"workflow_permanent_id": "wfp-other"},
|
||||
None,
|
||||
None,
|
||||
OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN,
|
||||
),
|
||||
("completed", {"status": WorkflowRunStatus.running}, None, None, OriginOutputRefusal.ORIGIN_UNSETTLED),
|
||||
("absent", {}, None, None, OriginOutputRefusal.UPSTREAM_ABSENT),
|
||||
("output_row_only", {}, None, None, OriginOutputRefusal.UPSTREAM_ABSENT),
|
||||
("failed", {}, None, None, OriginOutputRefusal.UPSTREAM_FAILED),
|
||||
(
|
||||
"completed",
|
||||
{},
|
||||
INERT_APPROVAL_WORKFLOW_YAML.replace("is authorized", "was approved"),
|
||||
None,
|
||||
OriginOutputRefusal.CHANGED_PRODUCER,
|
||||
),
|
||||
("unregistered", {}, None, None, OriginOutputRefusal.OUTPUT_UNAVAILABLE),
|
||||
("completed", {}, None, _SOURCE_STATUS_FIRST_YAML, OriginOutputRefusal.ORDER_UNPROVABLE),
|
||||
(
|
||||
"completed",
|
||||
{},
|
||||
_with_approval_export(INERT_APPROVAL_WORKFLOW_YAML, "object"),
|
||||
_with_approval_export(REPAIRED_APPROVAL_WORKFLOW_YAML, "array"),
|
||||
OriginOutputRefusal.CHANGED_PRODUCER,
|
||||
),
|
||||
(
|
||||
"completed",
|
||||
{},
|
||||
_with_workflow_prompt(INERT_APPROVAL_WORKFLOW_YAML, "Answer briefly."),
|
||||
_with_workflow_prompt(REPAIRED_APPROVAL_WORKFLOW_YAML, "Answer in full sentences."),
|
||||
OriginOutputRefusal.CHANGED_PRODUCER,
|
||||
),
|
||||
],
|
||||
ids=[
|
||||
"foreign_organization",
|
||||
"workflow_mismatch",
|
||||
"origin_still_running",
|
||||
"upstream_absent",
|
||||
"output_row_without_a_block_row",
|
||||
"upstream_failed",
|
||||
"changed_producer",
|
||||
"output_unavailable",
|
||||
"order_unprovable",
|
||||
"export_schema_edited_while_export_is_on",
|
||||
"workflow_system_prompt_edited",
|
||||
],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_unusable_origin_output_is_refused_by_name_without_its_value(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
rows: str,
|
||||
run_overrides: dict[str, object],
|
||||
origin_yaml: str | None,
|
||||
candidate_yaml: str | None,
|
||||
expected: OriginOutputRefusal,
|
||||
) -> None:
|
||||
ctx = await _origin_turn(
|
||||
monkeypatch, origin_yaml=origin_yaml or INERT_APPROVAL_WORKFLOW_YAML, rows=rows, **run_overrides
|
||||
)
|
||||
old, new = await _definitions(candidate_yaml or REPAIRED_APPROVAL_WORKFLOW_YAML)
|
||||
|
||||
labels, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
refusal = ctx.frontier_origin_output_refusal
|
||||
assert refusal is not None
|
||||
# A run that failed the ownership checks is never named back to the model.
|
||||
owned = expected is not OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN
|
||||
assert refusal.as_payload() == {
|
||||
"reason": expected.value,
|
||||
"block_label": "approval",
|
||||
"output_key": "approval_output",
|
||||
**({"origin_workflow_run_id": ORIGIN_RUN_ID} if owned else {}),
|
||||
}
|
||||
assert labels == ["source_status"]
|
||||
assert "approval" not in seed
|
||||
assert ctx.frontier_origin_reused_labels == []
|
||||
assert ORIGIN_OUTPUT_SENTINEL not in repr(ctx)
|
||||
|
||||
|
||||
_INTAKE_ORIGIN_YAML = _with_intake_before_approval(INERT_APPROVAL_WORKFLOW_YAML)
|
||||
_INTAKE_CANDIDATE_YAML = _with_intake_before_approval(REPAIRED_APPROVAL_WORKFLOW_YAML)
|
||||
_REVIEW_BLOCK = " - block_type: extraction\n label: review\n data_extraction_goal: Review the intake.\n"
|
||||
_APPROVAL_BLOCK_START = " - block_type: extraction\n label: approval\n"
|
||||
|
||||
|
||||
def _with_review_before_approval(workflow_yaml: str, *, intake_jumps_to_approval: bool = False) -> str:
|
||||
reviewed = workflow_yaml.replace(_APPROVAL_BLOCK_START, _REVIEW_BLOCK + _APPROVAL_BLOCK_START)
|
||||
if intake_jumps_to_approval:
|
||||
reviewed = reviewed.replace(" label: intake\n", " label: intake\n next_block_label: approval\n")
|
||||
return reviewed
|
||||
|
||||
|
||||
def _approval_after_intake(*intake_rows: dict[str, Any], approval_minute: int = 2) -> Callable[[Workflow], OriginRows]:
|
||||
def rows(origin: Workflow) -> OriginRows:
|
||||
return merge_origin_rows(
|
||||
*(origin_block_rows(origin, "intake", **row) for row in intake_rows),
|
||||
origin_block_rows(origin, "approval", value=_APPROVAL_VALUE, minute=approval_minute),
|
||||
)
|
||||
|
||||
return rows
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("origin_yaml", "candidate_yaml", "rows", "reason", "changed_label"),
|
||||
[
|
||||
(_INTAKE_ORIGIN_YAML, _INTAKE_CANDIDATE_YAML, _approval_after_intake(), "upstream_absent", "intake"),
|
||||
(
|
||||
_INTAKE_ORIGIN_YAML,
|
||||
_INTAKE_CANDIDATE_YAML,
|
||||
_approval_after_intake({"status": "failed", "minute": 1, "registered": False}),
|
||||
"upstream_failed",
|
||||
"intake",
|
||||
),
|
||||
(
|
||||
_INTAKE_ORIGIN_YAML,
|
||||
_INTAKE_CANDIDATE_YAML,
|
||||
_approval_after_intake({"minute": 5}),
|
||||
"order_unprovable",
|
||||
"intake",
|
||||
),
|
||||
(
|
||||
_with_review_before_approval(_INTAKE_ORIGIN_YAML, intake_jumps_to_approval=True),
|
||||
_with_review_before_approval(_INTAKE_CANDIDATE_YAML),
|
||||
_approval_after_intake({"minute": 1}),
|
||||
"order_unprovable",
|
||||
None,
|
||||
),
|
||||
],
|
||||
ids=[
|
||||
"partial_origin_with_only_the_producer_row",
|
||||
"block_before_the_producer_failed",
|
||||
"block_before_the_producer_ran_after_it",
|
||||
"origin_route_jumped_over_a_block_before_the_producer",
|
||||
],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_origin_that_did_not_run_the_producers_whole_prefix_first_is_refused(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
origin_yaml: str,
|
||||
candidate_yaml: str,
|
||||
rows: Callable[[Workflow], OriginRows],
|
||||
reason: str,
|
||||
changed_label: str | None,
|
||||
) -> None:
|
||||
ctx = await _origin_turn(monkeypatch, origin_yaml=origin_yaml, rows=rows)
|
||||
old, new = await _definitions(candidate_yaml)
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert "approval" not in seed
|
||||
refusal = ctx.frontier_origin_output_refusal
|
||||
assert refusal is not None
|
||||
assert refusal.as_payload() == {
|
||||
"reason": reason,
|
||||
"block_label": "approval",
|
||||
"output_key": "approval_output",
|
||||
"origin_workflow_run_id": ORIGIN_RUN_ID,
|
||||
**({"changed_label": changed_label} if changed_label else {}),
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_retried_producer_is_reused_from_its_latest_row_after_its_prefix(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
def rows(origin: Workflow) -> OriginRows:
|
||||
return merge_origin_rows(
|
||||
origin_block_rows(origin, "intake", minute=0),
|
||||
origin_block_rows(origin, "approval", status="failed", minute=1, registered=False),
|
||||
origin_block_rows(origin, "intake", minute=2),
|
||||
origin_block_rows(origin, "approval", value=_APPROVAL_VALUE, minute=3),
|
||||
)
|
||||
|
||||
ctx = await _origin_turn(monkeypatch, origin_yaml=_INTAKE_ORIGIN_YAML, rows=rows)
|
||||
old, new = await _definitions(_INTAKE_CANDIDATE_YAML)
|
||||
|
||||
_, seed, _, _ = _plan_frontier(ctx, ["source_status"], old, new)
|
||||
|
||||
assert ctx.frontier_origin_output_refusal is None
|
||||
assert seed == {"approval": _APPROVAL_VALUE}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_scrubbed_origin_input_never_proves_the_test_input_equal(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
ctx = await _origin_turn(monkeypatch)
|
||||
assert isinstance(ctx.repair_origin_outputs, OriginOutputSnapshot)
|
||||
ctx.repair_origin_outputs = replace(ctx.repair_origin_outputs, input_values={"request_id": SCRUBBED_VALUE})
|
||||
_, new = await _definitions()
|
||||
|
||||
refusal = frontier_module.origin_input_refusal(ctx, ["approval"], new, {"request_id": SCRUBBED_VALUE})
|
||||
|
||||
assert refusal is not None
|
||||
assert (refusal.reason, refusal.parameter_key) == (OriginOutputRefusal.CHANGED_INPUT, "request_id")
|
||||
|
|
|
|||
|
|
@ -8,21 +8,27 @@ would look exactly like a real one.
|
|||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timedelta
|
||||
from dataclasses import dataclass, field, fields
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock
|
||||
|
||||
import pytest
|
||||
from sqlalchemy.ext.asyncio import AsyncEngine
|
||||
from structlog.testing import capture_logs
|
||||
|
||||
from skyvern.constants import SCRUBBED_VALUE
|
||||
from skyvern.exceptions import WorkflowParameterNotFound, WorkflowRunNotFound
|
||||
from skyvern.forge import app
|
||||
from skyvern.forge.sdk.copilot.agent import run_copilot_agent
|
||||
from skyvern.forge.sdk.copilot.output_utils import sanitize_tool_result_for_llm
|
||||
from skyvern.forge.sdk.copilot.repair_origin_run import (
|
||||
OriginExecutionSettings,
|
||||
OriginOutputRefusal,
|
||||
OriginOutputSnapshot,
|
||||
RepairOriginRefusal,
|
||||
origin_block_outputs_from_rows,
|
||||
resolve_repair_origin_binding,
|
||||
seed_repair_origin_run,
|
||||
)
|
||||
|
|
@ -32,12 +38,24 @@ from skyvern.forge.sdk.db.enums import WorkflowRunTriggerType
|
|||
from skyvern.forge.sdk.db.models import WorkflowModel, WorkflowRunModel
|
||||
from skyvern.forge.sdk.schemas.workflow_copilot import WorkflowCopilotChatRequest
|
||||
from skyvern.forge.sdk.workflow.models.parameter import WorkflowParameter
|
||||
from skyvern.forge.sdk.workflow.models.workflow import WorkflowRunParameter, WorkflowRunStatus
|
||||
from skyvern.forge.sdk.workflow.models.workflow import (
|
||||
Workflow,
|
||||
WorkflowRun,
|
||||
WorkflowRunParameter,
|
||||
WorkflowRunStatus,
|
||||
)
|
||||
from skyvern.schemas.runs import ProxyLocation
|
||||
from tests.unit.copilot_test_helpers import (
|
||||
HARNESS_RUN_CREATED_AT,
|
||||
INERT_APPROVAL_WORKFLOW_YAML,
|
||||
ORIGIN_OUTPUT_SENTINEL,
|
||||
harness_run,
|
||||
inert_approval_workflow,
|
||||
install_get_run_results_harness,
|
||||
merge_origin_rows,
|
||||
origin_block_rows,
|
||||
origin_run_input,
|
||||
origin_run_row,
|
||||
run_result_block_row,
|
||||
stub_copilot_agent_loop,
|
||||
)
|
||||
|
|
@ -58,6 +76,8 @@ class _Ctx:
|
|||
last_run_binding_unavailable_reason: str | None = None
|
||||
repair_origin_input_values: tuple[tuple[WorkflowParameter, WorkflowRunParameter], ...] = field(default=())
|
||||
repair_origin_is_copilot_run: bool = False
|
||||
repair_origin_outputs: OriginOutputSnapshot | OriginOutputRefusal | None = None
|
||||
repair_origin_outputs_run_id: str | None = None
|
||||
|
||||
|
||||
ORIGIN_VALUES = [origin_run_input("resume", "resume_run_value")]
|
||||
|
|
@ -86,20 +106,19 @@ def _install_run(
|
|||
|
||||
monkeypatch.setattr(app, "WORKFLOW_SERVICE", SimpleNamespace(get_workflow_run=get_workflow_run), raising=False)
|
||||
monkeypatch.setattr(app.DATABASE.workflow_runs, "get_workflow_run_parameters", get_workflow_run_parameters)
|
||||
monkeypatch.setattr(app.DATABASE.workflows, "get_workflow", AsyncMock(return_value=None))
|
||||
return loaded_run_ids
|
||||
|
||||
|
||||
def _run(**overrides: object) -> SimpleNamespace:
|
||||
def _run(**overrides: object) -> WorkflowRun:
|
||||
fields: dict[str, object] = {
|
||||
"workflow_run_id": RUN,
|
||||
"organization_id": ORG,
|
||||
"workflow_permanent_id": WPID,
|
||||
"browser_session_id": RUN_BROWSER,
|
||||
"status": WorkflowRunStatus.failed,
|
||||
"copilot_session_id": None,
|
||||
}
|
||||
fields.update(overrides)
|
||||
return SimpleNamespace(**fields)
|
||||
return origin_run_row(**{**fields, **overrides})
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
|
|
@ -160,7 +179,7 @@ async def test_a_run_it_cannot_vouch_for_leaves_the_target_unavailable(
|
|||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_origin_input_values_load_only_after_the_ownership_checks(
|
||||
monkeypatch: pytest.MonkeyPatch, run: SimpleNamespace | Exception, loads_values: bool
|
||||
monkeypatch: pytest.MonkeyPatch, run: WorkflowRun | Exception, loads_values: bool
|
||||
) -> None:
|
||||
loaded_run_ids = _install_run(monkeypatch, run)
|
||||
ctx = _Ctx(repair_origin_input_values=STALE_VALUES)
|
||||
|
|
@ -833,3 +852,240 @@ def test_the_first_page_cursor_survives_a_head_cut_behind_a_large_run_field() ->
|
|||
page = run_execution.project_run_results_page(result, facts)
|
||||
|
||||
assert '"next_block_cursor": "wr-1:20"' in json.dumps(page)[:50_000]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_origin_rows_are_valued_the_way_verified_recording_values_them(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(app.WORKFLOW_SERVICE, "get_workflow_by_permanent_id", AsyncMock(return_value=None))
|
||||
workflow = await inert_approval_workflow(INERT_APPROVAL_WORKFLOW_YAML, workflow_id="w_origin")
|
||||
|
||||
row = {"row": ORIGIN_OUTPUT_SENTINEL}
|
||||
snapshot = origin_block_outputs_from_rows(
|
||||
workflow.workflow_definition,
|
||||
*merge_origin_rows(
|
||||
origin_block_rows(workflow, "approval", status="failed", minute=1, registered=False),
|
||||
origin_block_rows(workflow, "approval", minute=2, row_output=row, value=None),
|
||||
origin_block_rows(workflow, "approval", status="failed", minute=3, parent="wrb_loop", registered=False),
|
||||
origin_block_rows(workflow, "source_status", minute=4, row_output=row, registered=False),
|
||||
),
|
||||
)
|
||||
|
||||
approval = snapshot.outputs["approval"]
|
||||
assert (approval.status, approval.has_value, approval.value) == ("completed", True, None)
|
||||
source_status = snapshot.outputs["source_status"]
|
||||
assert (source_status.has_value, source_status.value) == (True, {"row": ORIGIN_OUTPUT_SENTINEL})
|
||||
assert ORIGIN_OUTPUT_SENTINEL not in repr(snapshot)
|
||||
assert ORIGIN_OUTPUT_SENTINEL not in repr(source_status)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("run", "expected_refusal", "loads_rows"),
|
||||
[
|
||||
(_run(), None, True),
|
||||
(_run(browser_session_id=None), None, True),
|
||||
(_run(status=WorkflowRunStatus.running), OriginOutputRefusal.ORIGIN_UNSETTLED, False),
|
||||
(WorkflowRunNotFound(RUN), OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN, False),
|
||||
(_run(organization_id="o_other"), OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN, False),
|
||||
(_run(workflow_permanent_id="wpid_other"), OriginOutputRefusal.FOREIGN_OR_MISMATCHED_ORIGIN, False),
|
||||
(RuntimeError("the run store is unreachable"), OriginOutputRefusal.OUTPUT_UNAVAILABLE, False),
|
||||
],
|
||||
ids=[
|
||||
"usable",
|
||||
"no_recorded_browser",
|
||||
"still_running",
|
||||
"run_not_found",
|
||||
"foreign_organization",
|
||||
"workflow_mismatch",
|
||||
"lookup_failed",
|
||||
],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_origin_outputs_load_only_from_an_owned_finished_run(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
run: WorkflowRun | Exception,
|
||||
expected_refusal: OriginOutputRefusal | None,
|
||||
loads_rows: bool,
|
||||
) -> None:
|
||||
monkeypatch.setattr(app.WORKFLOW_SERVICE, "get_workflow_by_permanent_id", AsyncMock(return_value=None))
|
||||
workflow = await inert_approval_workflow(INERT_APPROVAL_WORKFLOW_YAML, workflow_id="w_origin")
|
||||
_install_run(monkeypatch, run)
|
||||
get_workflow = AsyncMock(return_value=workflow)
|
||||
get_rows = AsyncMock(return_value=origin_block_rows(workflow, "approval")[0])
|
||||
monkeypatch.setattr(app.DATABASE.workflows, "get_workflow", get_workflow)
|
||||
monkeypatch.setattr(app.DATABASE.observer, "get_workflow_run_blocks", get_rows)
|
||||
monkeypatch.setattr(
|
||||
app.DATABASE.workflow_runs,
|
||||
"get_workflow_run_output_parameters",
|
||||
AsyncMock(return_value=origin_block_rows(workflow, "approval", value={"secret": ORIGIN_OUTPUT_SENTINEL})[1]),
|
||||
)
|
||||
ctx = _Ctx()
|
||||
|
||||
with capture_logs() as logs:
|
||||
await seed_repair_origin_run(ctx, workflow_run_id=RUN)
|
||||
|
||||
owned = expected_refusal in (None, OriginOutputRefusal.ORIGIN_UNSETTLED)
|
||||
assert ctx.repair_origin_outputs_run_id == (RUN if owned else None)
|
||||
assert isinstance(ctx.repair_origin_outputs, OriginOutputSnapshot) is loads_rows
|
||||
if not loads_rows:
|
||||
assert ctx.repair_origin_outputs is expected_refusal
|
||||
# Rows of a foreign or unsettled run are never read, not merely discarded.
|
||||
assert get_rows.await_count == (1 if loads_rows else 0)
|
||||
if loads_rows:
|
||||
get_workflow.assert_awaited_once_with(workflow_id="w_origin", organization_id=ORG)
|
||||
assert isinstance(ctx.repair_origin_outputs, OriginOutputSnapshot)
|
||||
assert ctx.repair_origin_outputs.outputs["approval"].value == {"secret": ORIGIN_OUTPUT_SENTINEL}
|
||||
assert ORIGIN_OUTPUT_SENTINEL not in repr(logs)
|
||||
assert ORIGIN_OUTPUT_SENTINEL not in repr(ctx)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("run_overrides", "parameters", "expected"),
|
||||
[
|
||||
({"copilot_session_id": "cs_test_run"}, ORIGIN_VALUES, None),
|
||||
({"debug_session_id": "ds_debugger"}, ORIGIN_VALUES, None),
|
||||
({}, RuntimeError("the parameter store is unreachable"), OriginOutputRefusal.OUTPUT_UNAVAILABLE),
|
||||
({"failure_reason": SCRUBBED_VALUE}, ORIGIN_VALUES, OriginOutputRefusal.OUTPUT_UNAVAILABLE),
|
||||
({}, [origin_run_input("resume", SCRUBBED_VALUE)], OriginOutputRefusal.OUTPUT_UNAVAILABLE),
|
||||
({"created_at": datetime(2020, 1, 1, tzinfo=UTC)}, ORIGIN_VALUES, OriginOutputRefusal.OUTPUT_UNAVAILABLE),
|
||||
({"script_run": {"script_id": "s_cached"}}, ORIGIN_VALUES, OriginOutputRefusal.CHANGED_EXECUTION_SETTINGS),
|
||||
({"parent_workflow_run_id": "wr_parent"}, ORIGIN_VALUES, OriginOutputRefusal.CHANGED_EXECUTION_SETTINGS),
|
||||
],
|
||||
ids=[
|
||||
"copilot_test_run_is_never_an_origin",
|
||||
"debugger_block_run_is_never_an_origin",
|
||||
"unreadable_inputs_never_reuse",
|
||||
"retention_scrubbed_failure_reason",
|
||||
"retention_scrubbed_input",
|
||||
"version_overwritten_after_the_run_started",
|
||||
"origin_ran_a_cached_script",
|
||||
"child_run_inherited_ancestor_prompts",
|
||||
],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_origin_outputs_need_a_real_run_and_its_recorded_inputs(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
run_overrides: dict[str, str | datetime],
|
||||
parameters: list[tuple[WorkflowParameter, WorkflowRunParameter]] | Exception,
|
||||
expected: OriginOutputRefusal | None,
|
||||
) -> None:
|
||||
monkeypatch.setattr(app.WORKFLOW_SERVICE, "get_workflow_by_permanent_id", AsyncMock(return_value=None))
|
||||
workflow = await inert_approval_workflow(INERT_APPROVAL_WORKFLOW_YAML, workflow_id="w_origin")
|
||||
_install_run(monkeypatch, _run(**run_overrides), parameters)
|
||||
monkeypatch.setattr(app.DATABASE.workflows, "get_workflow", AsyncMock(return_value=workflow))
|
||||
rows = origin_block_rows(workflow, "approval", value={"authorized": True})
|
||||
monkeypatch.setattr(app.DATABASE.observer, "get_workflow_run_blocks", AsyncMock(return_value=rows[0]))
|
||||
monkeypatch.setattr(
|
||||
app.DATABASE.workflow_runs, "get_workflow_run_output_parameters", AsyncMock(return_value=rows[1])
|
||||
)
|
||||
ctx = _Ctx()
|
||||
|
||||
await seed_repair_origin_run(ctx, workflow_run_id=RUN)
|
||||
|
||||
assert ctx.repair_origin_outputs is expected
|
||||
if expected is None:
|
||||
assert ctx.repair_origin_outputs_run_id is None
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_no_headers_and_empty_headers_are_the_same_execution_settings(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setattr(app.WORKFLOW_SERVICE, "get_workflow_by_permanent_id", AsyncMock(return_value=None))
|
||||
workflow = await inert_approval_workflow(INERT_APPROVAL_WORKFLOW_YAML, workflow_id="w_origin")
|
||||
|
||||
saved_by_copilot = OriginExecutionSettings.of(workflow.model_copy(update={"extra_http_headers": {}}))
|
||||
|
||||
assert saved_by_copilot == OriginExecutionSettings.of(
|
||||
workflow.model_copy(update={"extra_http_headers": None}), _run()
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_turn_opened_about_no_run_loads_no_origin_outputs(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
_install_run(monkeypatch, _run())
|
||||
ctx = _Ctx(repair_origin_outputs_run_id="wr_stale")
|
||||
|
||||
await seed_repair_origin_run(ctx, workflow_run_id=None)
|
||||
|
||||
assert (ctx.repair_origin_outputs_run_id, ctx.repair_origin_outputs) == (None, None)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_the_origin_runs_own_setting_overrides_win_over_its_versions(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setattr(app.WORKFLOW_SERVICE, "get_workflow_by_permanent_id", AsyncMock(return_value=None))
|
||||
workflow = await inert_approval_workflow(INERT_APPROVAL_WORKFLOW_YAML, workflow_id="w_origin")
|
||||
version = workflow.model_copy(
|
||||
update={"proxy_location": ProxyLocation.US_CA, "browser_profile_id": "bp_version", "model": {"name": "m"}}
|
||||
)
|
||||
_install_run(monkeypatch, _run(proxy_location=ProxyLocation.US_NY, extra_http_headers={"X-Run": "1"}))
|
||||
monkeypatch.setattr(app.DATABASE.workflows, "get_workflow", AsyncMock(return_value=version))
|
||||
monkeypatch.setattr(app.DATABASE.observer, "get_workflow_run_blocks", AsyncMock(return_value=[]))
|
||||
monkeypatch.setattr(app.DATABASE.workflow_runs, "get_workflow_run_output_parameters", AsyncMock(return_value=[]))
|
||||
ctx = _Ctx()
|
||||
|
||||
await seed_repair_origin_run(ctx, workflow_run_id=RUN)
|
||||
|
||||
assert isinstance(ctx.repair_origin_outputs, OriginOutputSnapshot)
|
||||
assert ctx.repair_origin_outputs.settings == OriginExecutionSettings(
|
||||
proxy_location=ProxyLocation.US_NY,
|
||||
browser_profile_id="bp_version",
|
||||
browser_profile_key=None,
|
||||
model={"name": "m"},
|
||||
extra_http_headers={"X-Run": "1"},
|
||||
)
|
||||
|
||||
|
||||
# Every other Workflow field is classified here as not changing what a block run outputs.
|
||||
_WORKFLOW_FIELDS_NOT_COMPARED = frozenset(
|
||||
{
|
||||
"workflow_id",
|
||||
"organization_id",
|
||||
"title",
|
||||
"workflow_permanent_id",
|
||||
"version",
|
||||
"is_saved_task",
|
||||
"is_template",
|
||||
"description",
|
||||
"workflow_definition",
|
||||
"webhook_callback_url",
|
||||
"totp_verification_url",
|
||||
"totp_identifier",
|
||||
"persist_browser_session",
|
||||
"reuse_browser_session",
|
||||
"mask_secrets",
|
||||
"pin_saved_session_ip",
|
||||
"status",
|
||||
"max_screenshot_scrolls",
|
||||
"max_elapsed_time_minutes",
|
||||
"cdp_connect_headers",
|
||||
"run_with",
|
||||
"browser_type",
|
||||
"ai_fallback",
|
||||
"cache_key",
|
||||
"adaptive_caching",
|
||||
"enable_self_healing",
|
||||
"code_version",
|
||||
"generate_script_on_terminal",
|
||||
"run_sequentially",
|
||||
"sequential_key",
|
||||
"folder_id",
|
||||
"import_error",
|
||||
"created_by",
|
||||
"edited_by",
|
||||
"copilot_authored",
|
||||
# Filled only by the workflow detail endpoint; never read at run time.
|
||||
"effective_default_engine",
|
||||
"original_created_by",
|
||||
"original_created_at",
|
||||
"created_at",
|
||||
"modified_at",
|
||||
"deleted_at",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def test_every_workflow_field_is_compared_across_origin_and_test_or_classified_as_not_shaping_output() -> None:
|
||||
compared = {setting.name for setting in fields(OriginExecutionSettings)}
|
||||
|
||||
assert compared <= set(Workflow.model_fields)
|
||||
assert set(Workflow.model_fields) - compared - _WORKFLOW_FIELDS_NOT_COMPARED == set()
|
||||
|
|
|
|||
|
|
@ -41,6 +41,7 @@ from skyvern.forge.sdk.copilot.blocker_signal import (
|
|||
contains_internal_machinery_leak,
|
||||
)
|
||||
from skyvern.forge.sdk.copilot.context import CopilotContext
|
||||
from skyvern.forge.sdk.copilot.repair_origin_run import OriginBlockOutput, OriginExecutionSettings, OriginOutputSnapshot
|
||||
from skyvern.forge.sdk.copilot.tools import (
|
||||
RUN_BLOCKS_SAFETY_CEILING_SECONDS,
|
||||
RUN_BLOCKS_STAGNATION_WINDOW_SECONDS,
|
||||
|
|
@ -58,7 +59,7 @@ from skyvern.forge.sdk.copilot.tools.run_execution import (
|
|||
)
|
||||
from skyvern.forge.sdk.copilot.turn_origin import TurnOrigin
|
||||
from skyvern.forge.sdk.schemas.workflow_runs import WorkflowRunBlock
|
||||
from skyvern.schemas.workflows import BlockType
|
||||
from skyvern.schemas.workflows import BlockStatus, BlockType
|
||||
from skyvern.webeye.actions.action_types import ActionType
|
||||
from skyvern.webeye.actions.actions import ActionStatus
|
||||
from tests.unit.copilot_test_helpers import SEARCH_THEN_SELECT_WORKFLOW_YAML
|
||||
|
|
@ -565,6 +566,9 @@ title: human approval example
|
|||
workflow_definition:
|
||||
parameters: []
|
||||
blocks:
|
||||
- block_type: wait
|
||||
label: request_access
|
||||
wait_sec: 1
|
||||
- block_type: human_interaction
|
||||
label: approve_login
|
||||
timeout_seconds: 3600
|
||||
|
|
@ -615,6 +619,17 @@ async def test_paused_run_is_reported_as_a_pause_and_left_running(monkeypatch: p
|
|||
ctx = make_copilot_ctx(browser_session_id="pbs_chat")
|
||||
ctx.staged_workflow = harness["workflow"]
|
||||
ctx.frontier_resume_session_id = "pbs_run"
|
||||
ctx.repair_origin_outputs_run_id = "wr_origin"
|
||||
ctx.repair_origin_outputs = OriginOutputSnapshot(
|
||||
definition=harness["workflow"].workflow_definition,
|
||||
outputs={
|
||||
"request_access": OriginBlockOutput(
|
||||
status=BlockStatus.completed, has_value=True, created_at=datetime.now(UTC), value={"sent": True}
|
||||
)
|
||||
},
|
||||
settings=OriginExecutionSettings.of(harness["workflow"]),
|
||||
)
|
||||
ctx.frontier_origin_reused_labels = ["request_access"]
|
||||
before = set(run_execution._DETACHED_CLEANUP_TASKS)
|
||||
|
||||
started = time.monotonic()
|
||||
|
|
@ -626,6 +641,8 @@ async def test_paused_run_is_reported_as_a_pause_and_left_running(monkeypatch: p
|
|||
assert result["data"]["control_signal"]["kind"] == "watchdog_paused", result
|
||||
assert "paused" in result["data"]["user_facing_summary"].lower()
|
||||
assert "uncertain" not in result["error"].lower()
|
||||
assert result["data"]["reused_origin_output_labels"] == ["request_access"]
|
||||
assert result["data"]["origin_workflow_run_id"] == "wr_origin"
|
||||
|
||||
harness["cancel_run_task"].assert_not_awaited()
|
||||
harness["cooperative_cancel"].assert_not_awaited()
|
||||
|
|
|
|||
|
|
@ -6052,6 +6052,7 @@ def _install_diagnose_run_lookup(
|
|||
browser_session_id="pbs-1",
|
||||
status=run_status,
|
||||
copilot_session_id=None,
|
||||
is_debug_session=False,
|
||||
)
|
||||
monkeypatch.setattr(app, "WORKFLOW_SERVICE", SimpleNamespace(get_workflow_run=AsyncMock(return_value=run)))
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue