processor

package
v1.19.0-rc.2 Latest Latest
Warning

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

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

Documentation

Overview

Package processor manages the lifecycle of every resource that flows through daprd: components, HTTP endpoints, MCP servers and subscriptions. It is structured as a three-level dapr/kit/events/loop hierarchy:

root  ->  per-resource category  ->  per-named-instance

Per-resource public entry points are defined in sibling files (components.go, httpendpoint.go, mcpserver.go, subscription.go); this file holds the shared scaffolding: Options, the Processor struct, construction, the runner manager that drives every loop, the Flush barrier and the default reporter.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func DefaultReporter added in v1.15.0

DefaultReporter is the default resource reporter for the registry. It does nothing.

Types

type BindingManager

type BindingManager interface {
	SendToOutputBinding(context.Context, string, *bindings.InvokeRequest) (*bindings.InvokeResponse, error)
	StartReadingFromBindings(context.Context) error
	StopReadingFromBindings(forever bool)
}

BindingManager exposes the binding sub-processor's read paths (SendToOutputBinding) and category-wide operations (StartReadingFromBindings, StopReadingFromBindings). Init and Close are driven by the bindings category loop.

type Options

type Options struct {
	ID                              string
	Namespace                       string
	Mode                            modes.DaprMode
	ActorsEnabled                   bool
	Actors                          actors.Interface
	IsHTTP                          bool
	Registry                        *registry.Registry
	ComponentStore                  *compstore.ComponentStore
	Meta                            *meta.Meta
	GlobalConfig                    *config.Configuration
	Resiliency                      resiliency.Provider
	GRPC                            *grpcmanager.Manager
	Channels                        *channels.Channels
	OperatorClient                  operatorv1.OperatorClient
	MiddlewareHTTP                  *http.HTTP
	Security                        security.Handler
	Outbox                          outbox.Outbox
	Adapter                         rtpubsub.Adapter
	AdapterStreamer                 rtpubsub.AdapterStreamer
	Reporter                        registry.Reporter
	ProgrammaticSubscriptionEnabled bool
	AppBindingOptionsTimeout        time.Duration
}

type Processor

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

Processor manages the lifecycle of all components, HTTP endpoints, MCP servers and subscription resources via a three level dapr/kit/events/loop hierarchy: root -> per-category -> per-named-instance.

Concurrency: every mutation to component lifecycle state happens inside one of the loops in this hierarchy. The pending-secret-store-dependents map (for components blocked on a secret store) is owned by the root loop. Sub-processor read paths (SendToOutputBinding, ActorStateStoreName, ProcessResource, ...) still call directly into the sub-processor; the sub-processors retain their internal mutexes to guard those concurrent read callers from concurrent lifecycle writes.

func New

func New(opts Options) *Processor

New constructs a Processor. The returned Processor is not running; call Process to start the loop hierarchy.

func (*Processor) AddPendingComponent added in v1.13.0

func (p *Processor) AddPendingComponent(ctx context.Context, comp compapi.Component) <-chan error

AddPendingComponent enqueues a component init and returns a buffered chan that receives exactly one error (nil on success). Returns nil if the processor is shut down.

func (*Processor) AddPendingEndpoint added in v1.13.0

func (p *Processor) AddPendingEndpoint(ctx context.Context, endpoint httpendpointsapi.HTTPEndpoint) <-chan error

AddPendingEndpoint enqueues an HTTP endpoint and returns a result chan.

func (*Processor) AddPendingMCPServer added in v1.18.0

func (p *Processor) AddPendingMCPServer(ctx context.Context, s mcpserverapi.MCPServer) <-chan error

AddPendingMCPServer enqueues an MCP server and returns a result chan.

func (*Processor) AddPendingSubscription added in v1.14.0

func (p *Processor) AddPendingSubscription(ctx context.Context, subscriptions ...subapi.Subscription) <-chan error

AddPendingSubscription enqueues one or more declarative subscriptions. Returns a chan that emits one error per submitted subscription, in order.

func (*Processor) Binding

func (p *Processor) Binding() BindingManager

func (*Processor) Close

func (p *Processor) Close(ctx context.Context, comp compapi.Component) error

Close synchronously closes a component. If Process is running, the close is routed through the loop; otherwise it runs inline. When routed through the loop the wait honours ctx so a caller is not blocked indefinitely if the loop is being torn down.

func (*Processor) CloseSubscription added in v1.14.0

func (p *Processor) CloseSubscription(ctx context.Context, sub *subapi.Subscription) error

CloseSubscription closes a declarative subscription synchronously.

func (*Processor) DeleteMCPServer added in v1.18.0

func (p *Processor) DeleteMCPServer(ctx context.Context, name string)

DeleteMCPServer removes an MCP server from the compstore and unregisters its in-process workflows. Routed through the root loop so the delete is serialised against concurrent Add events for the same name. Returns when the delete has been applied, when ctx is cancelled, or immediately if the processor is already shut down.

func (*Processor) Flush

func (p *Processor) Flush(ctx context.Context) error

Flush blocks until every Init currently enqueued in the root loop has completed. It is a no-op when the processor has already shut down. The Barrier event is queued unconditionally, so it is safe to call Flush concurrently with Process; the call simply waits until the root loop has processed every event queued ahead of the Barrier.

func (*Processor) Init

func (p *Processor) Init(ctx context.Context, comp compapi.Component) error

Init synchronously initialises a component. If Process is running, the init is routed through the loop hierarchy and waits for the result; otherwise the init runs inline. Tests that drive the processor without calling Process rely on the inline fallback.

func (*Processor) OnActorStateStoreChanged added in v1.18.3

func (p *Processor) OnActorStateStoreChanged()

OnActorStateStoreChanged notifies the actor runtime that the actor state store was added, removed, or replaced. Safe to call when no actor runtime is configured.

func (*Processor) Process added in v1.13.0

func (p *Processor) Process(ctx context.Context) error

Process runs the loop hierarchy until ctx is cancelled. It coordinates orderly shutdown by closing the root loop, which fans out to category and instance loops.

func (*Processor) ProcessMCPServerSecrets added in v1.18.2

func (p *Processor) ProcessMCPServerSecrets(ctx context.Context, s *mcpserverapi.MCPServer)

ProcessMCPServerSecrets resolves secretKeyRef and envRef entries in the transport headers (spec.endpoint.streamableHTTP.headers or spec.endpoint.sse.headers) and spec.endpoint.stdio.env using the configured secret store. Unlike components, MCPServer resources load after all secret store components are initialized, so secrets are available immediately. The underlying p.secret.ProcessResource (the secret manager) logs errors internally and resolves what it can; it does not return an error. Unresolvable secretKeyRef values remain as empty strings.

This does NOT resolve auth.oauth2.secretKeyRef: the OAuth2 client secret is fetched at connection time by the MCP worker (see pkg/runtime/mcp/auth) and is never written into the spec, so it does not take part in the hot-reload diff.

The hot-reload reconciler also calls this on a copy of an incoming spec before comparing it against the already-resolved stored copy, so an unchanged secret-ref server is not needlessly reloaded while a rotated secret value still triggers a reload.

func (*Processor) Secret added in v1.13.0

func (p *Processor) Secret() SecretManager

func (*Processor) SetInProcessWorkflows added in v1.18.0

func (p *Processor) SetInProcessWorkflows(r wfregistrar.Registrar)

SetInProcessWorkflows installs the in-process workflow wfregistrar. Used by the root loop's MCPServer add path to register workflows backing managed MCP resources.

func (*Processor) State

func (p *Processor) State() StateManager

func (*Processor) Subscriber added in v1.14.0

func (p *Processor) Subscriber() SubscribeManager

func (*Processor) WorkflowBackend added in v1.13.0

func (p *Processor) WorkflowBackend() WorkflowBackendManager

type SecretManager added in v1.13.0

type SecretManager interface {
	// TODO: update to return an error.
	ProcessResource(context.Context, meta.Resource) (bool, string)
}

SecretManager exposes the secret sub-processor's ProcessResource read API (used by reconciler for resolving component references). Init and Close are driven by the secret category loop.

type StateManager

type StateManager interface {
	ActorStateStoreName() (string, bool)
}

StateManager exposes the state sub-processor's read API. Lifecycle Init and Close are driven by the state category loop.

type SubscribeManager added in v1.14.0

type SubscribeManager interface {
	InitProgramaticSubscriptions(context.Context) error
	StartAppSubscriptions() error
	StopAppSubscriptions()
	StopAllSubscriptionsForever()
	ReloadDeclaredAppSubscription(name, pubsubName string) error
	StartStreamerSubscription(sub *subapi.Subscription, connectionID rtpubsub.ConnectionID) error
	StopStreamerSubscription(sub *subapi.Subscription, connectionID rtpubsub.ConnectionID)
	ReloadPubSub(string) error
	StopPubSub(string)
}

SubscribeManager exposes the subscriber-facing operations. Each method remains synchronous; the underlying subscriber holds an internal mutex that guards against concurrent access from the pubsub category loop and from external callers (e.g. gRPC subscription stream handlers).

type WorkflowBackendManager added in v1.13.0

type WorkflowBackendManager interface {
	Backend() (backend.Backend, bool)
}

Jump to

Keyboard shortcuts

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