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 ¶
- func DefaultReporter(context.Context, compapi.Component, *operatorv1.ResourceResult) error
- type BindingManager
- type Options
- type Processor
- func (p *Processor) AddPendingComponent(ctx context.Context, comp compapi.Component) <-chan error
- func (p *Processor) AddPendingEndpoint(ctx context.Context, endpoint httpendpointsapi.HTTPEndpoint) <-chan error
- func (p *Processor) AddPendingMCPServer(ctx context.Context, s mcpserverapi.MCPServer) <-chan error
- func (p *Processor) AddPendingSubscription(ctx context.Context, subscriptions ...subapi.Subscription) <-chan error
- func (p *Processor) Binding() BindingManager
- func (p *Processor) Close(ctx context.Context, comp compapi.Component) error
- func (p *Processor) CloseSubscription(ctx context.Context, sub *subapi.Subscription) error
- func (p *Processor) DeleteMCPServer(ctx context.Context, name string)
- func (p *Processor) Flush(ctx context.Context) error
- func (p *Processor) Init(ctx context.Context, comp compapi.Component) error
- func (p *Processor) OnActorStateStoreChanged()
- func (p *Processor) Process(ctx context.Context) error
- func (p *Processor) ProcessMCPServerSecrets(ctx context.Context, s *mcpserverapi.MCPServer)
- func (p *Processor) Secret() SecretManager
- func (p *Processor) SetInProcessWorkflows(r wfregistrar.Registrar)
- func (p *Processor) State() StateManager
- func (p *Processor) Subscriber() SubscribeManager
- func (p *Processor) WorkflowBackend() WorkflowBackendManager
- type SecretManager
- type StateManager
- type SubscribeManager
- type WorkflowBackendManager
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func DefaultReporter ¶ added in v1.15.0
func DefaultReporter(context.Context, compapi.Component, *operatorv1.ResourceResult) error
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 ¶
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
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
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 ¶
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
CloseSubscription closes a declarative subscription synchronously.
func (*Processor) DeleteMCPServer ¶ added in v1.18.0
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 ¶
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 ¶
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
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 ¶
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).