Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/database/_migrations/versions/postgresql/2022_10_12_102048_22b7cb02e593_add_state_timestamp.py: 78%

33 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1"""Add state_timestamp 

2 

3Revision ID: 22b7cb02e593 

4Revises: 2d5e000696f1 

5Create Date: 2022-10-12 10:20:48.760447 

6 

7""" 

8 

9import sqlalchemy as sa 

10from alembic import op 

11 

12import prefect 

13 

14# revision identifiers, used by Alembic. 

15revision = "22b7cb02e593" 

16down_revision = "2d5e000696f1" 

17branch_labels = None 

18depends_on = None 

19 

20 

21def upgrade(): 

22 with op.batch_alter_table("flow_run", schema=None) as batch_op: 

23 batch_op.add_column( 

24 sa.Column( 

25 "state_timestamp", 

26 prefect.server.utilities.database.Timestamp(timezone=True), 

27 nullable=True, 

28 ) 

29 ) 

30 batch_op.create_index( 

31 "ix_flow_run__state_timestamp", ["state_timestamp"], unique=False 

32 ) 

33 

34 with op.batch_alter_table("task_run", schema=None) as batch_op: 

35 batch_op.add_column( 

36 sa.Column( 

37 "state_timestamp", 

38 prefect.server.utilities.database.Timestamp(timezone=True), 

39 nullable=True, 

40 ) 

41 ) 

42 batch_op.create_index( 

43 "ix_task_run__state_timestamp", ["state_timestamp"], unique=False 

44 ) 

45 

46 update_flow_run_state_timestamp_in_batches = """ 

47 WITH null_flow_run_state_timestamp_cte as (SELECT id from flow_run where state_timestamp is null and state_id is not null limit 500) 

48 UPDATE flow_run 

49 SET state_timestamp = flow_run_state.timestamp 

50 FROM flow_run_state, null_flow_run_state_timestamp_cte 

51 WHERE flow_run.state_id = flow_run_state.id 

52 AND flow_run.id = null_flow_run_state_timestamp_cte.id; 

53 """ 

54 

55 update_task_run_state_timestamp_in_batches = """ 

56 WITH null_task_run_state_timestamp_cte as (SELECT id from task_run where state_timestamp is null and state_id is not null limit 500) 

57 UPDATE task_run 

58 SET state_timestamp = task_run_state.timestamp 

59 FROM task_run_state, null_task_run_state_timestamp_cte 

60 WHERE task_run.state_id = task_run_state.id 

61 AND task_run.id = null_task_run_state_timestamp_cte.id; 

62 """ 

63 

64 with op.get_context().autocommit_block(): 

65 conn = op.get_bind() 

66 while True: 

67 # execute until we've backfilled all flow run state timestamps 

68 # autocommit mode will commit each time `execute` is called 

69 result = conn.execute(sa.text(update_flow_run_state_timestamp_in_batches)) 

70 if result.rowcount <= 0: 70 ↛ 66line 70 didn't jump to line 66 because the condition on line 70 was always true

71 break 

72 

73 while True: 

74 # execute until we've backfilled all task run state timestamps 

75 # autocommit mode will commit each time `execute` is called 

76 result = conn.execute(sa.text(update_task_run_state_timestamp_in_batches)) 

77 if result.rowcount <= 0: 77 ↛ 73line 77 didn't jump to line 73 because the condition on line 77 was always true

78 break 

79 

80 

81def downgrade(): 

82 with op.batch_alter_table("task_run", schema=None) as batch_op: 

83 batch_op.drop_index("ix_task_run__state_timestamp") 

84 batch_op.drop_column("state_timestamp") 

85 

86 with op.batch_alter_table("flow_run", schema=None) as batch_op: 

87 batch_op.drop_index("ix_flow_run__state_timestamp") 

88 batch_op.drop_column("state_timestamp")