workflow

package
v1.19.0-rc.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 1, 2026 License: Apache-2.0 Imports: 29 Imported by: 0

Documentation

Index

Constants

View Source
const ClusteredDeploymentConfig = `` /* 180-byte string literal not displayed */

ClusteredDeploymentConfig is a Configuration manifest enabling the WorkflowsClusteredDeployment preview feature.

Variables

This section is empty.

Functions

func ClusteredDeploymentFromEnv

func ClusteredDeploymentFromEnv() bool

ClusteredDeploymentFromEnv reports whether the suite is running with DAPR_INTEGRATION_WORKFLOW_CLUSTERED set truthy, which enables the WorkflowsClusteredDeployment feature flag on every daprd built by this harness unless a test overrides it with WithClusteredDeployment.

func FastPathFromEnv

func FastPathFromEnv() bool

FastPathFromEnv reports whether the suite is running with DAPR_INTEGRATION_WORKFLOW_FASTPATH set truthy, which enables the WorkflowsFastPath feature flag on every daprd built by this harness unless a test overrides it with WithFastPath.

func SchedulerPlacementFromEnv

func SchedulerPlacementFromEnv() bool

SchedulerPlacementFromEnv reports whether DAPR_INTEGRATION_SCHEDULER_PLACEMENT is set truthy, which has the scheduler serve actor placement for every harness built by this package. Tests which pick a topology, or drive the placement service, keep their choice.

func SigningFromEnv

func SigningFromEnv() bool

SigningFromEnv reports whether the suite is running with DAPR_INTEGRATION_WORKFLOW_SIGNING set truthy. WorkflowHistorySigning requires mTLS (daprd fatals otherwise), so signing mode forces a Sentry per workflow, exactly as WithMTLS does, unless a test overrides it with WithSigning.

Types

type ConnectedWorker

type ConnectedWorker struct {
	Client   *client.TaskHubGrpcClient
	Observer *WorkItemObserver
	// contains filtered or unexported fields
}

ConnectedWorker is a backend worker connected to a daprd through a dedicated gRPC connection whose lifetime the caller controls. Its Observer records the work items the sidecar streams to it. Use Disconnect to drop the worker's stream (e.g. to model a crash or reconnect); a fresh ConnectWorkerN call reconnects as a new, cold stream.

func (*ConnectedWorker) Disconnect

func (c *ConnectedWorker) Disconnect(t *testing.T)

Disconnect tears down the worker's stream and connection. It is idempotent.

type Option

type Option func(*options)

func WithAddActivity added in v1.16.0

func WithAddActivity(t *testing.T, name string, a func(task.ActivityContext) (any, error)) Option

func WithAddActivityN

func WithAddActivityN(t *testing.T, index int, name string, a func(task.ActivityContext) (any, error)) Option

func WithAddOrchestrator added in v1.16.0

func WithAddOrchestrator(t *testing.T, name string, or func(*task.WorkflowContext) (any, error)) Option

func WithAddWorkflowN added in v1.18.0

func WithAddWorkflowN(t *testing.T, index int, name string, or func(*task.WorkflowContext) (any, error)) Option

func WithClusteredDeployment

func WithClusteredDeployment(enabled bool) Option

func WithDaprdOptions added in v1.16.0

func WithDaprdOptions(index int, opts ...daprd.Option) Option

func WithDaprds added in v1.15.6

func WithDaprds(daprds int) Option

func WithFastPath

func WithFastPath(enabled bool) Option

WithFastPath explicitly enables or disables the WorkflowsFastPath feature flag on every daprd in the workflow, overriding the DAPR_INTEGRATION_WORKFLOW_FASTPATH environment variable.

func WithHistorySigning added in v1.18.4

func WithHistorySigning(t *testing.T) Option

WithHistorySigning enables the WorkflowHistorySigning feature flag on every daprd in the workflow. History signing needs the Sentry-issued workload identity for its attestation and signing keys, so this also enables the mTLS setup of WithMTLS. Prefer this over WithMTLS in tests that are about signing behavior, so the intent is explicit at the call site.

func WithMTLS added in v1.18.0

func WithMTLS(t *testing.T) Option

WithMTLS spins up a Sentry process for mTLS and enables the WorkflowHistorySigning feature flag on every daprd in the workflow.

func WithNoDB added in v1.17.0

func WithNoDB() Option

func WithPlacementOptions added in v1.18.0

func WithPlacementOptions(opts ...placement.Option) Option

func WithPlacementService

func WithPlacementService() Option

WithPlacementService pins actor placement to the placement service, overriding DAPR_INTEGRATION_SCHEDULER_PLACEMENT.

func WithSchedulerAddress added in v1.17.7

func WithSchedulerAddress(addr string) Option

WithSchedulerAddress overrides the address used for the daprd's --scheduler-host-address flag. Use this to point daprd at a proxy that fronts the real scheduler.

func WithSchedulerInstance added in v1.17.7

func WithSchedulerInstance(sched *scheduler.Scheduler) Option

WithSchedulerInstance lets a test supply a pre-constructed scheduler. The framework uses this scheduler instead of creating its own and skips adding it to its process list (the caller is responsible for that). Combine with WithSchedulerAddress when interposing a proxy.

func WithSchedulerOptions added in v1.18.0

func WithSchedulerOptions(opts ...scheduler.Option) Option

func WithSchedulerPlacement

func WithSchedulerPlacement() Option

WithClusteredDeployment explicitly enables or disables the WorkflowsClusteredDeployment feature flag on every daprd in the workflow, overriding the DAPR_INTEGRATION_WORKFLOW_CLUSTERED environment variable. WithSchedulerPlacement serves actor placement from the scheduler: no placement process runs.

func WithSentryInstance

func WithSentryInstance(sen *sentry.Sentry) Option

WithSentryInstance lets a test supply a pre-constructed Sentry, implying mTLS. The framework uses it for placement, scheduler and daprd identity instead of creating its own, and skips adding it to its process list (the caller is responsible for that). Required when combining WithSchedulerInstance with mTLS/signing so the caller-built scheduler and proxy share the same trust chain.

func WithSigning

func WithSigning(enabled bool) Option

WithSigning explicitly enables or disables workflow history signing mode (mTLS with a Sentry plus the WorkflowHistorySigning feature flag), overriding the DAPR_INTEGRATION_WORKFLOW_SIGNING environment variable. WithSigning(false) disables the feature flag only: mTLS requested via WithMTLS or WithSentryInstance stays on.

func WithSigningDisabledN added in v1.18.0

func WithSigningDisabledN(index int) Option

WithSigningDisabledN excludes the daprd at the given index from having the WorkflowHistorySigning feature flag set. Has no effect without WithMTLS or WithHistorySigning.

type WorkItemObserver

type WorkItemObserver struct {
	// contains filtered or unexported fields
}

WorkItemObserver records the workflow work items a worker receives over its GetWorkItems stream, classified by whether the sidecar sent the committed history in full or as a stateful-history delta (a WorkflowRequest carrying a CachedHistory message). It is installed as a gRPC stream interceptor by ConnectWorkerN and is safe for concurrent use.

func (*WorkItemObserver) Deltas

func (o *WorkItemObserver) Deltas() int

Deltas returns the total number of delta (cached-history) work items received across all instances.

func (*WorkItemObserver) DeltasFor

func (o *WorkItemObserver) DeltasFor(instanceID string) int

DeltasFor returns the number of delta work items received for a single instance.

func (*WorkItemObserver) FullSends

func (o *WorkItemObserver) FullSends() int

FullSends returns the total number of full-history work items received across all instances.

func (*WorkItemObserver) FullSendsFor

func (o *WorkItemObserver) FullSendsFor(instanceID string) int

FullSendsFor returns the number of full-history work items received for a single instance.

func (*WorkItemObserver) GetInstanceHistoryCalls

func (o *WorkItemObserver) GetInstanceHistoryCalls() int

GetInstanceHistoryCalls returns how many times the worker fell back to the GetInstanceHistory RPC, i.e. how many stateful-history cache misses it recovered from. In steady state with a warm cache this stays zero.

func (*WorkItemObserver) HealthPings

func (o *WorkItemObserver) HealthPings() int

HealthPings returns the number of HealthPing work items received.

func (*WorkItemObserver) WorkItemStreams

func (o *WorkItemObserver) WorkItemStreams() int

WorkItemStreams returns how many GetWorkItems streams the worker opened, so one more than the number of times it reconnected.

type Workflow

type Workflow struct {
	// contains filtered or unexported fields
}

func New

func New(t *testing.T, fopts ...Option) *Workflow

func NewClustered added in v1.17.13

func NewClustered(t *testing.T, daprds int, extraDaprdOpts ...daprd.Option) *Workflow

NewClustered returns a Workflow whose daprds share a single app ID with WorkflowsClusteredDeployment enabled, representing a clustered deployment behind a load balancer.

func (*Workflow) ActivityActorType added in v1.18.4

func (w *Workflow) ActivityActorType(index int) string

ActivityActorType returns the activity actor type registered by the daprd at the given index (default namespace).

func (*Workflow) ActorTypesCount

func (w *Workflow) ActorTypesCount() int

ActorTypesCount returns the number of actor types a daprd in this workflow registers when a worker is connected: workflow, activity and retentioner, plus the executor rendezvous type in clustered deployment mode.

func (*Workflow) BackendClient

func (w *Workflow) BackendClient(t *testing.T, ctx context.Context) *client.TaskHubGrpcClient

func (*Workflow) BackendClientN added in v1.16.0

func (w *Workflow) BackendClientN(t *testing.T, ctx context.Context, index int) *client.TaskHubGrpcClient

BackendClient returns a backend client for the specified index

func (*Workflow) Cleanup

func (w *Workflow) Cleanup(t *testing.T)

func (*Workflow) ClusteredDeployment

func (w *Workflow) ClusteredDeployment() bool

ClusteredDeployment reports whether every daprd in this workflow runs with the WorkflowsClusteredDeployment feature flag enabled. Tests use this to branch assertions which differ between the two modes.

func (*Workflow) ConnectWorker

ConnectWorker connects a worker to daprd index 0. See ConnectWorkerN.

func (*Workflow) ConnectWorkerN

func (w *Workflow) ConnectWorkerN(t *testing.T, ctx context.Context, index int, r *task.TaskRegistry, opts ...client.TaskHubGrpcClientOption) *ConnectedWorker

ConnectWorkerN connects a backend worker to the daprd at the given index, serving the supplied task registry, through a gRPC stream interceptor that observes every work item it receives. Unlike BackendClientN it does not wait for the worker to register (use WaitForConnectedWorkersN), so tests can model connect/disconnect timing precisely. The worker is torn down on test cleanup.

func (*Workflow) DB added in v1.15.0

func (w *Workflow) DB() *sqlite.SQLite

func (*Workflow) Dapr added in v1.15.0

func (w *Workflow) Dapr() *daprd.Daprd

func (*Workflow) DaprN added in v1.15.6

func (w *Workflow) DaprN(i int) *daprd.Daprd

func (*Workflow) FastPath

func (w *Workflow) FastPath() bool

FastPath reports whether every daprd in this workflow runs with the WorkflowsFastPath feature flag enabled. Tests use this to branch assertions which differ between the two modes.

func (*Workflow) GRPCClient

func (w *Workflow) GRPCClient(t *testing.T, ctx context.Context) rtv1.DaprClient

func (*Workflow) GRPCClientN added in v1.16.0

func (w *Workflow) GRPCClientN(t *testing.T, ctx context.Context, index int) rtv1.DaprClient

GRPCClientForApp returns a GRPC client for the specified app index

func (*Workflow) HasPlacement

func (w *Workflow) HasPlacement() bool

HasPlacement reports whether a standalone placement service runs, rather than placement served by the scheduler.

func (*Workflow) JoinOptions

func (w *Workflow) JoinOptions(t *testing.T) []daprd.Option

JoinOptions returns the options an extra daprd needs to join this harness's cluster in the same modes: the feature flags and, under mTLS, the Sentry wiring.

func (*Workflow) ManagementClient added in v1.18.3

func (w *Workflow) ManagementClient(t *testing.T, ctx context.Context) *client.TaskHubGrpcClient

ManagementClient returns a backend client for daprd index 0 that is not a worker. See ManagementClientN.

func (*Workflow) ManagementClientN added in v1.18.3

func (w *Workflow) ManagementClientN(t *testing.T, ctx context.Context, index int) *client.TaskHubGrpcClient

ManagementClientN returns a backend client connected to the daprd at the given index for control-plane operations only (scheduling, waiting, raising events, fetching history). It does not start a work-item listener, so it never executes workflows and never advertises any worker capability. Its connection is stable for the lifetime of the test, independent of any ConnectedWorker, so it can drive instances while workers connect and disconnect around it.

func (*Workflow) Metrics

func (w *Workflow) Metrics(t *testing.T, ctx context.Context) map[string]float64

func (*Workflow) Placement added in v1.18.0

func (w *Workflow) Placement() *placement.Placement

func (*Workflow) PlacementVersion

func (w *Workflow) PlacementVersion(t *testing.T, ctx context.Context) uint64

PlacementVersion returns a counter which advances whenever the active placement authority completes a dissemination: the default namespace table version of the placement service, or the scheduler's dissemination count. Only successive values compare, the units differ per authority.

func (*Workflow) Registry added in v1.15.0

func (w *Workflow) Registry() *task.TaskRegistry

func (*Workflow) RegistryN added in v1.16.0

func (w *Workflow) RegistryN(index int) *task.TaskRegistry

Registry returns the registry for a specific index

func (*Workflow) ResetRegistry added in v1.17.0

func (w *Workflow) ResetRegistry(t *testing.T)

func (*Workflow) Run

func (w *Workflow) Run(t *testing.T, ctx context.Context)

func (*Workflow) Scheduler added in v1.17.0

func (w *Workflow) Scheduler() *scheduler.Scheduler

func (*Workflow) SchedulerPlacement

func (w *Workflow) SchedulerPlacement() bool

SchedulerPlacement reports whether the scheduler serves actor placement for this harness, in which case Placement returns nil.

func (*Workflow) Sentry added in v1.18.0

func (w *Workflow) Sentry() *sentry.Sentry

func (*Workflow) Signing

func (w *Workflow) Signing() bool

Signing reports whether this workflow runs with mTLS active and the WorkflowHistorySigning feature flag enabled on its daprds (except any excluded via WithSigningDisabledN). Tests use this to branch assertions which differ when history signing is active.

func (*Workflow) StrayFire added in v1.18.4

func (w *Workflow) StrayFire(t *testing.T, ctx context.Context, index int, instanceID string, mtls bool)

StrayFire schedules a stray new-event reminder against the workflow actor hosted by daprd index, driving its empty-inbox path, and waits until the scheduler has delivered it.

func (*Workflow) WaitForConnectedWorkers

func (w *Workflow) WaitForConnectedWorkers(t *testing.T, ctx context.Context, count int)

WaitForConnectedWorkers waits for at least count workers to be registered with daprd index 0's actor runtime. See WaitForConnectedWorkersN.

func (*Workflow) WaitForConnectedWorkersN

func (w *Workflow) WaitForConnectedWorkersN(t *testing.T, ctx context.Context, index, count int)

WaitForConnectedWorkersN waits for at least count workers to be registered with the actor runtime of the daprd at the given index, and for the workflow actor runtime itself to be ready.

func (*Workflow) WaitForNoConnectedWorkers added in v1.17.13

func (w *Workflow) WaitForNoConnectedWorkers(t *testing.T, ctx context.Context)

func (*Workflow) WaitForNoConnectedWorkersN added in v1.17.13

func (w *Workflow) WaitForNoConnectedWorkersN(t *testing.T, ctx context.Context, index int)

func (*Workflow) WaitUntilRunning

func (w *Workflow) WaitUntilRunning(t *testing.T, ctx context.Context)

func (*Workflow) WorkflowActorType added in v1.18.4

func (w *Workflow) WorkflowActorType(index int) string

WorkflowActorType returns the orchestrator actor type registered by the daprd at the given index (default namespace).

func (*Workflow) WorkflowClient added in v1.17.0

func (w *Workflow) WorkflowClient(t *testing.T, ctx context.Context) *workflow.Client

func (*Workflow) WorkflowClientN added in v1.17.0

func (w *Workflow) WorkflowClientN(t *testing.T, ctx context.Context, index int) *workflow.Client

func (*Workflow) WriteWorkflowState added in v1.18.4

func (w *Workflow) WriteWorkflowState(t *testing.T, ctx context.Context, index int, instanceID string, generation uint64, history, inbox []*protos.HistoryEvent)

WriteWorkflowState writes a fabricated durable workflow state for the given instance straight into the SQLite actor state store: the history and inbox event rows plus the metadata row describing them, in the exact key layout the workflow state loader reads. Existing history and inbox rows for the instance are deleted first so the metadata lengths stay authoritative. daprd in-memory caches are not touched; pair this with a scheduler-driven reminder or a fresh actor activation to make daprd observe the rows.

Directories

Path Synopsis

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL