Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 8 additions & 8 deletions dapr/ext/databricks/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -78,14 +78,14 @@ identity layer on top of the existing `dapr.ext.workflow.DaprWorkflowClient`.

```python
from dapr.ext.databricks import (
register_workflow_sink, # high-level: registers a foreach_batch_sink
DaprWorkflowBatchHandler, # low-level: reusable batch-processing handler
WorkflowSinkConfig, # typed config, if constructing a handler directly
default_row_mapper, # the default Row -> dict mapper
DaprDatabricksError, # base exception
SinkConfigurationError, # bad register_workflow_sink() config
MissingBusinessKeyError, # configured id_field/id_fields absent or null on a row
DaprDatabricksSinkError, # a micro-batch could not be durably handed off
register_workflow_sink, # high-level: registers a foreach_batch_sink
DaprWorkflowBatchHandler, # low-level: reusable batch-processing handler
WorkflowSinkConfig, # typed config, if constructing a handler directly
default_row_mapper, # the default Row -> dict mapper
DaprDatabricksError, # base exception
SinkConfigurationError, # bad register_workflow_sink() config
MissingBusinessKeyError, # configured id_field/id_fields absent or null on a row
DaprDatabricksSinkError, # a micro-batch could not be durably handed off
)
```

Expand Down
36 changes: 17 additions & 19 deletions dapr/ext/databricks/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@ register_workflow_sink(
namespace='orders',
)


@dp.append_flow(target='order_actions', name='order_actions_flow')
def order_actions_flow():
return spark.readStream.table('validated_orders')
Expand Down Expand Up @@ -81,7 +82,7 @@ register_workflow_sink(
register_workflow_sink(
name='transfers',
workflow='process_transfer',
instance_id_factory=lambda row, batch_id: f"{row['account_id']}:{row['transaction_id']}",
instance_id_factory=lambda row, batch_id: f'{row["account_id"]}:{row["transaction_id"]}',
)
```

Expand All @@ -108,6 +109,7 @@ handler = DaprWorkflowBatchHandler(
WorkflowSinkConfig(name='orders', workflow='process_order', id_field='order_id')
)


@dp.foreach_batch_sink(name='orders')
def orders_handler(df, batch_id):
handler.process(df, batch_id)
Expand Down Expand Up @@ -171,24 +173,20 @@ scheduling latency. Business record contents are never logged.

```python
register_workflow_sink(
name='orders', # sink name (foreach_batch_sink name)
workflow='process_order', # registered Dapr Workflow name

id_field='order_id', # business key column (mutually exclusive with the two below)
id_fields=None, # composite business key columns
instance_id_factory=None, # (row, batch_id) -> business-key component

namespace='orders', # identity partition; part of the instance ID
generation='v1', # identity epoch; bump to replay after a full refresh

input_mapper=None, # Row -> dict; defaults to a JSON-safe asDict()
metadata=True, # wrap input as {'data': ..., 'metadata': {...}}

max_in_flight=8, # bounded concurrent schedule_new_workflow calls
max_records_per_batch=None, # optional hard cap; fails the batch, never truncates silently

host=None, port=None, # Dapr endpoint; defaults to the standard SDK env/settings
allow_batch_position_identity=False, # required opt-in to run without any id_field/id_fields/instance_id_factory
name='orders', # sink name (foreach_batch_sink name)
workflow='process_order', # registered Dapr Workflow name
id_field='order_id', # business key column (mutually exclusive with the two below)
id_fields=None, # composite business key columns
instance_id_factory=None, # (row, batch_id) -> business-key component
namespace='orders', # identity partition; part of the instance ID
generation='v1', # identity epoch; bump to replay after a full refresh
input_mapper=None, # Row -> dict; defaults to a JSON-safe asDict()
metadata=True, # wrap input as {'data': ..., 'metadata': {...}}
max_in_flight=8, # bounded concurrent schedule_new_workflow calls
max_records_per_batch=None, # optional hard cap; fails the batch, never truncates silently
host=None,
port=None, # Dapr endpoint; defaults to the standard SDK env/settings
allow_batch_position_identity=False, # required opt-in to run without any id_field/id_fields/instance_id_factory
)
```

Expand Down
6 changes: 4 additions & 2 deletions dapr/ext/fastapi/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,10 @@ Wraps a FastAPI instance to add Dapr pub/sub event handling.
app = FastAPI()
dapr_app = DaprApp(app, router_tags=['PubSub']) # router_tags optional, default ['PubSub']

@dapr_app.subscribe(pubsub='pubsub', topic='orders', route='/handle-order',
metadata={}, dead_letter_topic=None)

@dapr_app.subscribe(
pubsub='pubsub', topic='orders', route='/handle-order', metadata={}, dead_letter_topic=None
)
def handle_order(event_data):
return {'status': 'ok'}
```
Expand Down
6 changes: 4 additions & 2 deletions dapr/ext/flask/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,10 @@ Wraps a Flask instance to add Dapr pub/sub event handling.
app = Flask('myapp')
dapr_app = DaprApp(app)

@dapr_app.subscribe(pubsub='pubsub', topic='orders', route='/handle-order',
metadata={}, dead_letter_topic=None)

@dapr_app.subscribe(
pubsub='pubsub', topic='orders', route='/handle-order', metadata={}, dead_letter_topic=None
)
def handle_order():
event_data = request.json
return 'ok'
Expand Down
1 change: 1 addition & 0 deletions dapr/ext/flask/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ from dapr.ext.flask import DaprApp
app = Flask('myapp')
dapr_app = DaprApp(app)


@dapr_app.subscribe(pubsub='pubsub', topic='some_topic', route='/some_endpoint')
def my_event_handler():
# request.data contains the pubsub event
Expand Down
49 changes: 28 additions & 21 deletions dapr/ext/grpc/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,17 +25,17 @@ Installed via the `grpc` extra on core dapr: `pip install "dapr[grpc]"`.

```python
from dapr.ext.grpc import (
App, # Main entry point — decorator-based gRPC server
Rule, # CEL-based topic rule with priority
SubscriptionMessage, # Event type received by pub/sub topic handlers (preferred)
InvokeMethodRequest, # Request object for service invocation handlers
InvokeMethodResponse, # Response object for service invocation handlers
BindingRequest, # Request object for input binding handlers
TopicEventResponse, # Response object for pub/sub handlers
Job, # Job definition for scheduler
JobEvent, # Job event received by handler
FailurePolicy, # ABC for job failure policies
DropFailurePolicy, # Drop on failure (no retry)
App, # Main entry point — decorator-based gRPC server
Rule, # CEL-based topic rule with priority
SubscriptionMessage, # Event type received by pub/sub topic handlers (preferred)
InvokeMethodRequest, # Request object for service invocation handlers
InvokeMethodResponse, # Response object for service invocation handlers
BindingRequest, # Request object for input binding handlers
TopicEventResponse, # Response object for pub/sub handlers
Job, # Job definition for scheduler
JobEvent, # Job event received by handler
FailurePolicy, # ABC for job failure policies
DropFailurePolicy, # Drop on failure (no retry)
ConstantFailurePolicy, # Retry with constant interval
)
```
Expand All @@ -51,22 +51,29 @@ The central entry point. Creates a gRPC server and provides decorators for handl
```python
app = App()


@app.method('method_name')
def handle_method(request: InvokeMethodRequest) -> InvokeMethodResponse:
...
def handle_method(request: InvokeMethodRequest) -> InvokeMethodResponse: ...


@app.subscribe(
pubsub_name='pubsub',
topic='orders',
metadata={},
dead_letter_topic=None,
rule=Rule('event.type == "order"', priority=1),
disable_topic_validation=False,
)
def handle_event(event: SubscriptionMessage) -> Optional[TopicEventResponse]: ...

@app.subscribe(pubsub_name='pubsub', topic='orders', metadata={}, dead_letter_topic=None,
rule=Rule('event.type == "order"', priority=1), disable_topic_validation=False)
def handle_event(event: SubscriptionMessage) -> Optional[TopicEventResponse]:
...

@app.binding('binding_name')
def handle_binding(request: BindingRequest) -> None:
...
def handle_binding(request: BindingRequest) -> None: ...


@app.job_event('job_name')
def handle_job(event: JobEvent) -> None:
...
def handle_job(event: JobEvent) -> None: ...


app.register_health_check(lambda: None) # Not a decorator — direct registration
```
Expand Down
13 changes: 9 additions & 4 deletions dapr/ext/rag/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -118,10 +118,15 @@ per-adapter optional-import guards).

```python
from dapr.ext.rag import (
DurableRAGPipeline, PipelineConfig,
S3Source, AzureBlobSource,
UnstructuredParser, TextSplitter, OpenAIEmbedder,
PgVectorStore, PineconeVectorStore,
DurableRAGPipeline,
PipelineConfig,
S3Source,
AzureBlobSource,
UnstructuredParser,
TextSplitter,
OpenAIEmbedder,
PgVectorStore,
PineconeVectorStore,
ActiveVersionResolver,
)

Expand Down
9 changes: 8 additions & 1 deletion dapr/ext/rag/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,14 @@ pip install "dapr[rag,rag-azure,rag-pinecone]" # Azure Blob + Pinecone
```

```python
from dapr.ext.rag import DurableRAGPipeline, S3Source, UnstructuredParser, TextSplitter, OpenAIEmbedder, PgVectorStore
from dapr.ext.rag import (
DurableRAGPipeline,
S3Source,
UnstructuredParser,
TextSplitter,
OpenAIEmbedder,
PgVectorStore,
)

pipeline = DurableRAGPipeline(
source=S3Source(bucket='company-docs', prefix='policies/'),
Expand Down
6 changes: 3 additions & 3 deletions dapr/ext/strands/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,9 +31,9 @@ Extends both `RepositorySessionManager` and `SessionRepository` from the Strands
manager = DaprSessionManager(
session_id='my-session',
state_store_name='statestore',
dapr_client=client, # DaprClient instance
ttl=3600, # Optional: TTL in seconds
consistency='eventual', # 'eventual' (default) or 'strong'
dapr_client=client, # DaprClient instance
ttl=3600, # Optional: TTL in seconds
consistency='eventual', # 'eventual' (default) or 'strong'
)
```

Expand Down
18 changes: 9 additions & 9 deletions dapr/ext/workflow/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,16 +76,16 @@ All public symbols are exported from `dapr.ext.workflow`:

```python
from dapr.ext.workflow import (
WorkflowRuntime, # Registration & lifecycle (start/shutdown)
DaprWorkflowClient, # Sync client for scheduling/managing workflows
DaprWorkflowContext, # Passed to workflow functions as first arg
WorkflowRuntime, # Registration & lifecycle (start/shutdown)
DaprWorkflowClient, # Sync client for scheduling/managing workflows
DaprWorkflowContext, # Passed to workflow functions as first arg
WorkflowActivityContext, # Passed to activity functions as first arg
WorkflowState, # Snapshot of a workflow instance's state
WorkflowStatus, # Enum: UNKNOWN, RUNNING, COMPLETED, FAILED, TERMINATED, PENDING, SUSPENDED, STALLED
when_all, # Parallel combinator — wait for all tasks
when_any, # Race combinator — wait for first task
alternate_name, # Decorator to set a custom registration name
RetryPolicy, # Retry config for activities/child workflows
WorkflowState, # Snapshot of a workflow instance's state
WorkflowStatus, # Enum: UNKNOWN, RUNNING, COMPLETED, FAILED, TERMINATED, PENDING, SUSPENDED, STALLED
when_all, # Parallel combinator — wait for all tasks
when_any, # Race combinator — wait for first task
alternate_name, # Decorator to set a custom registration name
RetryPolicy, # Retry config for activities/child workflows
)

# Async client:
Expand Down
1 change: 1 addition & 0 deletions dapr/ext/workflow/docs/concurrency.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ cap:
```python
_shared_client: httpx.AsyncClient | None = None


def _get_client() -> httpx.AsyncClient:
global _shared_client
if _shared_client is None:
Expand Down
Loading
Loading