proxy

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: 32 Imported by: 0

Documentation

Index

Constants

View Source
const (
	MethodScheduleJob        = "ScheduleJob"
	MethodDeleteJob          = "DeleteJob"
	MethodGetJob             = "GetJob"
	MethodListJobs           = "ListJobs"
	MethodDeleteByMetadata   = "DeleteByMetadata"
	MethodDeleteByNamePrefix = "DeleteByNamePrefix"
	MethodWatchHosts         = "WatchHosts"
	MethodWatchJobs          = "WatchJobs"
	MethodReportActorTypes   = "ReportActorTypes"
)

Method names used with ArmFailures. Using typed constants keeps tests from drifting if a method is renamed.

Variables

This section is empty.

Functions

This section is empty.

Types

type Option

type Option func(*Proxy)

Option configures the proxy.

func WithSentry

func WithSentry(t *testing.T, sen *sentry.Sentry, namespace, appID string) Option

WithSentry makes the proxy a full mTLS member of the control plane. It serves daprd with a leaf certificate minted from the sentry CA carrying the scheduler control-plane SPIFFE identity, and dials the upstream scheduler presenting the identity of the daprd app whose requests it forwards (the scheduler authorizes each request against the caller's SPIFFE identity, so the proxy must impersonate the app, which therefore must run with the given fixed app ID). The upstream scheduler must be built with scheduler.WithSentry and an ID matching its certificate DNS names.

type Proxy

type Proxy struct {
	schedulerv1pb.UnimplementedSchedulerServer
	// contains filtered or unexported fields
}

func New

func New(t *testing.T, sched *scheduler.Scheduler, fopts ...Option) *Proxy

New returns a proxy that wraps the given scheduler. Run blocks until the upstream is reachable, so the scheduler must precede the proxy in framework process ordering. daprd should be configured with daprd.WithSchedulerAddresses(proxy.Address()) instead of pointing at the scheduler directly.

func (*Proxy) Address

func (p *Proxy) Address() string

func (*Proxy) ArmFailures

func (p *Proxy) ArmFailures(method string, n int, code codes.Code, notify chan struct{})

ArmFailures arms the proxy to fail the next n requests to the given method with the supplied gRPC status code. n=0 disarms. If notify is non-nil it is closed the first time a matching call is failed.

Choose the code with care: the daprd-side CreateReminderWithRetry transparently retries codes.Unavailable / codes.DeadlineExceeded with exponential backoff (up to a minute), so injecting those codes only stalls the call instead of surfacing the failure. Use codes.Internal / codes.Aborted or another non-transient code when the test wants the failure to propagate to the orchestrator's error path.

func (*Proxy) ArmNamedFailures added in v1.18.4

func (p *Proxy) ArmNamedFailures(method, nameContains string, n int, code codes.Code, notify chan struct{})

ArmNamedFailures is ArmFailures restricted to requests whose job name contains nameContains; requests for other jobs pass through untouched.

func (*Proxy) Cleanup

func (p *Proxy) Cleanup(t *testing.T)

func (*Proxy) FailedCount

func (p *Proxy) FailedCount() int

func (*Proxy) GetJob

func (*Proxy) Partition

func (p *Proxy) Partition(t *testing.T)

Partition blackholes all daprd connections to this proxy and closes the upstream, so the scheduler evicts the hosts. It cannot be undone.

func (*Proxy) Port

func (p *Proxy) Port() int

func (*Proxy) ReportActorTypes

func (p *Proxy) ReportActorTypes(stream schedulerv1pb.Scheduler_ReportActorTypesServer) error

func (*Proxy) Run

func (p *Proxy) Run(t *testing.T, ctx context.Context)

func (*Proxy) WatchHosts

WatchHosts forwards host updates from the upstream scheduler, rewriting every Host.Address so daprd reconnects to the proxy rather than to the real scheduler on every refresh. Without this rewrite daprd would bypass the proxy as soon as the first host list arrived.

func (*Proxy) WatchJobs

func (p *Proxy) WatchJobs(stream schedulerv1pb.Scheduler_WatchJobsServer) error

WatchJobs bidirectionally forwards messages between the daprd-side stream and the upstream scheduler stream. When either direction errors we cancel the shared context to unblock the other goroutine and drain its error so no goroutine leaks.

Jump to

Keyboard shortcuts

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