temporal-python-testing - integration-testing
发布时间:2026/10/2 18:02:55来源:尧图网络
使用模拟活动进行集成测试使用模拟外部依赖、错误注入和复杂场景测试工作流的全面模式。活动模拟策略目的: 在不调用真实外部服务的情况下测试工作流编排逻辑基本模拟模式importpytestfromtemporalio.testingimportWorkflowEnvironmentfromtemporalio.workerimportWorkerfromunittest.mockimportMockpytest.mark.asyncioasyncdeftest_workflow_with_mocked_activity(workflow_env):Mock activity to test workflow logic# Create mock activitymock_activityMock(return_valuemocked-result)workflow.defnclassWorkflowWithActivity:workflow.runasyncdefrun(self,input:str)-str:resultawaitworkflow.execute_activity(process_external_data,input,start_to_close_timeouttimedelta(seconds10),)returnfprocessed:{result}asyncwithWorker(workflow_env.client,task_queuetest,workflows[WorkflowWithActivity],activities[mock_activity],# Use mock instead of real activity):resultawaitworkflow_env.client.execute_workflow(WorkflowWithActivity.run,test-input,idwf-mock,task_queuetest,)assertresultprocessed: mocked-resultmock_activity.assert_called_once()动态模拟响应基于场景的模拟pytest.mark.asyncioasyncdeftest_workflow_multiple_mock_scenarios(workflow_env):Test different workflow paths with dynamic mocks# Mock returns different values based on inputdefdynamic_activity(input:str)-str:ifinputerror-case:raiseApplicationError(Validation failed,non_retryableTrue)returnfprocessed-{input}workflow.defnclassDynamicWorkflow:workflow.runasyncdefrun(self,input:str)-str:try:resultawaitworkflow.execute_activity(dynamic_activity,input,start_to_close_timeouttimedelta(seconds10),)returnfsuccess:{result}exceptApplicationErrorase:returnferror:{e.message}asyncwithWorker(workflow_env.client,task_queuetest,workflows[DynamicWorkflow],activities[dynamic_activity],):# Test success pathresult_successawaitworkflow_env.client.execute_workflow(DynamicWorkflow.run,valid-input,idwf-success,task_queuetest,)assertresult_successsuccess: processed-valid-input# Test error pathresult_errorawaitworkflow_env.client.execute_workflow(DynamicWorkflow.run,error-case,idwf-error,task_queuetest,)assertValidation failedinresult_error错误注入模式测试瞬时失败重试行为pytest.mark.asyncioasyncdeftest_workflow_transient_errors(workflow_env):Test retry logic with controlled failuresattempt_count0activity.defnasyncdeftransient_activity()-str:nonlocalattempt_count attempt_count1ifattempt_count3:raiseException(fTransient error{attempt_count})returnsuccess-after-retriesworkflow.defnclassRetryWorkflow:workflow.runasyncdefrun(self)-str:returnawaitworkflow.execute_activity(transient_activity,start_to_close_timeouttimedelta(seconds10),retry_policyRetryPolicy(initial_intervaltimedelta(milliseconds10),maximum_attempts5,backoff_coefficient1.0,),)asyncwithWorker(workflow_env.client,task_queuetest,workflows[RetryWorkflow],activities[transient_activity],):resultawaitworkflow_env.client.execute_workflow(RetryWorkflow.run,idretry-wf,task_queuetest,)assertresultsuccess-after-retriesassertattempt_count3测试不可重试错误业务验证失败pytest.mark.asyncioasyncdeftest_workflow_non_retryable_error(workflow_env):Test handling of permanent failuresactivity.defnasyncdefvalidation_activity(input:dict)-str:ifnotinput.get(valid):raiseApplicationError(Invalid input,non_retryableTrue,# Dont retry validation errors)returnvalidatedworkflow.defnclassValidationWorkflow:workflow.runasyncdefrun(self,input:dict)-str:try:returnawaitworkflow.execute_activity(validation_activity,input,start_to_close_timeouttimedelta(seconds10),)exceptApplicationErrorase:returnfvalidation-failed:{e.message}asyncwithWorker(workflow_env.client,task_queuetest,workflows[ValidationWorkflow],activities[validation_activity],):resultawaitworkflow_env.client.execute_workflow(ValidationWorkflow.run,{valid:False},idvalidation-wf,task_queuetest,)assertvalidation-failedinresult多活动工作流测试顺序活动模式pytest.mark.asyncioasyncdeftest_workflow_sequential_activities(workflow_env):Test workflow orchestrating multiple activitiesactivity_calls[]activity.defnasyncdefstep_1(input:str)-str:activity_calls.append(step_1)returnf{input}-step1activity.defnasyncdefstep_2(input:str)-str:activity_calls.append(step_2)returnf{input}-step2activity.defnasyncdefstep_3(input:str)-str:activity_calls.append(step_3)returnf{input}-step3workflow.defnclassSequentialWorkflow:workflow.runasyncdefrun(self,input:str)-str:result_1awaitworkflow.execute_activity(step_1,input,start_to_close_timeouttimedelta(seconds10),)result_2awaitworkflow.execute_activity(step_2,result_1,start_to_close_timeouttimedelta(seconds10),)result_3awaitworkflow.execute_activity(step_3,result_2,start_to_close_timeouttimedelta(seconds10),)returnresult_3asyncwithWorker(workflow_env.client,task_queuetest,workflows[SequentialWorkflow],activities[step_1,step_2,step_3],):resultawaitworkflow_env.client.execute_workflow(SequentialWorkflow.run,start,idseq-wf,task_queuetest,)assertresultstart-step1-step2-step3assertactivity_calls[step_1,step_2,step_3]并行活动模式pytest.mark.asyncioasyncdeftest_workflow_parallel_activities(workflow_env):Test concurrent activity executionactivity.defnasyncdefparallel_task(task_id:int)-str:returnftask-{task_id}workflow.defnclassParallelWorkflow:workflow.runasyncdefrun(self,task_count:int)-list[str]:# Execute activities in paralleltasks[workflow.execute_activity(parallel_task,i,start_to_close_timeouttimedelta(seconds10),)foriinrange(task_count)]returnawaitasyncio.gather(*tasks)asyncwithWorker(workflow_env.client,task_queuetest,workflows[ParallelWorkflow],activities[parallel_task],):resultawaitworkflow_env.client.execute_workflow(ParallelWorkflow.run,3,idparallel-wf,task_queuetest,)assertresult[task-0,task-1,task-2]信号和查询测试信号处理器pytest.mark.asyncioasyncdeftest_workflow_signals(workflow_env):Test workflow signal handlingworkflow.defnclassSignalWorkflow:def__init__(self)-None:self._statusinitializedworkflow.runasyncdefrun(self)-str:# Wait for completion signalawaitworkflow.wait_condition(lambda:self._statuscompleted)returnself._statusworkflow.signalasyncdefupdate_status(self,new_status:str)-None:self._statusnew_statusworkflow.querydefget_status(self)-str:returnself._statusasyncwithWorker(workflow_env.client,task_queuetest,workflows[SignalWorkflow],):# Start workflowhandleawaitworkflow_env.client.start_workflow(SignalWorkflow.run,idsignal-wf,task_queuetest,)# Verify initial state via queryinitial_statusawaithandle.query(SignalWorkflow.get_status)assertinitial_statusinitialized# Send signalawaithandle.signal(SignalWorkflow.update_status,processing)# Verify updated stateupdated_statusawaithandle.query(SignalWorkflow.get_status)assertupdated_statusprocessing# Complete workflowawaithandle.signal(SignalWorkflow.update_status,completed)resultawaithandle.result()assertresultcompleted覆盖率策略工作流逻辑覆盖率目标: 工作流决策逻辑覆盖率 ≥80%# Test all branchespytest.mark.parametrize(condition,expected,[(True,branch-a),(False,branch-b),])asyncdeftest_workflow_branches(workflow_env,condition,expected):Ensure all code paths are tested# Test implementationpass活动覆盖率目标: 活动逻辑覆盖率 ≥80%# Test activity edge casespytest.mark.parametrize(input,expected,[(valid,success),(,empty-input-error),(None,null-input-error),])asyncdeftest_activity_edge_cases(activity_env,input,expected):Test activity error handling# Test implementationpass集成测试组织测试结构tests/ ├── integration/ │ ├── conftest.py # Shared fixtures │ ├── test_order_workflow.py # Order processing tests │ ├── test_payment_workflow.py # Payment tests │ └── test_fulfillment_workflow.py ├── unit/ │ ├── test_order_activities.py │ └── test_payment_activities.py └── fixtures/ └── test_data.py # Test data builders共享夹具# conftest.pyimportpytestfromtemporalio.testingimportWorkflowEnvironmentpytest.fixture(scopesession)asyncdefworkflow_env():Session-scoped environment for integration testsenvawaitWorkflowEnvironment.start_time_skipping()yieldenvawaitenv.shutdown()pytest.fixturedefmock_payment_service():Mock external payment servicereturnMock()pytest.fixturedefmock_inventory_service():Mock external inventory servicereturnMock()最佳实践模拟外部依赖: 测试中绝不调用真实 API测试错误场景: 验证补偿和重试逻辑并行测试: 使用 pytest-xdist 加快测试运行隔离测试: 每个测试应该独立清晰断言: 验证结果和副作用覆盖率目标: 关键工作流 ≥80%快速执行: 使用时间跳跃避免真实延迟其他资源模拟策略: docs.temporal.io/develop/python/testing-suitepytest 最佳实践: docs.pytest.org/en/stable/goodpractices.htmlPython SDK 示例: github.com/temporalio/samples-python
网站建设高端定制企业官网