Class: Smith::Workflow

Inherits:
Object
  • Object
show all
Includes:
ArtifactIntegration, BudgetIntegration, DSL, DataVolumePolicy, DeadlineEnforcement, DefinitionIdentity, Durability, EventIntegration, Execution, GuardrailIntegration, MessageAdmissionBoundary, Persistence, SplitStepPersistence
Defined in:
lib/smith/workflow.rb,
lib/smith/workflow/dsl.rb,
lib/smith/workflow/claim.rb,
lib/smith/workflow/graph.rb,
lib/smith/workflow/router.rb,
lib/smith/workflow/parallel.rb,
lib/smith/workflow/pipeline.rb,
lib/smith/workflow/execution.rb,
lib/smith/workflow/graph_dsl.rb,
lib/smith/workflow/branch_env.rb,
lib/smith/workflow/durability.rb,
lib/smith/workflow/identifier.rb,
lib/smith/workflow/run_result.rb,
lib/smith/workflow/transition.rb,
lib/smith/workflow/persistence.rb,
lib/smith/workflow/usage_entry.rb,
lib/smith/workflow/agent_result.rb,
lib/smith/workflow/graph/report.rb,
lib/smith/workflow/step_context.rb,
lib/smith/workflow/graph/metrics.rb,
lib/smith/workflow/graph/targets.rb,
lib/smith/workflow/message_batch.rb,
lib/smith/workflow/prepared_step.rb,
lib/smith/workflow/process_local.rb,
lib/smith/workflow/composite/plan.rb,
lib/smith/workflow/failure_record.rb,
lib/smith/workflow/composite/enums.rb,
lib/smith/workflow/composite/error.rb,
lib/smith/workflow/composite/input.rb,
lib/smith/workflow/execution_frame.rb,
lib/smith/workflow/graph/reference.rb,
lib/smith/workflow/graph/validator.rb,
lib/smith/workflow/retry_execution.rb,
lib/smith/workflow/step_completion.rb,
lib/smith/workflow/string_snapshot.rb,
lib/smith/workflow/composite/branch.rb,
lib/smith/workflow/fanout_execution.rb,
lib/smith/workflow/graph/diagnostic.rb,
lib/smith/workflow/nested_execution.rb,
lib/smith/workflow/worker_execution.rb,
lib/smith/workflow/composite/effects.rb,
lib/smith/workflow/composite/payload.rb,
lib/smith/workflow/composite/planner.rb,
lib/smith/workflow/composite/reducer.rb,
lib/smith/workflow/event_integration.rb,
lib/smith/workflow/message_admission.rb,
lib/smith/workflow/budget_integration.rb,
lib/smith/workflow/composite/contract.rb,
lib/smith/workflow/data_volume_policy.rb,
lib/smith/workflow/deterministic_step.rb,
lib/smith/workflow/graph/reachability.rb,
lib/smith/workflow/optimization_state.rb,
lib/smith/workflow/parallel_execution.rb,
lib/smith/workflow/composite/reduction.rb,
lib/smith/workflow/definition_identity.rb,
lib/smith/workflow/evaluator_optimizer.rb,
lib/smith/workflow/failure_record_text.rb,
lib/smith/workflow/orchestration_state.rb,
lib/smith/workflow/orchestrator_worker.rb,
lib/smith/workflow/artifact_integration.rb,
lib/smith/workflow/deadline_enforcement.rb,
lib/smith/workflow/composite/budget_math.rb,
lib/smith/workflow/composite/outcome_set.rb,
lib/smith/workflow/composite/preparation.rb,
lib/smith/workflow/failure_reconstructor.rb,
lib/smith/workflow/graph/diagnostic_path.rb,
lib/smith/workflow/graph/fanout_contract.rb,
lib/smith/workflow/guardrail_integration.rb,
lib/smith/workflow/parallel/cancellation.rb,
lib/smith/workflow/composite/value_budget.rb,
lib/smith/workflow/failure_record_restore.rb,
lib/smith/workflow/graph/contract_helpers.rb,
lib/smith/workflow/guarded_step_execution.rb,
lib/smith/workflow/parallel_agent_binding.rb,
lib/smith/workflow/prepared_step_dispatch.rb,
lib/smith/workflow/prepared_step_recovery.rb,
lib/smith/workflow/split_step_persistence.rb,
lib/smith/workflow/deterministic_execution.rb,
lib/smith/workflow/failure_detail_snapshot.rb,
lib/smith/workflow/graph/runtime_readiness.rb,
lib/smith/workflow/graph/state_diagnostics.rb,
lib/smith/workflow/parallel/root_execution.rb,
lib/smith/workflow/thread_context_snapshot.rb,
lib/smith/workflow/composite/branch_failure.rb,
lib/smith/workflow/composite/branch_outcome.rb,
lib/smith/workflow/composite/error_evidence.rb,
lib/smith/workflow/composite/payload_digest.rb,
lib/smith/workflow/composite/plan_integrity.rb,
lib/smith/workflow/failure_record_validator.rb,
lib/smith/workflow/message_value_normalizer.rb,
lib/smith/workflow/transition_actionability.rb,
lib/smith/workflow/composite/branch_contract.rb,
lib/smith/workflow/execution_result_snapshot.rb,
lib/smith/workflow/graph/transition_contract.rb,
lib/smith/workflow/graph/transition_snapshot.rb,
lib/smith/workflow/parallel/nested_execution.rb,
lib/smith/workflow/prepared_branch_execution.rb,
lib/smith/workflow/composite/branch_execution.rb,
lib/smith/workflow/composite/budget_allocator.rb,
lib/smith/workflow/composite/effects_baseline.rb,
lib/smith/workflow/graph/execution_successors.rb,
lib/smith/workflow/message_admission_boundary.rb,
lib/smith/workflow/parallel/execution_context.rb,
lib/smith/workflow/composite/effects_preflight.rb,
lib/smith/workflow/graph/identifier_projection.rb,
lib/smith/workflow/graph/optimization_contract.rb,
lib/smith/workflow/composite/execution_contract.rb,
lib/smith/workflow/execution_binding_resolution.rb,
lib/smith/workflow/graph/orchestration_contract.rb,
lib/smith/workflow/graph/transition_diagnostics.rb,
lib/smith/workflow/parallel/cancellation_signal.rb,
lib/smith/workflow/composite/effects_application.rb,
lib/smith/workflow/composite/outcome_accumulator.rb,
lib/smith/workflow/graph/retry_policy_diagnostic.rb,
lib/smith/workflow/prepared_step_execution_scope.rb,
lib/smith/workflow/composite/encoded_value_budget.rb,
lib/smith/workflow/graph/reachability_diagnostics.rb,
lib/smith/workflow/graph/runtime_readiness_report.rb,
lib/smith/workflow/prepared_step_execution_result.rb,
lib/smith/workflow/graph/runtime_readiness_metrics.rb,
lib/smith/workflow/split_step_persistence/boundary.rb,
lib/smith/workflow/split_step_persistence/payloads.rb,
lib/smith/workflow/split_step_persistence/recovery.rb,
lib/smith/workflow/composite/branch_budget_contract.rb,
lib/smith/workflow/composite/fanout_branch_contract.rb,
lib/smith/workflow/split_step_persistence/execution.rb,
lib/smith/workflow/definition_identity/class_methods.rb,
lib/smith/workflow/graph/runtime_binding_diagnostics.rb,
lib/smith/workflow/graph/runtime_readiness_traversal.rb,
lib/smith/workflow/split_step_persistence/checkpoint.rb,
lib/smith/workflow/graph/nested_readiness_diagnostics.rb,
lib/smith/workflow/split_step_persistence/inheritance.rb,
lib/smith/workflow/split_step_persistence/preparation.rb,
lib/smith/workflow/graph/transition_contract_attributes.rb,
lib/smith/workflow/prepared_step_execution_authorization.rb,
lib/smith/workflow/split_step_persistence/boundary_reset.rb,
lib/smith/workflow/split_step_persistence/dispatch_claim.rb,
lib/smith/workflow/split_step_persistence/state_snapshot.rb,
lib/smith/workflow/graph/runtime_readiness_report_builder.rb,
lib/smith/workflow/split_step_persistence/checkpoint_state.rb,
lib/smith/workflow/composite_branch_execution_authorization.rb,
lib/smith/workflow/graph/runtime_binding_diagnostic_builder.rb,
lib/smith/workflow/graph/transition_contract_configurations.rb,
lib/smith/workflow/split_step_persistence/dispatch_boundary.rb,
lib/smith/workflow/split_step_persistence/preparation_claim.rb,
lib/smith/workflow/split_step_persistence/recovery_boundary.rb,
lib/smith/workflow/split_step_persistence/subclass_boundary.rb,
lib/smith/workflow/split_step_persistence/composite_execution.rb,
lib/smith/workflow/split_step_persistence/definition_boundary.rb,
lib/smith/workflow/split_step_persistence/execution_lifecycle.rb,
lib/smith/workflow/split_step_persistence/preparation_payload.rb,
lib/smith/workflow/split_step_persistence/transition_contract.rb,
lib/smith/workflow/graph/transition_optimization_configuration.rb,
lib/smith/workflow/split_step_persistence/preparation_recovery.rb,
lib/smith/workflow/split_step_persistence/restart_safe_adapter.rb,
lib/smith/workflow/split_step_persistence/transaction_identity.rb,
lib/smith/workflow/split_step_persistence/composite_preparation.rb,
lib/smith/workflow/split_step_persistence/dispatch_confirmation.rb,
lib/smith/workflow/split_step_persistence/dispatch_verification.rb,
lib/smith/workflow/split_step_persistence/execution_verification.rb,
lib/smith/workflow/split_step_persistence/recovery_class_methods.rb,
lib/smith/workflow/split_step_persistence/execution_authorization.rb,
lib/smith/workflow/split_step_persistence/replace_exact_signature.rb,
lib/smith/workflow/split_step_persistence/canonical_payload_digest.rb,
lib/smith/workflow/split_step_persistence/composite_branch_effects.rb,
lib/smith/workflow/split_step_persistence/composite_branch_outcome.rb,
lib/smith/workflow/split_step_persistence/execution_result_capture.rb,
lib/smith/workflow/split_step_persistence/composite_branch_execution.rb,
lib/smith/workflow/split_step_persistence/execution_binding_snapshot.rb,
lib/smith/workflow/split_step_persistence/execution_binding_collector.rb,
lib/smith/workflow/split_step_persistence/execution_workflow_snapshot.rb,
lib/smith/workflow/split_step_persistence/transition_contract_freezer.rb,
lib/smith/workflow/split_step_persistence/composite_reduction_execution.rb,
lib/smith/workflow/split_step_persistence/transition_contract_signature.rb,
lib/smith/workflow/split_step_persistence/composite_branch_authorization.rb,
lib/smith/workflow/split_step_persistence/execution_authorization_issuance.rb,
lib/smith/workflow/split_step_persistence/transition_contract_structured_values.rb

Defined Under Namespace

Modules: ArtifactIntegration, BudgetIntegration, Claim, Composite, DSL, DataVolumePolicy, DeadlineEnforcement, DefinitionIdentity, DeterministicExecution, Durability, EvaluatorOptimizer, EventIntegration, Execution, ExecutionBindingResolution, FanoutExecution, GuardedStepExecution, GuardrailIntegration, MessageAdmissionBoundary, NestedExecution, OrchestratorWorker, ParallelExecution, Persistence, RetryExecution, SplitStepPersistence, StepCompletion, StepContext, TransitionActionability Classes: AgentResult, BranchEnv, CompositeBranchExecutionAuthorization, DeterministicStep, ExecutionFrame, ExecutionResultSnapshot, FailureDetailSnapshot, FailureReconstructor, FailureRecord, FailureRecordRestore, FailureRecordText, FailureRecordValidator, Graph, Identifier, MessageAdmission, MessageBatch, MessageValueNormalizer, OptimizationState, OrchestrationState, Parallel, ParallelAgentBinding, Pipeline, PreparedStep, PreparedStepDispatch, PreparedStepExecutionAuthorization, PreparedStepExecutionResult, PreparedStepExecutionScope, PreparedStepRecovery, Router, RunResult, StringSnapshot, Transition, UsageEntry, WorkerExecution

Constant Summary collapse

DEFAULT_MAX_TRANSITIONS =
100

Constants included from DataVolumePolicy

DataVolumePolicy::LIGHTWEIGHT_SCALARS

Constants included from BudgetIntegration

BudgetIntegration::AGENT_DIM_MAP, BudgetIntegration::BUDGET_DIMENSIONS, BudgetIntegration::COST_DIMENSIONS, BudgetIntegration::TOKEN_DIMENSIONS

Constants included from SplitStepPersistence

SplitStepPersistence::NO_SPLIT_TRANSITION

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Methods included from MessageAdmissionBoundary

#append_session_messages!

Methods included from Durability

#advance_persisted!, #clear_persisted!, #clear_step_in_progress!, included, #mark_step_in_progress!, #persist!, #run_persisted!

Methods included from Persistence

#to_state

Methods included from DefinitionIdentity

included

Methods included from DSL

included

Methods included from SplitStepPersistence

prepended

Methods included from SplitStepPersistence::Checkpoint

#complete_persisted_step!, #persist!

Methods included from SplitStepPersistence::CompositeExecution

#execute_prepared_composite_branch!, #prepare_composite_step!, #reduce_prepared_composite_step!

Methods included from SplitStepPersistence::ExecutionAuthorizationIssuance

#authorize_prepared_step_execution!

Methods included from SplitStepPersistence::ExecutionAuthorization

#release_prepared_step_execution!

Methods included from SplitStepPersistence::DispatchConfirmation

#confirm_prepared_step_dispatch!

Methods included from SplitStepPersistence::DispatchClaim

#claim_prepared_step_dispatch!

Methods included from SplitStepPersistence::Preparation

#confirm_prepared_step!, #prepare_persisted_step!, #prepared_persisted_step

Methods included from SplitStepPersistence::StateSnapshot

#to_state

Methods included from SplitStepPersistence::Boundary

#advance_persisted!, #clear_persisted!, #run_persisted!

Constructor Details

#initialize(context: {}, ledger: nil, created_at: nil) ⇒ Workflow

Returns a new instance of Workflow.



41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
# File 'lib/smith/workflow.rb', line 41

def initialize(context: {}, ledger: nil, created_at: nil)
  @state = self.class.initial_state
  @context = context
  @step_count = 0
  @next_transition_name = nil
  @ledger = ledger || build_ledger
  @created_at = created_at || Time.now.utc.iso8601
  @updated_at = @created_at
  @total_cost = 0.0
  @total_tokens = 0
  @outcome = nil
  # Eager init for usage tracking. Both `@usage_mutex` (lazy
  # init at the call site would race across parallel fan-out
  # branches) and the durable per-call/output/failure fields
  # must be present before any agent recording fires.
  # `restore_state` mirrors these inits because `from_state` uses
  # `allocate` and bypasses `initialize` — see persistence.rb.
  @usage_entries = []
  @usage_mutex = Mutex.new
  @last_output = nil
  # Durable attribution of the most recent serial `execute :agent` step:
  # { model:, provider: } of the model that actually served it (post
  # fallback resolution), exposed to the following deterministic step as
  # `last_agent_model` / `last_agent_provider`. `@pending_agent_execution`
  # is the transient per-step carrier from `execute_serial_step` to
  # `complete_step`; it is never persisted.
  @last_agent_execution = nil
  @pending_agent_execution = nil
  @last_failed_step = nil
  # Optimistic-locking version. Incremented on each persist!; restored
  # from the persisted payload. Adapters that support store_versioned
  # raise Smith::PersistenceVersionConflict when expected_version
  # doesn't match the stored payload's version (i.e., a concurrent
  # write occurred between this process's restore and persist).
  @persistence_version = 0
  # Digest of the seed_messages produced at construction time.
  # Compared on restore against the live builder's output when
  # seed_validation is :warn or :strict; nil when no seed builder
  # ran or its output was empty.
  @seed_digest = nil
  @seed_message_count = 0
  # Idempotency marker stamped between persist-before-advance and
  # persist-after-advance under idempotency_mode :strict; restored
  # workflows with the marker set raise
  # Smith::StepInProgressOnRestore. Lax mode leaves it false.
  @step_in_progress = false
  # Process-local phase for the split-step persistence API. This is not
  # serialized: strict restore rejects the durable step_in_progress marker
  # before an uncertain step can be resumed.
  @split_step_phase = nil
  @split_step_transition_name = nil
  @split_step_mutex = Mutex.new
  # Set of context keys recorded via deterministic step write_context
  # writes. Used by persist :auto Context mode to compute the
  # persisted-context slice. Seeded from the Context class's
  # also: declaration so explicit input keys round-trip.
  @persisted_keys = ::Set.new(initial_persist_auto_seed)
  @persisted_keys_mutex = Mutex.new
  initialize_tool_result_state
  seed_initial_session_messages
end

Instance Attribute Details

#last_prepared_inputObject (readonly)

Returns the value of attribute last_prepared_input.



35
36
37
# File 'lib/smith/workflow.rb', line 35

def last_prepared_input
  @last_prepared_input
end

#ledgerObject (readonly)

Returns the value of attribute ledger.



35
36
37
# File 'lib/smith/workflow.rb', line 35

def ledger
  @ledger
end

#stateObject (readonly)

Returns the value of attribute state.



35
36
37
# File 'lib/smith/workflow.rb', line 35

def state
  @state
end

Class Method Details

.graphObject



5
6
7
8
9
10
11
12
# File 'lib/smith/workflow/graph_dsl.rb', line 5

def self.graph
  Graph.new(
    workflow_class: self,
    initial_state: initial_state,
    states: @states || [],
    transitions: @transitions || {}
  )
end

.runtime_readinessObject



18
19
20
# File 'lib/smith/workflow/graph_dsl.rb', line 18

def self.runtime_readiness
  graph.runtime_readiness
end

.validate_graphObject



14
15
16
# File 'lib/smith/workflow/graph_dsl.rb', line 14

def self.validate_graph
  graph.validate
end

Instance Method Details

#advance!Object



107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
# File 'lib/smith/workflow.rb', line 107

def advance!
  advance_claim = SplitStepPersistence.instance_method(:claim_split_step_advance!).bind_call(self)
  ensure_transition_budget!
  @step_work_started = false

  transition = SplitStepPersistence.instance_method(:resolve_split_step_advance_transition).bind_call(self)
  transition = resolve_transition if transition.equal?(SplitStepPersistence::NO_SPLIT_TRANSITION)
  return if transition.nil?

  @step_work_started = true
  step_result = if advance_claim == :split_step
                  Execution.instance_method(:execute_step).bind_call(self, transition)
                else
                  execute_step(transition)
                end
  @step_count += 1
  @updated_at = Time.now.utc.iso8601
  record_step_snapshot(step_result)
  step_result
rescue UnresolvedTransitionError => e
  step_result = GuardrailIntegration
                .instance_method(:handle_unresolved_transition_failure)
                .bind_call(self, e)
  record_step_snapshot(step_result)
  step_result
ensure
  SplitStepPersistence.instance_method(:release_split_step_advance!).bind_call(self, advance_claim) if advance_claim
end

#done?Boolean

Returns:

  • (Boolean)


154
155
156
# File 'lib/smith/workflow.rb', line 154

def done?
  state_named?(:done)
end

#failed?Boolean

Returns:

  • (Boolean)


158
159
160
# File 'lib/smith/workflow.rb', line 158

def failed?
  state_named?(:failed)
end

#pending_transition_nameObject



146
147
148
# File 'lib/smith/workflow.rb', line 146

def pending_transition_name
  @next_transition_name || self.class.first_transition_from(@state)&.name
end

#persisted_keysObject



103
104
105
# File 'lib/smith/workflow.rb', line 103

def persisted_keys
  @persisted_keys.dup.freeze
end

#record_persisted_key!(key) ⇒ Object



162
163
164
165
166
# File 'lib/smith/workflow.rb', line 162

def record_persisted_key!(key)
  @persisted_keys_mutex.synchronize do
    @persisted_keys << key.to_sym
  end
end

#run!Object



136
137
138
139
140
141
142
143
144
# File 'lib/smith/workflow.rb', line 136

def run!
  SplitStepPersistence.instance_method(:ensure_split_step_execution_allowed!).bind_call(self)
  steps = []
  until terminal?
    step = advance!
    steps << step if step
  end
  build_run_result(steps)
end

#session_messagesObject



37
38
39
# File 'lib/smith/workflow.rb', line 37

def session_messages
  snapshot_session_messages
end

#terminal?Boolean

Returns:

  • (Boolean)


150
151
152
# File 'lib/smith/workflow.rb', line 150

def terminal?
  !self.class.transition_from?(@state) && @next_transition_name.nil?
end

#usage_entriesObject

Public, read-only view of the per-provider-call usage ledger, so hosts can diff usage across a step boundary without paying a full to_state serialization. The entries themselves are frozen; the returned array is a frozen copy taken under the recording mutex. Tolerates allocated but unrestored instances (no mutex yet), which have no entries.



173
174
175
176
177
178
# File 'lib/smith/workflow.rb', line 173

def usage_entries
  mutex = @usage_mutex
  return [].freeze unless mutex

  mutex.synchronize { @usage_entries.dup.freeze }
end