Documentation
¶
Index ¶
- Constants
- func ClusteredDeploymentFromEnv() bool
- func FastPathFromEnv() bool
- func SchedulerPlacementFromEnv() bool
- func SigningFromEnv() bool
- type ConnectedWorker
- type Option
- func WithAddActivity(t *testing.T, name string, a func(task.ActivityContext) (any, error)) Option
- func WithAddActivityN(t *testing.T, index int, name string, ...) Option
- func WithAddOrchestrator(t *testing.T, name string, or func(*task.WorkflowContext) (any, error)) Option
- func WithAddWorkflowN(t *testing.T, index int, name string, ...) Option
- func WithClusteredDeployment(enabled bool) Option
- func WithDaprdOptions(index int, opts ...daprd.Option) Option
- func WithDaprds(daprds int) Option
- func WithFastPath(enabled bool) Option
- func WithHistorySigning(t *testing.T) Option
- func WithMTLS(t *testing.T) Option
- func WithNoDB() Option
- func WithPlacementOptions(opts ...placement.Option) Option
- func WithPlacementService() Option
- func WithSchedulerAddress(addr string) Option
- func WithSchedulerInstance(sched *scheduler.Scheduler) Option
- func WithSchedulerOptions(opts ...scheduler.Option) Option
- func WithSchedulerPlacement() Option
- func WithSentryInstance(sen *sentry.Sentry) Option
- func WithSigning(enabled bool) Option
- func WithSigningDisabledN(index int) Option
- type WorkItemObserver
- func (o *WorkItemObserver) Deltas() int
- func (o *WorkItemObserver) DeltasFor(instanceID string) int
- func (o *WorkItemObserver) FullSends() int
- func (o *WorkItemObserver) FullSendsFor(instanceID string) int
- func (o *WorkItemObserver) GetInstanceHistoryCalls() int
- func (o *WorkItemObserver) HealthPings() int
- func (o *WorkItemObserver) WorkItemStreams() int
- type Workflow
- func (w *Workflow) ActivityActorType(index int) string
- func (w *Workflow) ActorTypesCount() int
- func (w *Workflow) BackendClient(t *testing.T, ctx context.Context) *client.TaskHubGrpcClient
- func (w *Workflow) BackendClientN(t *testing.T, ctx context.Context, index int) *client.TaskHubGrpcClient
- func (w *Workflow) Cleanup(t *testing.T)
- func (w *Workflow) ClusteredDeployment() bool
- func (w *Workflow) ConnectWorker(t *testing.T, ctx context.Context, r *task.TaskRegistry, ...) *ConnectedWorker
- func (w *Workflow) ConnectWorkerN(t *testing.T, ctx context.Context, index int, r *task.TaskRegistry, ...) *ConnectedWorker
- func (w *Workflow) DB() *sqlite.SQLite
- func (w *Workflow) Dapr() *daprd.Daprd
- func (w *Workflow) DaprN(i int) *daprd.Daprd
- func (w *Workflow) FastPath() bool
- func (w *Workflow) GRPCClient(t *testing.T, ctx context.Context) rtv1.DaprClient
- func (w *Workflow) GRPCClientN(t *testing.T, ctx context.Context, index int) rtv1.DaprClient
- func (w *Workflow) HasPlacement() bool
- func (w *Workflow) JoinOptions(t *testing.T) []daprd.Option
- func (w *Workflow) ManagementClient(t *testing.T, ctx context.Context) *client.TaskHubGrpcClient
- func (w *Workflow) ManagementClientN(t *testing.T, ctx context.Context, index int) *client.TaskHubGrpcClient
- func (w *Workflow) Metrics(t *testing.T, ctx context.Context) map[string]float64
- func (w *Workflow) Placement() *placement.Placement
- func (w *Workflow) PlacementVersion(t *testing.T, ctx context.Context) uint64
- func (w *Workflow) Registry() *task.TaskRegistry
- func (w *Workflow) RegistryN(index int) *task.TaskRegistry
- func (w *Workflow) ResetRegistry(t *testing.T)
- func (w *Workflow) Run(t *testing.T, ctx context.Context)
- func (w *Workflow) Scheduler() *scheduler.Scheduler
- func (w *Workflow) SchedulerPlacement() bool
- func (w *Workflow) Sentry() *sentry.Sentry
- func (w *Workflow) Signing() bool
- func (w *Workflow) StrayFire(t *testing.T, ctx context.Context, index int, instanceID string, mtls bool)
- func (w *Workflow) WaitForConnectedWorkers(t *testing.T, ctx context.Context, count int)
- func (w *Workflow) WaitForConnectedWorkersN(t *testing.T, ctx context.Context, index, count int)
- func (w *Workflow) WaitForNoConnectedWorkers(t *testing.T, ctx context.Context)
- func (w *Workflow) WaitForNoConnectedWorkersN(t *testing.T, ctx context.Context, index int)
- func (w *Workflow) WaitUntilRunning(t *testing.T, ctx context.Context)
- func (w *Workflow) WorkflowActorType(index int) string
- func (w *Workflow) WorkflowClient(t *testing.T, ctx context.Context) *workflow.Client
- func (w *Workflow) WorkflowClientN(t *testing.T, ctx context.Context, index int) *workflow.Client
- func (w *Workflow) WriteWorkflowState(t *testing.T, ctx context.Context, index int, instanceID string, ...)
Constants ¶
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 WithAddActivityN ¶
func WithAddOrchestrator ¶ added in v1.16.0
func WithAddWorkflowN ¶ added in v1.18.0
func WithClusteredDeployment ¶
func WithDaprdOptions ¶ added in v1.16.0
func WithDaprds ¶ added in v1.15.6
func WithFastPath ¶
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
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
WithMTLS spins up a Sentry process for mTLS and enables the WorkflowHistorySigning feature flag on every daprd in the workflow.
func WithPlacementOptions ¶ added in v1.18.0
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
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
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 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 ¶
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 ¶
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
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 NewClustered ¶ added in v1.17.13
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
ActivityActorType returns the activity actor type registered by the daprd at the given index (default namespace).
func (*Workflow) ActorTypesCount ¶
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 (*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) ClusteredDeployment ¶
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 ¶
func (w *Workflow) ConnectWorker(t *testing.T, ctx context.Context, r *task.TaskRegistry, opts ...client.TaskHubGrpcClientOption) *ConnectedWorker
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) FastPath ¶
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 (*Workflow) GRPCClientN ¶ added in v1.16.0
GRPCClientForApp returns a GRPC client for the specified app index
func (*Workflow) HasPlacement ¶
HasPlacement reports whether a standalone placement service runs, rather than placement served by the scheduler.
func (*Workflow) JoinOptions ¶
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
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) PlacementVersion ¶
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 (*Workflow) SchedulerPlacement ¶
SchedulerPlacement reports whether the scheduler serves actor placement for this harness, in which case Placement returns nil.
func (*Workflow) Signing ¶
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 ¶
WaitForConnectedWorkers waits for at least count workers to be registered with daprd index 0's actor runtime. See WaitForConnectedWorkersN.
func (*Workflow) WaitForConnectedWorkersN ¶
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 (*Workflow) WaitForNoConnectedWorkersN ¶ added in v1.17.13
func (*Workflow) WaitUntilRunning ¶
func (*Workflow) WorkflowActorType ¶ added in v1.18.4
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 (*Workflow) WorkflowClientN ¶ added in v1.17.0
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.