Documentation
¶
Index ¶
- Variables
- func NewDynamoDBIntegration(backend dynamodbpkg.StorageBackend) asl.DynamoDBIntegration
- func NewSNSIntegration(backend sns.StorageBackend) asl.SNSIntegration
- func NewSQSIntegration(backend sqs.StorageBackend) asl.SQSIntegration
- type Activity
- type ActivityTask
- type AliasRoutingConfig
- type CloudWatchLogsLogGroup
- type EBStartExecutionAdapter
- type EncryptionConfiguration
- type Execution
- type Handler
- func (h *Handler) ChaosOperations() []string
- func (h *Handler) ChaosRegions() []string
- func (h *Handler) ChaosServiceName() string
- func (h *Handler) ExtractOperation(c *echo.Context) string
- func (h *Handler) ExtractResource(c *echo.Context) string
- func (h *Handler) GetSupportedOperations() []string
- func (h *Handler) Handler() echo.HandlerFunc
- func (h *Handler) MatchPriority() int
- func (h *Handler) Name() string
- func (h *Handler) Reset()
- func (h *Handler) Restore(ctx context.Context, data []byte) error
- func (h *Handler) RouteMatcher() service.Matcher
- func (h *Handler) Shutdown(ctx context.Context)
- func (h *Handler) Snapshot(ctx context.Context) []byte
- func (h *Handler) StartWorker(ctx context.Context) error
- type HistoryEvent
- type InMemoryBackend
- func (b *InMemoryBackend) AccountID() string
- func (b *InMemoryBackend) CreateActivity(ctx context.Context, name string) (*Activity, error)
- func (b *InMemoryBackend) CreateStateMachine(ctx context.Context, name, definition, roleArn, smType string) (*StateMachine, error)
- func (b *InMemoryBackend) CreateStateMachineAlias(smARN, name, description string, routing []AliasRoutingConfig) (*StateMachineAlias, error)
- func (b *InMemoryBackend) DeleteActivity(activityArn string) error
- func (b *InMemoryBackend) DeleteStateMachine(arn string) error
- func (b *InMemoryBackend) DeleteStateMachineAlias(aliasARN string) error
- func (b *InMemoryBackend) DeleteStateMachineVersion(versionARN string) error
- func (b *InMemoryBackend) DescribeActivity(activityArn string) (*Activity, error)
- func (b *InMemoryBackend) DescribeExecution(executionArn string) (*Execution, error)
- func (b *InMemoryBackend) DescribeMapRun(mapRunARN string) (*MapRun, error)
- func (b *InMemoryBackend) DescribeStateMachine(arn string) (*StateMachine, error)
- func (b *InMemoryBackend) DescribeStateMachineAlias(aliasARN string) (*StateMachineAlias, error)
- func (b *InMemoryBackend) DescribeStateMachineForExecution(executionARN string) (*StateMachine, error)
- func (b *InMemoryBackend) DescribeStateMachineVersion(versionARN string) (*StateMachineVersion, error)
- func (b *InMemoryBackend) Destroy()
- func (b *InMemoryBackend) GetActivityTask(ctx context.Context, activityArn, _ string) (*ActivityTask, error)
- func (b *InMemoryBackend) GetExecutionHistory(executionArn, nextToken string, maxResults int, reverseOrder bool) ([]HistoryEvent, string, error)
- func (b *InMemoryBackend) InvokeActivity(ctx context.Context, activityArn, inputJSON string, heartbeatSeconds int) (string, error)
- func (b *InMemoryBackend) ListActivities(ctx context.Context, nextToken string, maxResults int) ([]Activity, string, error)
- func (b *InMemoryBackend) ListExecutions(stateMachineArn, statusFilter, nextToken string, maxResults int) ([]Execution, string, error)
- func (b *InMemoryBackend) ListMapRuns(executionARN, nextToken string, maxResults int) ([]MapRun, string, error)
- func (b *InMemoryBackend) ListStateMachineAliases(smARN, nextToken string, maxResults int) ([]StateMachineAlias, string, error)
- func (b *InMemoryBackend) ListStateMachineVersions(smARN, nextToken string, maxResults int) ([]StateMachineVersion, string, error)
- func (b *InMemoryBackend) ListStateMachines(ctx context.Context, nextToken string, maxResults int) ([]StateMachine, string, error)
- func (b *InMemoryBackend) OnMapRunEnd(mapRunARN, status string, succeeded, failed, total int)
- func (b *InMemoryBackend) OnMapRunStart(executionARN, stateName string, maxConcurrency, itemCount int) string
- func (b *InMemoryBackend) PruneExecutions(_ context.Context) int
- func (b *InMemoryBackend) PublishStateMachineVersion(smARN, description, revisionID string) (*StateMachineVersion, error)
- func (b *InMemoryBackend) RedriveExecution(executionARN string) (*Execution, error)
- func (b *InMemoryBackend) Region() string
- func (b *InMemoryBackend) Reset()
- func (b *InMemoryBackend) Restore(ctx context.Context, data []byte) error
- func (b *InMemoryBackend) SendTaskFailure(taskToken, errCode, cause string) error
- func (b *InMemoryBackend) SendTaskHeartbeat(taskToken string) error
- func (b *InMemoryBackend) SendTaskSuccess(taskToken, output string) error
- func (b *InMemoryBackend) SetDynamoDBIntegration(ddb asl.DynamoDBIntegration)
- func (b *InMemoryBackend) SetECSIntegration(ecs asl.ECSIntegration)
- func (b *InMemoryBackend) SetEventBridgeIntegration(eb asl.EventBridgeIntegration)
- func (b *InMemoryBackend) SetGlueIntegration(glue asl.GlueIntegration)
- func (b *InMemoryBackend) SetLambdaInvoker(invoker asl.LambdaInvoker)
- func (b *InMemoryBackend) SetSNSIntegration(sns asl.SNSIntegration)
- func (b *InMemoryBackend) SetSQSIntegration(sqs asl.SQSIntegration)
- func (b *InMemoryBackend) SetSettings(s Settings)
- func (b *InMemoryBackend) SetStateMachineConfigurations(arn string, tracing *TracingConfiguration, logging *LoggingConfiguration, ...) error
- func (b *InMemoryBackend) Snapshot(ctx context.Context) []byte
- func (b *InMemoryBackend) StartExecution(stateMachineArn, name, input string) (*Execution, error)
- func (b *InMemoryBackend) StartSyncExecution(stateMachineArn, name, input string) (*SyncExecutionResult, error)
- func (b *InMemoryBackend) StopExecution(executionArn, errCode, cause string) error
- func (b *InMemoryBackend) SweepTaskTokens() int
- func (b *InMemoryBackend) UpdateMapRun(mapRunARN string, maxConcurrency int, toleratedFailureCount int, ...) (*MapRun, error)
- func (b *InMemoryBackend) UpdateStateMachine(smARN, definition, roleArn string) (float64, error)
- func (b *InMemoryBackend) UpdateStateMachineAlias(aliasARN, description string, routing []AliasRoutingConfig) (*StateMachineAlias, error)
- func (b *InMemoryBackend) WaitForTaskToken(ctx context.Context, taskToken string, heartbeatSeconds int) (string, error)
- type Janitor
- type LoggingConfiguration
- type LoggingDestination
- type MapRun
- type MapRunItemCounts
- type Provider
- type Settings
- type SettingsProvider
- type StateEnteredEventDetails
- type StateExitedEventDetails
- type StateMachine
- type StateMachineAlias
- type StateMachineVersion
- type StorageBackend
- type SyncExecutionResult
- type TaskFailedEventDetails
- type TaskScheduledEventDetails
- type TaskSucceededEventDetails
- type TracingConfiguration
Constants ¶
This section is empty.
Variables ¶
var ( ErrStateMachineAlreadyExists = errors.New("StateMachineAlreadyExists") ErrStateMachineDoesNotExist = errors.New("StateMachineDoesNotExist") ErrStateMachineVersionDoesNotExist = errors.New("StateMachineVersionDoesNotExist") ErrStateMachineAliasAlreadyExists = errors.New("StateMachineAliasAlreadyExists") ErrStateMachineAliasDoesNotExist = errors.New("StateMachineAliasDoesNotExist") ErrExecutionAlreadyExists = errors.New("ExecutionAlreadyExists") ErrExecutionDoesNotExist = errors.New("ExecutionDoesNotExist") ErrExecutionNotRedrivable = errors.New("ExecutionNotRedrivable") ErrInvalidDefinition = errors.New("InvalidDefinition") ErrInvalidExecutionType = errors.New("InvalidExecutionType") ErrStateMachineTypeNotSupported = errors.New("StateMachineTypeNotSupported") ErrInvalidRoleArn = errors.New("InvalidArn") ErrInvalidName = errors.New("InvalidName") ErrInvalidRoutingConfiguration = errors.New("InvalidRoutingConfiguration") ErrTagPolicyViolation = errors.New("TagPolicyViolation") ErrActivityAlreadyExists = errors.New("ActivityAlreadyExists") ErrActivityDoesNotExist = errors.New("ActivityDoesNotExist") ErrTaskTokenNotFound = errors.New("TaskTokenNotFound") ErrTaskTokenAlreadyExists = errors.New("TaskTokenAlreadyExists") ErrActivityTaskFailed = errors.New("ActivityTaskFailed") ErrHeartbeatTimeout = errors.New("States.HeartbeatTimeout") ErrInvalidExecutionInput = errors.New("InvalidExecutionInput") ErrValidation = errors.New("ValidationException") ErrMapRunDoesNotExist = errors.New("MapRunDoesNotExist") )
var ErrNilAppContext = errors.New("stepfunctions: nil app context")
ErrNilAppContext is returned when Init is called with a nil AppContext.
Functions ¶
func NewDynamoDBIntegration ¶
func NewDynamoDBIntegration(backend dynamodbpkg.StorageBackend) asl.DynamoDBIntegration
NewDynamoDBIntegration creates a new DynamoDB integration adapter.
func NewSNSIntegration ¶
func NewSNSIntegration(backend sns.StorageBackend) asl.SNSIntegration
NewSNSIntegration creates a new SNS integration adapter.
func NewSQSIntegration ¶
func NewSQSIntegration(backend sqs.StorageBackend) asl.SQSIntegration
NewSQSIntegration creates a new SQS integration adapter.
Types ¶
type Activity ¶
type Activity struct {
Name string `json:"name"`
ActivityArn string `json:"activityArn"`
CreationDate float64 `json:"creationDate"`
}
Activity represents an AWS Step Functions activity resource.
type ActivityTask ¶
ActivityTask represents a task polled from an activity queue.
type AliasRoutingConfig ¶
type AliasRoutingConfig struct {
StateMachineVersionArn string `json:"stateMachineVersionArn"`
Weight int `json:"weight"`
}
AliasRoutingConfig represents a weighted routing target for a state machine alias.
type CloudWatchLogsLogGroup ¶
type CloudWatchLogsLogGroup struct {
LogGroupArn string `json:"logGroupArn,omitempty"`
}
CloudWatchLogsLogGroup names a destination CloudWatch Logs log group.
type EBStartExecutionAdapter ¶
type EBStartExecutionAdapter struct {
B *InMemoryBackend
}
EBStartExecutionAdapter wraps InMemoryBackend to satisfy the EventBridge StepFunctionsExecutor interface (which returns only error, not *Execution).
func (*EBStartExecutionAdapter) StartExecution ¶
func (a *EBStartExecutionAdapter) StartExecution(stateMachineARN, name, input string) error
StartExecution starts an execution and discards the Execution struct, returning only the error so it fits the eventbridge.StepFunctionsExecutor interface.
type EncryptionConfiguration ¶
type EncryptionConfiguration struct {
KMSKeyID string `json:"kmsKeyId,omitempty"`
Type string `json:"type,omitempty"`
KMSDataKeyReusePeriodSeconds int `json:"kmsDataKeyReusePeriodSeconds,omitempty"`
}
EncryptionConfiguration configures KMS encryption for a state machine.
type Execution ¶
type Execution struct {
RedriveDate *float64 `json:"redriveDate,omitempty"`
StopDate *float64 `json:"stopDate,omitempty"`
Status string `json:"status"`
ExecutionArn string `json:"executionArn"`
StateMachineArn string `json:"stateMachineArn"`
Name string `json:"name"`
Input string `json:"input,omitempty"`
Output string `json:"output,omitempty"`
Error string `json:"error,omitempty"`
Cause string `json:"cause,omitempty"`
StartDate float64 `json:"startDate"`
RedriveCount int `json:"redriveCount,omitempty"`
// contains filtered or unexported fields
}
Execution represents a state machine execution.
history holds the execution's history events inline (Phase 3.3: this replaces the backend's former separate `map[string]*HistoryEvent` history map). It is deliberately unexported: *Execution is returned directly as the wire body for DescribeExecution/StartExecution/etc., and AWS's real DescribeExecution response has no "history" field -- history is only ever retrieved via GetExecutionHistory. Being unexported also means it is skipped by every encoding/json.Marshal of *Execution, including the one store.Table.Snapshot would otherwise perform; persistence.go's executionSnapshot DTO adds it back as an ordinary exported field solely for the on-disk snapshot round trip. See persistence.go for details.
type Handler ¶
type Handler struct {
Backend StorageBackend
DefaultRegion string
// contains filtered or unexported fields
}
Handler is the Echo HTTP service handler for Step Functions operations.
func NewHandler ¶
func NewHandler(backend StorageBackend) *Handler
NewHandler creates a new Step Functions handler.
func (*Handler) ChaosOperations ¶
ChaosOperations returns all operations that can be fault-injected.
func (*Handler) ChaosRegions ¶
ChaosRegions returns all regions this Step Functions instance handles.
func (*Handler) ChaosServiceName ¶
ChaosServiceName returns the lowercase AWS service name for fault rule matching.
func (*Handler) ExtractOperation ¶
ExtractOperation extracts the operation name from the X-Amz-Target header.
func (*Handler) ExtractResource ¶
ExtractResource extracts the resource name from the request body.
func (*Handler) GetSupportedOperations ¶
GetSupportedOperations returns all mocked Step Functions operations.
func (*Handler) Handler ¶
func (h *Handler) Handler() echo.HandlerFunc
Handler returns the Echo handler function for Step Functions requests.
func (*Handler) MatchPriority ¶
MatchPriority returns the routing priority for the Step Functions handler.
func (*Handler) Reset ¶
func (h *Handler) Reset()
Reset clears all in-memory state from the backend. It is used by the POST /_gopherstack/reset endpoint for CI pipelines and rapid local development.
func (*Handler) Restore ¶
Restore implements persistence.Persistable. It restores both the backend state and the handler's tag map.
func (*Handler) RouteMatcher ¶
RouteMatcher returns a matcher for Step Functions requests.
func (*Handler) Shutdown ¶
Shutdown implements service.Shutdowner. It cancels all running execution goroutines and releases associated resources. Destroy() calls each cancel func synchronously and returns immediately; the ASL goroutines exit asynchronously. If ctx expires before Destroy returns, Shutdown returns early so the process shutdown is not blocked.
type HistoryEvent ¶
type HistoryEvent struct {
StateEnteredEventDetails *StateEnteredEventDetails `json:"stateEnteredEventDetails,omitempty"`
StateExitedEventDetails *StateExitedEventDetails `json:"stateExitedEventDetails,omitempty"`
TaskScheduledEventDetails *TaskScheduledEventDetails `json:"taskScheduledEventDetails,omitempty"`
TaskSucceededEventDetails *TaskSucceededEventDetails `json:"taskSucceededEventDetails,omitempty"`
TaskFailedEventDetails *TaskFailedEventDetails `json:"taskFailedEventDetails,omitempty"`
Type string `json:"type"` // e.g. "ExecutionStarted", "ExecutionSucceeded"
Timestamp float64 `json:"timestamp"`
ID int64 `json:"id"`
PreviousEventID int64 `json:"previousEventId"`
}
HistoryEvent represents a single event in execution history.
type InMemoryBackend ¶
type InMemoryBackend struct {
// contains filtered or unexported fields
}
InMemoryBackend implements StorageBackend using in-memory maps.
Phase 3.3: the resource maps (stateMachines, executions, activities, versions, aliases, mapRuns) are *store.Table[T] registered once on registry -- see store_setup.go. Execution history moved from a separate map onto Execution.history (inline, unexported); see the doc comment on Execution. Fields whose key is not a pure, immutable function of the stored value's own fields (or that hold non-serializable runtime state: channels, cancel funcs, timers) remain plain maps -- see store_setup.go's package comment and each field's own comment below for why.
func NewInMemoryBackend ¶
func NewInMemoryBackend() *InMemoryBackend
NewInMemoryBackend creates a new InMemoryBackend with default configuration.
func NewInMemoryBackendWithConfig ¶
func NewInMemoryBackendWithConfig(accountID, region string) *InMemoryBackend
NewInMemoryBackendWithConfig creates a new InMemoryBackend with given account and region. Use NewInMemoryBackendWithContext to bind execution goroutines to a parent context.
func NewInMemoryBackendWithContext ¶
func NewInMemoryBackendWithContext( svcCtx context.Context, accountID, region string, ) *InMemoryBackend
NewInMemoryBackendWithContext creates a new InMemoryBackend whose execution goroutines derive their contexts from svcCtx. When svcCtx is cancelled (e.g. on server shutdown), all running executions are also cancelled. If svcCtx is nil, context.Background is used.
func (*InMemoryBackend) AccountID ¶
func (b *InMemoryBackend) AccountID() string
AccountID returns the AWS account ID this backend is configured for.
func (*InMemoryBackend) CreateActivity ¶
CreateActivity creates a new activity resource in the caller's region.
func (*InMemoryBackend) CreateStateMachine ¶
func (b *InMemoryBackend) CreateStateMachine( ctx context.Context, name, definition, roleArn, smType string, ) (*StateMachine, error)
CreateStateMachine creates and stores a new state machine in the caller's region.
func (*InMemoryBackend) CreateStateMachineAlias ¶
func (b *InMemoryBackend) CreateStateMachineAlias( smARN, name, description string, routing []AliasRoutingConfig, ) (*StateMachineAlias, error)
CreateStateMachineAlias creates a named routing alias for one or more state machine versions.
func (*InMemoryBackend) DeleteActivity ¶
func (b *InMemoryBackend) DeleteActivity(activityArn string) error
DeleteActivity deletes an activity and closes its pending task queue.
func (*InMemoryBackend) DeleteStateMachine ¶
func (b *InMemoryBackend) DeleteStateMachine(arn string) error
DeleteStateMachine marks a state machine as DELETING then removes it.
func (*InMemoryBackend) DeleteStateMachineAlias ¶
func (b *InMemoryBackend) DeleteStateMachineAlias(aliasARN string) error
DeleteStateMachineAlias removes a state machine alias.
func (*InMemoryBackend) DeleteStateMachineVersion ¶
func (b *InMemoryBackend) DeleteStateMachineVersion(versionARN string) error
DeleteStateMachineVersion removes a specific version.
func (*InMemoryBackend) DescribeActivity ¶
func (b *InMemoryBackend) DescribeActivity(activityArn string) (*Activity, error)
DescribeActivity returns activity details.
func (*InMemoryBackend) DescribeExecution ¶
func (b *InMemoryBackend) DescribeExecution(executionArn string) (*Execution, error)
DescribeExecution returns details for a single execution.
func (*InMemoryBackend) DescribeMapRun ¶
func (b *InMemoryBackend) DescribeMapRun(mapRunARN string) (*MapRun, error)
DescribeMapRun returns details for a Map Run.
func (*InMemoryBackend) DescribeStateMachine ¶
func (b *InMemoryBackend) DescribeStateMachine(arn string) (*StateMachine, error)
DescribeStateMachine returns details for a single state machine.
func (*InMemoryBackend) DescribeStateMachineAlias ¶
func (b *InMemoryBackend) DescribeStateMachineAlias(aliasARN string) (*StateMachineAlias, error)
DescribeStateMachineAlias returns details for a state machine alias.
func (*InMemoryBackend) DescribeStateMachineForExecution ¶
func (b *InMemoryBackend) DescribeStateMachineForExecution( executionARN string, ) (*StateMachine, error)
DescribeStateMachineForExecution returns the state machine definition that was active when the given execution was started.
func (*InMemoryBackend) DescribeStateMachineVersion ¶
func (b *InMemoryBackend) DescribeStateMachineVersion( versionARN string, ) (*StateMachineVersion, error)
DescribeStateMachineVersion returns details for a specific version.
func (*InMemoryBackend) Destroy ¶
func (b *InMemoryBackend) Destroy()
Destroy cancels all running execution goroutines and releases resources.
func (*InMemoryBackend) GetActivityTask ¶
func (b *InMemoryBackend) GetActivityTask( ctx context.Context, activityArn, _ string, ) (*ActivityTask, error)
GetActivityTask long-polls for a pending task (up to 60 seconds). Returns an empty ActivityTask (TaskToken="") if no task is available — AWS-compatible behavior.
func (*InMemoryBackend) GetExecutionHistory ¶
func (b *InMemoryBackend) GetExecutionHistory( executionArn, nextToken string, maxResults int, reverseOrder bool, ) ([]HistoryEvent, string, error)
GetExecutionHistory returns history events for an execution.
func (*InMemoryBackend) InvokeActivity ¶
func (b *InMemoryBackend) InvokeActivity( ctx context.Context, activityArn, inputJSON string, heartbeatSeconds int, ) (string, error)
InvokeActivity implements asl.ActivityInvoker. It enqueues a task for the activity and blocks until a worker calls SendTaskSuccess or SendTaskFailure, or the context is cancelled. If heartbeatSeconds > 0, the task fails with ErrHeartbeatTimeout if no SendTaskHeartbeat call arrives within the interval.
func (*InMemoryBackend) ListActivities ¶
func (b *InMemoryBackend) ListActivities( ctx context.Context, nextToken string, maxResults int, ) ([]Activity, string, error)
ListActivities returns activities in the caller's region with optional pagination.
func (*InMemoryBackend) ListExecutions ¶
func (b *InMemoryBackend) ListExecutions( stateMachineArn, statusFilter, nextToken string, maxResults int, ) ([]Execution, string, error)
ListExecutions returns executions for a state machine with optional pagination.
func (*InMemoryBackend) ListMapRuns ¶
func (b *InMemoryBackend) ListMapRuns( executionARN, nextToken string, maxResults int, ) ([]MapRun, string, error)
ListMapRuns returns all MapRuns for an execution.
func (*InMemoryBackend) ListStateMachineAliases ¶
func (b *InMemoryBackend) ListStateMachineAliases( smARN, nextToken string, maxResults int, ) ([]StateMachineAlias, string, error)
ListStateMachineAliases returns all aliases for a state machine.
func (*InMemoryBackend) ListStateMachineVersions ¶
func (b *InMemoryBackend) ListStateMachineVersions( smARN, nextToken string, maxResults int, ) ([]StateMachineVersion, string, error)
ListStateMachineVersions returns all versions for a state machine.
func (*InMemoryBackend) ListStateMachines ¶
func (b *InMemoryBackend) ListStateMachines( ctx context.Context, nextToken string, maxResults int, ) ([]StateMachine, string, error)
ListStateMachines returns state machines in the caller's region with optional pagination.
func (*InMemoryBackend) OnMapRunEnd ¶
func (b *InMemoryBackend) OnMapRunEnd(mapRunARN, status string, succeeded, failed, total int)
OnMapRunEnd implements asl.MapRunNotifier.
func (*InMemoryBackend) OnMapRunStart ¶
func (b *InMemoryBackend) OnMapRunStart( executionARN, stateName string, maxConcurrency, itemCount int, ) string
OnMapRunStart implements asl.MapRunNotifier.
func (*InMemoryBackend) PruneExecutions ¶
func (b *InMemoryBackend) PruneExecutions(_ context.Context) int
PruneExecutions removes executions and history older than the retention period.
func (*InMemoryBackend) PublishStateMachineVersion ¶
func (b *InMemoryBackend) PublishStateMachineVersion( smARN, description, revisionID string, ) (*StateMachineVersion, error)
PublishStateMachineVersion creates an immutable snapshot version of a state machine.
func (*InMemoryBackend) RedriveExecution ¶
func (b *InMemoryBackend) RedriveExecution(executionARN string) (*Execution, error)
RedriveExecution re-runs a FAILED or ABORTED execution starting from its last known state. AWS Step Functions re-runs from the last state that was reached before failure. In this implementation we restart the entire execution with the original input (AWS parity for STANDARD executions).
func (*InMemoryBackend) Region ¶
func (b *InMemoryBackend) Region() string
Region returns the AWS region this backend is configured for.
func (*InMemoryBackend) Reset ¶
func (b *InMemoryBackend) Reset()
Reset clears all in-memory state from the backend. It is used by the POST /_gopherstack/reset endpoint for CI pipelines and rapid local development. Running executions are cancelled before state is cleared.
func (*InMemoryBackend) Restore ¶
func (b *InMemoryBackend) Restore(ctx context.Context, data []byte) error
Restore loads backend state from a JSON snapshot. It implements persistence.Persistable. Service integrations (Lambda, SQS, SNS, DynamoDB) are not restored — they are re-wired by the CLI.
func (*InMemoryBackend) SendTaskFailure ¶
func (b *InMemoryBackend) SendTaskFailure(taskToken, errCode, cause string) error
SendTaskFailure signals failure of an activity task.
func (*InMemoryBackend) SendTaskHeartbeat ¶
func (b *InMemoryBackend) SendTaskHeartbeat(taskToken string) error
SendTaskHeartbeat resets the heartbeat timer for an activity task.
func (*InMemoryBackend) SendTaskSuccess ¶
func (b *InMemoryBackend) SendTaskSuccess(taskToken, output string) error
SendTaskSuccess signals successful completion of an activity task with output.
func (*InMemoryBackend) SetDynamoDBIntegration ¶
func (b *InMemoryBackend) SetDynamoDBIntegration(ddb asl.DynamoDBIntegration)
SetDynamoDBIntegration configures the DynamoDB integration for Task states.
func (*InMemoryBackend) SetECSIntegration ¶
func (b *InMemoryBackend) SetECSIntegration(ecs asl.ECSIntegration)
SetECSIntegration configures the ECS integration.
func (*InMemoryBackend) SetEventBridgeIntegration ¶
func (b *InMemoryBackend) SetEventBridgeIntegration(eb asl.EventBridgeIntegration)
SetEventBridgeIntegration configures the EventBridge integration.
func (*InMemoryBackend) SetGlueIntegration ¶
func (b *InMemoryBackend) SetGlueIntegration(glue asl.GlueIntegration)
SetGlueIntegration configures the Glue integration.
func (*InMemoryBackend) SetLambdaInvoker ¶
func (b *InMemoryBackend) SetLambdaInvoker(invoker asl.LambdaInvoker)
SetLambdaInvoker configures the Lambda invoker for Task states.
func (*InMemoryBackend) SetSNSIntegration ¶
func (b *InMemoryBackend) SetSNSIntegration(sns asl.SNSIntegration)
SetSNSIntegration configures the SNS integration for Task states.
func (*InMemoryBackend) SetSQSIntegration ¶
func (b *InMemoryBackend) SetSQSIntegration(sqs asl.SQSIntegration)
SetSQSIntegration configures the SQS integration for Task states.
func (*InMemoryBackend) SetSettings ¶
func (b *InMemoryBackend) SetSettings(s Settings)
SetSettings updates the backend settings.
func (*InMemoryBackend) SetStateMachineConfigurations ¶
func (b *InMemoryBackend) SetStateMachineConfigurations( arn string, tracing *TracingConfiguration, logging *LoggingConfiguration, encryption *EncryptionConfiguration, ) error
SetStateMachineConfigurations sets optional tracing, logging, and encryption configuration for a state machine. Any nil argument leaves the corresponding field unchanged.
func (*InMemoryBackend) Snapshot ¶
func (b *InMemoryBackend) Snapshot(ctx context.Context) []byte
Snapshot serialises the backend state to JSON. It implements persistence.Persistable. Any execution still RUNNING at snapshot time is promoted to TIMED_OUT so that Restore() never encounters non-terminal executions without a running goroutine.
func (*InMemoryBackend) StartExecution ¶
func (b *InMemoryBackend) StartExecution(stateMachineArn, name, input string) (*Execution, error)
StartExecution creates an execution and runs the ASL interpreter asynchronously.
func (*InMemoryBackend) StartSyncExecution ¶
func (b *InMemoryBackend) StartSyncExecution( stateMachineArn, name, input string, ) (*SyncExecutionResult, error)
StartSyncExecution executes an EXPRESS state machine synchronously and returns the result.
func (*InMemoryBackend) StopExecution ¶
func (b *InMemoryBackend) StopExecution(executionArn, errCode, cause string) error
StopExecution marks a RUNNING execution as ABORTED. AWS behaviour: idempotent on already-terminal executions — returns success without mutation.
func (*InMemoryBackend) SweepTaskTokens ¶
func (b *InMemoryBackend) SweepTaskTokens() int
SweepTaskTokens evicts task tokens that have exceeded the TTL without a worker response, signalling the blocked InvokeActivity goroutine so it is not leaked. Returns evicted count.
func (*InMemoryBackend) UpdateMapRun ¶
func (b *InMemoryBackend) UpdateMapRun( mapRunARN string, maxConcurrency int, toleratedFailureCount int, toleratedFailurePercentage float64, ) (*MapRun, error)
UpdateMapRun updates concurrency/tolerated-failure settings for a Map Run.
func (*InMemoryBackend) UpdateStateMachine ¶
func (b *InMemoryBackend) UpdateStateMachine(smARN, definition, roleArn string) (float64, error)
UpdateStateMachine updates a state machine's definition and/or roleArn. It returns the update timestamp (Unix epoch seconds).
func (*InMemoryBackend) UpdateStateMachineAlias ¶
func (b *InMemoryBackend) UpdateStateMachineAlias( aliasARN, description string, routing []AliasRoutingConfig, ) (*StateMachineAlias, error)
UpdateStateMachineAlias updates an alias's description and/or routing configuration.
func (*InMemoryBackend) WaitForTaskToken ¶
func (b *InMemoryBackend) WaitForTaskToken( ctx context.Context, taskToken string, heartbeatSeconds int, ) (string, error)
WaitForTaskToken registers a callback token and blocks until terminal callback. It returns ErrTaskTokenAlreadyExists when token already exists, ErrHeartbeatTimeout when heartbeatSeconds elapses without heartbeat/success/failure, or ctx.Err() on cancellation.
type Janitor ¶
type Janitor struct {
Backend *InMemoryBackend
Interval time.Duration
TaskTimeout time.Duration
}
Janitor is the Step Functions background worker that manages resource cleanup.
func NewJanitor ¶
func NewJanitor(backend *InMemoryBackend, settings Settings) *Janitor
NewJanitor creates a new Janitor for the given backend.
type LoggingConfiguration ¶
type LoggingConfiguration struct {
Level string `json:"level,omitempty"`
Destinations []LoggingDestination `json:"destinations,omitempty"`
IncludeExecutionData bool `json:"includeExecutionData,omitempty"`
}
LoggingConfiguration controls CloudWatch Logs export for a state machine.
type LoggingDestination ¶
type LoggingDestination struct {
CloudWatchLogsLogGroup *CloudWatchLogsLogGroup `json:"cloudWatchLogsLogGroup,omitempty"`
}
LoggingDestination references a CloudWatch Logs log group destination.
type MapRun ¶
type MapRun struct {
StopDate *float64 `json:"stopDate,omitempty"`
RedriveDate *float64 `json:"redriveDate,omitempty"`
MapRunArn string `json:"mapRunArn"`
ExecutionArn string `json:"executionArn"`
StateMachineArn string `json:"stateMachineArn"`
Status string `json:"status"`
ItemCounts MapRunItemCounts `json:"itemCounts"`
StartDate float64 `json:"startDate"`
ToleratedFailurePercentage float64 `json:"toleratedFailurePercentage,omitempty"`
MaxConcurrency int `json:"maxConcurrency,omitempty"`
ToleratedFailureCount int `json:"toleratedFailureCount,omitempty"`
RedriveCount int `json:"redriveCount,omitempty"`
}
MapRun represents an AWS Step Functions Map Run (a Map state parallel execution group).
type MapRunItemCounts ¶
type MapRunItemCounts struct {
Total int `json:"total"`
Succeeded int `json:"succeeded"`
Failed int `json:"failed"`
Pending int `json:"pending"`
Running int `json:"running"`
Aborted int `json:"aborted"`
TimedOut int `json:"timedOut"`
ResultsWritten int `json:"resultsWritten"`
FailuresNotRedrivable int `json:"failuresNotRedrivable,omitempty"`
PendingRedrive int `json:"pendingRedrive,omitempty"`
}
MapRunItemCounts holds item-level counts for a Map Run.
type Provider ¶
type Provider struct{}
Provider implements service.Provider for the Step Functions service.
func (*Provider) Init ¶
func (p *Provider) Init(ctx *service.AppContext) (service.Registerable, error)
Init initializes the Step Functions service backend and handler.
type Settings ¶
type Settings struct {
// ExecutionRetention is how long execution history is kept before being pruned.
// Defaults to 24 hours for local mock stability.
ExecutionRetention time.Duration `` //nolint:lll // tags
/* 142-byte string literal not displayed */
// JanitorInterval is how often the background worker runs.
JanitorInterval time.Duration `` //nolint:lll // tags
/* 133-byte string literal not displayed */
// TaskTokenTTL is the maximum time a task token may remain in tasksByToken without
// a worker response before the janitor evicts it. Defaults to 1 hour.
TaskTokenTTL time.Duration `` //nolint:lll // tags
/* 146-byte string literal not displayed */
}
Settings holds configurable settings for the Step Functions service.
func DefaultSettings ¶
func DefaultSettings() Settings
DefaultSettings returns Settings with sensible defaults.
type SettingsProvider ¶
type SettingsProvider interface {
GetStepFunctionsSettings() Settings
}
SettingsProvider is implemented by config objects that supply Step Functions settings.
type StateEnteredEventDetails ¶
type StateEnteredEventDetails struct {
Name string `json:"name"`
Input string `json:"input,omitempty"`
}
StateEnteredEventDetails holds details for state-entered events.
type StateExitedEventDetails ¶
type StateExitedEventDetails struct {
Name string `json:"name"`
Output string `json:"output,omitempty"`
}
StateExitedEventDetails holds details for state-exited events.
type StateMachine ¶
type StateMachine struct {
EncryptionConfiguration *EncryptionConfiguration `json:"encryptionConfiguration,omitempty"`
TracingConfiguration *TracingConfiguration `json:"tracingConfiguration,omitempty"`
LoggingConfiguration *LoggingConfiguration `json:"loggingConfiguration,omitempty"`
Name string `json:"name"`
StateMachineArn string `json:"stateMachineArn"`
Type string `json:"type"`
Status string `json:"status"`
Definition string `json:"definition"`
RoleArn string `json:"roleArn"`
CreationDate float64 `json:"creationDate"`
UpdatedDate float64 `json:"updatedDate,omitempty"`
}
StateMachine represents a Step Functions state machine. Timestamp fields use float64 (Unix epoch seconds) as required by the AWS JSON 1.0 protocol used by Step Functions.
type StateMachineAlias ¶
type StateMachineAlias struct {
StateMachineAliasArn string `json:"stateMachineAliasArn"`
Name string `json:"name"`
Description string `json:"description,omitempty"`
RoutingConfiguration []AliasRoutingConfig `json:"routingConfiguration"`
CreationDate float64 `json:"creationDate"`
UpdatedDate float64 `json:"updatedDate,omitempty"`
}
StateMachineAlias represents a routing alias for one or more state machine versions.
type StateMachineVersion ¶
type StateMachineVersion struct {
StateMachineVersionArn string `json:"stateMachineVersionArn"`
StateMachineArn string `json:"stateMachineArn"`
Name string `json:"name"`
Definition string `json:"definition"`
RoleArn string `json:"roleArn"`
Type string `json:"type"`
Status string `json:"status"`
Description string `json:"description,omitempty"`
RevisionID string `json:"revisionId,omitempty"`
CreationDate float64 `json:"creationDate"`
}
StateMachineVersion represents an immutable versioned snapshot of a state machine.
type StorageBackend ¶
type StorageBackend interface {
CreateStateMachine(
ctx context.Context,
name, definition, roleArn, smType string,
) (*StateMachine, error)
DeleteStateMachine(arn string) error
ListStateMachines(
ctx context.Context,
nextToken string,
maxResults int,
) ([]StateMachine, string, error)
DescribeStateMachine(arn string) (*StateMachine, error)
UpdateStateMachine(arn, definition, roleArn string) (float64, error)
PublishStateMachineVersion(smARN, description, revisionID string) (*StateMachineVersion, error)
DescribeStateMachineVersion(versionARN string) (*StateMachineVersion, error)
DeleteStateMachineVersion(versionARN string) error
ListStateMachineVersions(
smARN, nextToken string,
maxResults int,
) ([]StateMachineVersion, string, error)
CreateStateMachineAlias(
smARN, name, description string,
routing []AliasRoutingConfig,
) (*StateMachineAlias, error)
UpdateStateMachineAlias(
aliasARN, description string,
routing []AliasRoutingConfig,
) (*StateMachineAlias, error)
DeleteStateMachineAlias(aliasARN string) error
DescribeStateMachineAlias(aliasARN string) (*StateMachineAlias, error)
ListStateMachineAliases(
smARN, nextToken string,
maxResults int,
) ([]StateMachineAlias, string, error)
StartExecution(stateMachineArn, name, input string) (*Execution, error)
StartSyncExecution(stateMachineArn, name, input string) (*SyncExecutionResult, error)
StopExecution(executionArn, errCode, cause string) error
RedriveExecution(executionARN string) (*Execution, error)
DescribeExecution(executionArn string) (*Execution, error)
DescribeStateMachineForExecution(executionARN string) (*StateMachine, error)
ListExecutions(
stateMachineArn, statusFilter, nextToken string,
maxResults int,
) ([]Execution, string, error)
GetExecutionHistory(
executionArn, nextToken string,
maxResults int,
reverseOrder bool,
) ([]HistoryEvent, string, error)
CreateActivity(ctx context.Context, name string) (*Activity, error)
DeleteActivity(activityArn string) error
DescribeActivity(activityArn string) (*Activity, error)
ListActivities(
ctx context.Context,
nextToken string,
maxResults int,
) ([]Activity, string, error)
GetActivityTask(ctx context.Context, activityArn, workerName string) (*ActivityTask, error)
SendTaskSuccess(taskToken, output string) error
SendTaskFailure(taskToken, errCode, cause string) error
SendTaskHeartbeat(taskToken string) error
SetStateMachineConfigurations(
arn string,
tracing *TracingConfiguration,
logging *LoggingConfiguration,
encryption *EncryptionConfiguration,
) error
DescribeMapRun(mapRunARN string) (*MapRun, error)
UpdateMapRun(
mapRunARN string,
maxConcurrency int,
toleratedFailureCount int,
toleratedFailurePercentage float64,
) (*MapRun, error)
ListMapRuns(executionARN, nextToken string, maxResults int) ([]MapRun, string, error)
}
StorageBackend is the interface for a Step Functions in-memory store.
type SyncExecutionResult ¶
type SyncExecutionResult struct {
ExecutionArn string `json:"executionArn"`
StateMachineArn string `json:"stateMachineArn"`
Name string `json:"name"`
Status string `json:"status"`
Input string `json:"input,omitempty"`
Output string `json:"output,omitempty"`
Error string `json:"error,omitempty"`
Cause string `json:"cause,omitempty"`
StopDate float64 `json:"stopDate"`
StartDate float64 `json:"startDate"`
}
SyncExecutionResult holds the result of a synchronous (EXPRESS) execution.
type TaskFailedEventDetails ¶
type TaskFailedEventDetails struct {
Error string `json:"error,omitempty"`
Cause string `json:"cause,omitempty"`
Resource string `json:"resource,omitempty"`
}
TaskFailedEventDetails holds details for TaskFailed history events.
type TaskScheduledEventDetails ¶
type TaskScheduledEventDetails struct {
Resource string `json:"resource,omitempty"`
}
TaskScheduledEventDetails holds details for TaskScheduled history events.
type TaskSucceededEventDetails ¶
type TaskSucceededEventDetails struct {
Resource string `json:"resource,omitempty"`
Output string `json:"output,omitempty"`
}
TaskSucceededEventDetails holds details for TaskSucceeded history events.
type TracingConfiguration ¶
type TracingConfiguration struct {
Enabled bool `json:"enabled"`
}
TracingConfiguration controls AWS X-Ray tracing for a state machine.
Source Files
¶
- activities.go
- aliases.go
- eb_integration.go
- errors.go
- execution_history.go
- executions.go
- handler.go
- handler_activities.go
- handler_aliases.go
- handler_executions.go
- handler_map_runs.go
- handler_state_machine_versions.go
- handler_state_machines.go
- handler_tags.go
- handler_util.go
- integrations.go
- interfaces.go
- janitor.go
- map_runs.go
- models.go
- persistence.go
- provider.go
- settings.go
- state_machine_versions.go
- state_machines.go
- store.go
- store_setup.go
Directories
¶
| Path | Synopsis |
|---|---|
|
Package asl implements an interpreter for the Amazon States Language (ASL) used by AWS Step Functions to define state machine workflows.
|
Package asl implements an interpreter for the Amazon States Language (ASL) used by AWS Step Functions to define state machine workflows. |