|
17 | 17 | from typing import List, Optional |
18 | 18 | from unittest import mock |
19 | 19 |
|
| 20 | +import grpc |
20 | 21 | from dapr.ext.workflow.dapr_workflow_context import DaprWorkflowContext |
21 | 22 | from dapr.ext.workflow.workflow_activity_context import WorkflowActivityContext |
22 | 23 | from dapr.ext.workflow.workflow_runtime import WorkflowRuntime, alternate_name |
@@ -46,6 +47,63 @@ def add_named_activity(self, name: str, fn): |
46 | 47 | self._activity_fns[name] = fn |
47 | 48 |
|
48 | 49 |
|
| 50 | +class WorkflowRuntimeTimeoutInterceptorTest(unittest.TestCase): |
| 51 | + def setUp(self): |
| 52 | + listActivities.clear() |
| 53 | + listOrchestrators.clear() |
| 54 | + self._registry_patch = mock.patch( |
| 55 | + 'dapr.ext.workflow._durabletask.worker._Registry', |
| 56 | + return_value=FakeTaskHubGrpcWorker(), |
| 57 | + ) |
| 58 | + self._registry_patch.start() |
| 59 | + |
| 60 | + def tearDown(self): |
| 61 | + mock.patch.stopall() |
| 62 | + |
| 63 | + def test_timeout_interceptor_is_added(self): |
| 64 | + with mock.patch( |
| 65 | + 'dapr.ext.workflow._durabletask.worker.TaskHubGrpcWorker' |
| 66 | + ) as mock_worker_cls: |
| 67 | + WorkflowRuntime() |
| 68 | + mock_worker_cls.assert_called_once() |
| 69 | + call_kwargs = mock_worker_cls.call_args[1] |
| 70 | + interceptors = call_kwargs['interceptors'] |
| 71 | + self.assertEqual(len(interceptors), 1) |
| 72 | + from dapr.clients.grpc.interceptors import DaprClientTimeoutInterceptor |
| 73 | + |
| 74 | + self.assertIsInstance(interceptors[0], DaprClientTimeoutInterceptor) |
| 75 | + |
| 76 | + def test_timeout_interceptor_with_custom_interceptors(self): |
| 77 | + custom_interceptor = mock.MagicMock(spec=grpc.UnaryUnaryClientInterceptor) |
| 78 | + with mock.patch( |
| 79 | + 'dapr.ext.workflow._durabletask.worker.TaskHubGrpcWorker' |
| 80 | + ) as mock_worker_cls: |
| 81 | + WorkflowRuntime(interceptors=[custom_interceptor]) |
| 82 | + call_kwargs = mock_worker_cls.call_args[1] |
| 83 | + interceptors = call_kwargs['interceptors'] |
| 84 | + self.assertEqual(len(interceptors), 2) |
| 85 | + from dapr.clients.grpc.interceptors import DaprClientTimeoutInterceptor |
| 86 | + |
| 87 | + self.assertIs(interceptors[0], custom_interceptor) |
| 88 | + self.assertIsInstance(interceptors[1], DaprClientTimeoutInterceptor) |
| 89 | + |
| 90 | + def test_timeout_interceptor_preserves_custom_interceptor_order(self): |
| 91 | + custom1 = mock.MagicMock(spec=grpc.UnaryUnaryClientInterceptor) |
| 92 | + custom2 = mock.MagicMock(spec=grpc.UnaryStreamClientInterceptor) |
| 93 | + with mock.patch( |
| 94 | + 'dapr.ext.workflow._durabletask.worker.TaskHubGrpcWorker' |
| 95 | + ) as mock_worker_cls: |
| 96 | + WorkflowRuntime(interceptors=[custom1, custom2]) |
| 97 | + call_kwargs = mock_worker_cls.call_args[1] |
| 98 | + interceptors = call_kwargs['interceptors'] |
| 99 | + self.assertEqual(len(interceptors), 3) |
| 100 | + from dapr.clients.grpc.interceptors import DaprClientTimeoutInterceptor |
| 101 | + |
| 102 | + self.assertIs(interceptors[0], custom1) |
| 103 | + self.assertIs(interceptors[1], custom2) |
| 104 | + self.assertIsInstance(interceptors[2], DaprClientTimeoutInterceptor) |
| 105 | + |
| 106 | + |
49 | 107 | class WorkflowRuntimeTest(unittest.TestCase): |
50 | 108 | def setUp(self): |
51 | 109 | listActivities.clear() |
@@ -765,3 +823,9 @@ def my_act(ctx, order: Optional[Order]): |
765 | 823 | wrapper = self.fake_registry._activity_fns['optional_no_default_act'] |
766 | 824 |
|
767 | 825 | self.assertIsNone(wrapper(mock.MagicMock(), None)) |
| 826 | + wrapper = self.fake_registry._activity_fns['optional_no_default_act'] |
| 827 | + |
| 828 | + self.assertIsNone(wrapper(mock.MagicMock(), None)) |
| 829 | + wrapper = self.fake_registry._activity_fns['optional_no_default_act'] |
| 830 | + |
| 831 | + self.assertIsNone(wrapper(mock.MagicMock(), None)) |
0 commit comments