2445
}
2446
2448
>
// Create workflow execution.
2449
>
workflowBranchToken, workflowSnapshot, _ := s.CreateWorkflow(
2450
>
rand.Int63(),
2451
>
enumsspb.WORKFLOW_EXECUTION_STATE_CREATED,
2452
>
enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING,
2453
>
rand.Int63(),
2454
>
)
2455
>
2456
>
// Create CHASM execution with same ID.
2457
>
// Do NOT use CreateCHASMExecution() here, which uses the runID as CreateWorkflow().
2458
>
chasmRunID := uuid.New().String()
2459
>
chasmArchetypeID := rand.Uint32()
2460
>
chasmSnapshot, events := RandomSnapshot(
2461
>
s.T(),
2462
>
s.NamespaceID,
2463
>
s.WorkflowID,
2464
>
chasmRunID,
2465
>
common.FirstEventID,
2466
>
rand.Int63(),
2467
>
enumsspb.WORKFLOW_EXECUTION_STATE_CREATED,
2468
>
enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING,
2469
>
rand.Int63(),
2470
>
nil,
2471
>
)
2472
>
_, err := s.ExecutionManager.CreateWorkflowExecution(s.Ctx, &p.CreateWorkflowExecutionRequest{
2473
>
ShardID: s.ShardID,
2474
>
RangeID: s.RangeID,
2475
>
Mode: p.CreateWorkflowModeBrandNew,
2476
>
2477
>
PreviousRunID: "",
2478
>
PreviousLastWriteVersion: 0,
2479
>
2480
>
ArchetypeID: chasmArchetypeID,
2481
>
2482
>
NewWorkflowSnapshot: *chasmSnapshot,
2483
>
NewWorkflowEvents: events,
2484
>
})
2485
>
s.NoError(err)
2486
>
2487
>
// Update Workflow execution.
2488
>
workflowMutation, workflowEvents := RandomMutation(
2489
>
s.T(),
2490
>
s.NamespaceID,
2491
>
s.WorkflowID,
2492
>
s.RunID,
2493
>
workflowSnapshot.NextEventID,
2494
>
rand.Int63(),
2495
>
enumsspb.WORKFLOW_EXECUTION_STATE_RUNNING,
2496
>
enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING,
2497
>
workflowSnapshot.DBRecordVersion+1,
2498
>
workflowBranchToken,
2499
>
)
2500
>
_, err = s.ExecutionManager.UpdateWorkflowExecution(s.Ctx, &p.UpdateWorkflowExecutionRequest{
2501
>
ShardID: s.ShardID,
2502
>
RangeID: s.RangeID,
2503
>
Mode: p.UpdateWorkflowModeUpdateCurrent,
2504
>
2505
>
ArchetypeID: chasm.WorkflowArchetypeID,
2506
>
2507
>
UpdateWorkflowMutation: *workflowMutation,
2508
>
UpdateWorkflowEvents: workflowEvents,
2509
>
2510
>
NewWorkflowSnapshot: nil,
2511
>
NewWorkflowEvents: nil,
2512
>
})
2513
>
s.NoError(err)
2514
>
2515
>
// Update CHASM execution.
2516
>
chasmMutation, chasmEvents := RandomMutation(
2517
>
s.T(),
2518
>
s.NamespaceID,
2519
>
s.WorkflowID,
2520
>
chasmRunID,
2521
>
chasmSnapshot.NextEventID,
2522
>
rand.Int63(),
2523
>
enumsspb.WORKFLOW_EXECUTION_STATE_RUNNING,
2524
>
enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING,
2525
>
chasmSnapshot.DBRecordVersion+1,
2526
>
nil, // No branch token for CHASM
2527
>
)
2528
>
_, err = s.ExecutionManager.UpdateWorkflowExecution(s.Ctx, &p.UpdateWorkflowExecutionRequest{
2529
>
ShardID: s.ShardID,
2530
>
RangeID: s.RangeID,
2531
>
Mode: p.UpdateWorkflowModeUpdateCurrent,
2532
>
2533
>
ArchetypeID: chasmArchetypeID,
2534
>
2535
>
UpdateWorkflowMutation: *chasmMutation,
2536
>
UpdateWorkflowEvents: chasmEvents,
2537
>
2538
>
NewWorkflowSnapshot: nil,
2539
>
NewWorkflowEvents: nil,
2540
>
})
2541
>
s.NoError(err)
2542
>
2543
>
// Validate current execution for both archetypes.
2544
>
resp, err := s.ExecutionManager.GetCurrentExecution(s.Ctx, &p.GetCurrentExecutionRequest{
2545
>
ShardID: s.ShardID,
2546
>
NamespaceID: s.NamespaceID,
2547
>
WorkflowID: s.WorkflowID,
2548
>
ArchetypeID: chasm.WorkflowArchetypeID,
2549
>
})
2550
>
s.NoError(err)
2551
>
s.Equal(s.RunID, resp.RunID)
2552
>
resp, err = s.ExecutionManager.GetCurrentExecution(s.Ctx, &p.GetCurrentExecutionRequest{
2553
>
ShardID: s.ShardID,
2554
>
NamespaceID: s.NamespaceID,
2555
>
WorkflowID: s.WorkflowID,
2556
>
ArchetypeID: chasmArchetypeID,
2557
>
})
2558
>
s.NoError(err)
2559
>
s.Equal(chasmRunID, resp.RunID)
2560
>
2561
>
// Validate concrete execution for both archetypes.
2562
>
s.AssertMSEqualWithDB(chasm.WorkflowArchetypeID, workflowSnapshot, workflowMutation)
2563
>
s.AssertMSEqualWithDB(chasmArchetypeID, chasmSnapshot, chasmMutation)
2564
>
2565
>
// Delete workflow execution.
2566
>
err = s.ExecutionManager.DeleteCurrentWorkflowExecution(s.Ctx, &p.DeleteCurrentWorkflowExecutionRequest{
2567
>
ShardID: s.ShardID,
2568
>
NamespaceID: s.NamespaceID,
2569
>
WorkflowID: s.WorkflowID,
2570
>
RunID: s.RunID,
2571
>
ArchetypeID: chasm.WorkflowArchetypeID,
2572
>
})
2573
>
s.NoError(err)
2574
>
err = s.ExecutionManager.DeleteWorkflowExecution(s.Ctx, &p.DeleteWorkflowExecutionRequest{
2575
>
ShardID: s.ShardID,
2576
>
NamespaceID: s.NamespaceID,
2577
>
WorkflowID: s.WorkflowID,
2578
>
RunID: s.RunID,
2579
>
ArchetypeID: chasm.WorkflowArchetypeID,
2580
>
})
2581
>
s.NoError(err)
2582
>
2583
>
s.AssertMissingFromDB(s.NamespaceID, s.WorkflowID, s.RunID, chasm.WorkflowArchetypeID)
2584
>
s.AssertMSEqualWithDB(chasmArchetypeID, chasmSnapshot, chasmMutation)
2585
>
2586
>
// Delete CHASM execution.
2587
>
err = s.ExecutionManager.DeleteCurrentWorkflowExecution(s.Ctx, &p.DeleteCurrentWorkflowExecutionRequest{
2588
>
ShardID: s.ShardID,
2589
>
NamespaceID: s.NamespaceID,
2590
>
WorkflowID: s.WorkflowID,
2591
>
RunID: chasmRunID,
2592
>
ArchetypeID: chasmArchetypeID,
2593
>
})
2594
>
s.NoError(err)
2595
>
err = s.ExecutionManager.DeleteWorkflowExecution(s.Ctx, &p.DeleteWorkflowExecutionRequest{
2596
>
ShardID: s.ShardID,
2597
>
NamespaceID: s.NamespaceID,
2598
>
WorkflowID: s.WorkflowID,
2599
>
RunID: chasmRunID,
2600
>
ArchetypeID: chasmArchetypeID,
2601
>
})
2602
>
s.NoError(err)
2603
>
2604
>
s.AssertMissingFromDB(s.NamespaceID, s.WorkflowID, chasmRunID, chasmArchetypeID)
2605
>
}
2606
2607
func (s *ExecutionMutableStateSuite) CreateWorkflow(