1"""initial migration
2
3Revision ID: 25f4b90a7a42
4Revises:
5Create Date: 2022-01-20 12:21:27.508018
6
7"""
8
9from typing import Dict, List, Union
10
11import sqlalchemy as sa
12from alembic import op
13from sqlalchemy import Text
14
15import prefect
16from prefect.server.utilities.schemas import PrefectBaseModel
17
18
19class DataDocument(PrefectBaseModel):
20 """
21 DataDocuments were deprecated in September 2022 and this stub is included here
22 to simplify removal from the library.
23 """
24
25 encoding: str
26 blob: bytes
27
28
29# revision identifiers, used by Alembic.
30revision = "25f4b90a7a42"
31down_revision = None
32branch_labels = None
33depends_on = None
34
35
36def upgrade():
37 # Create tables
38 op.create_table(
39 "flow",
40 sa.Column(
41 "id",
42 prefect.server.utilities.database.UUID(),
43 server_default=sa.text("(GEN_RANDOM_UUID())"),
44 nullable=False,
45 ),
46 sa.Column(
47 "created",
48 prefect.server.utilities.database.Timestamp(timezone=True),
49 server_default=sa.text("CURRENT_TIMESTAMP"),
50 nullable=False,
51 ),
52 sa.Column(
53 "updated",
54 prefect.server.utilities.database.Timestamp(timezone=True),
55 server_default=sa.text("CURRENT_TIMESTAMP"),
56 nullable=False,
57 ),
58 sa.Column("name", sa.String(), nullable=False),
59 sa.Column(
60 "tags",
61 prefect.server.utilities.database.JSON(astext_type=Text()),
62 server_default="[]",
63 nullable=False,
64 ),
65 sa.PrimaryKeyConstraint("id", name=op.f("pk_flow")),
66 sa.UniqueConstraint("name", name=op.f("uq_flow__name")),
67 )
68 op.create_index(op.f("ix_flow__updated"), "flow", ["updated"], unique=False)
69 op.create_table(
70 "log",
71 sa.Column(
72 "id",
73 prefect.server.utilities.database.UUID(),
74 server_default=sa.text("(GEN_RANDOM_UUID())"),
75 nullable=False,
76 ),
77 sa.Column(
78 "created",
79 prefect.server.utilities.database.Timestamp(timezone=True),
80 server_default=sa.text("CURRENT_TIMESTAMP"),
81 nullable=False,
82 ),
83 sa.Column(
84 "updated",
85 prefect.server.utilities.database.Timestamp(timezone=True),
86 server_default=sa.text("CURRENT_TIMESTAMP"),
87 nullable=False,
88 ),
89 sa.Column("name", sa.String(), nullable=False),
90 sa.Column("level", sa.SmallInteger(), nullable=False),
91 sa.Column(
92 "flow_run_id", prefect.server.utilities.database.UUID(), nullable=False
93 ),
94 sa.Column(
95 "task_run_id", prefect.server.utilities.database.UUID(), nullable=True
96 ),
97 sa.Column("message", sa.Text(), nullable=False),
98 sa.Column(
99 "timestamp",
100 prefect.server.utilities.database.Timestamp(timezone=True),
101 nullable=False,
102 ),
103 sa.PrimaryKeyConstraint("id", name=op.f("pk_log")),
104 )
105 op.create_index(op.f("ix_log__flow_run_id"), "log", ["flow_run_id"], unique=False)
106 op.create_index(op.f("ix_log__level"), "log", ["level"], unique=False)
107 op.create_index(op.f("ix_log__task_run_id"), "log", ["task_run_id"], unique=False)
108 op.create_index(op.f("ix_log__timestamp"), "log", ["timestamp"], unique=False)
109 op.create_index(op.f("ix_log__updated"), "log", ["updated"], unique=False)
110 op.create_table(
111 "concurrency_limit",
112 sa.Column(
113 "id",
114 prefect.server.utilities.database.UUID(),
115 server_default=sa.text("(GEN_RANDOM_UUID())"),
116 nullable=False,
117 ),
118 sa.Column(
119 "created",
120 prefect.server.utilities.database.Timestamp(timezone=True),
121 server_default=sa.text("CURRENT_TIMESTAMP"),
122 nullable=False,
123 ),
124 sa.Column(
125 "updated",
126 prefect.server.utilities.database.Timestamp(timezone=True),
127 server_default=sa.text("CURRENT_TIMESTAMP"),
128 nullable=False,
129 ),
130 sa.Column("tag", sa.String(), nullable=False),
131 sa.Column("concurrency_limit", sa.Integer(), nullable=False),
132 sa.Column(
133 "active_slots",
134 prefect.server.utilities.database.JSON(astext_type=Text()),
135 server_default="[]",
136 nullable=False,
137 ),
138 sa.PrimaryKeyConstraint("id", name=op.f("pk_concurrency_limit")),
139 )
140 op.create_index(
141 op.f("ix_concurrency_limit__tag"), "concurrency_limit", ["tag"], unique=True
142 )
143 op.create_index(
144 op.f("ix_concurrency_limit__updated"),
145 "concurrency_limit",
146 ["updated"],
147 unique=False,
148 )
149 op.create_table(
150 "saved_search",
151 sa.Column(
152 "id",
153 prefect.server.utilities.database.UUID(),
154 server_default=sa.text("(GEN_RANDOM_UUID())"),
155 nullable=False,
156 ),
157 sa.Column(
158 "created",
159 prefect.server.utilities.database.Timestamp(timezone=True),
160 server_default=sa.text("CURRENT_TIMESTAMP"),
161 nullable=False,
162 ),
163 sa.Column(
164 "updated",
165 prefect.server.utilities.database.Timestamp(timezone=True),
166 server_default=sa.text("CURRENT_TIMESTAMP"),
167 nullable=False,
168 ),
169 sa.Column("name", sa.String(), nullable=False),
170 sa.Column(
171 "filters",
172 prefect.server.utilities.database.JSON(astext_type=Text()),
173 server_default="[]",
174 nullable=False,
175 ),
176 sa.PrimaryKeyConstraint("id", name=op.f("pk_saved_search")),
177 sa.UniqueConstraint("name", name=op.f("uq_saved_search__name")),
178 )
179 op.create_index(
180 op.f("ix_saved_search__updated"), "saved_search", ["updated"], unique=False
181 )
182 op.create_table(
183 "task_run_state_cache",
184 sa.Column(
185 "id",
186 prefect.server.utilities.database.UUID(),
187 server_default=sa.text("(GEN_RANDOM_UUID())"),
188 nullable=False,
189 ),
190 sa.Column(
191 "created",
192 prefect.server.utilities.database.Timestamp(timezone=True),
193 server_default=sa.text("CURRENT_TIMESTAMP"),
194 nullable=False,
195 ),
196 sa.Column(
197 "updated",
198 prefect.server.utilities.database.Timestamp(timezone=True),
199 server_default=sa.text("CURRENT_TIMESTAMP"),
200 nullable=False,
201 ),
202 sa.Column("cache_key", sa.String(), nullable=False),
203 sa.Column(
204 "cache_expiration",
205 prefect.server.utilities.database.Timestamp(timezone=True),
206 nullable=True,
207 ),
208 sa.Column(
209 "task_run_state_id",
210 prefect.server.utilities.database.UUID(),
211 nullable=False,
212 ),
213 sa.PrimaryKeyConstraint("id", name=op.f("pk_task_run_state_cache")),
214 )
215 op.create_index(
216 op.f("ix_task_run_state_cache__updated"),
217 "task_run_state_cache",
218 ["updated"],
219 unique=False,
220 )
221 op.create_table(
222 "deployment",
223 sa.Column(
224 "id",
225 prefect.server.utilities.database.UUID(),
226 server_default=sa.text("(GEN_RANDOM_UUID())"),
227 nullable=False,
228 ),
229 sa.Column(
230 "created",
231 prefect.server.utilities.database.Timestamp(timezone=True),
232 server_default=sa.text("CURRENT_TIMESTAMP"),
233 nullable=False,
234 ),
235 sa.Column(
236 "updated",
237 prefect.server.utilities.database.Timestamp(timezone=True),
238 server_default=sa.text("CURRENT_TIMESTAMP"),
239 nullable=False,
240 ),
241 sa.Column("name", sa.String(), nullable=False),
242 sa.Column(
243 "schedule",
244 prefect.server.utilities.database.Pydantic(
245 prefect.server.schemas.schedules.SCHEDULE_TYPES
246 ),
247 nullable=True,
248 ),
249 sa.Column(
250 "is_schedule_active", sa.Boolean(), server_default="1", nullable=False
251 ),
252 sa.Column(
253 "tags",
254 prefect.server.utilities.database.JSON(astext_type=Text()),
255 server_default="[]",
256 nullable=False,
257 ),
258 sa.Column(
259 "parameters",
260 prefect.server.utilities.database.JSON(astext_type=Text()),
261 server_default="{}",
262 nullable=False,
263 ),
264 sa.Column(
265 "flow_data",
266 prefect.server.utilities.database.Pydantic(DataDocument),
267 nullable=True,
268 ),
269 sa.Column("flow_runner_type", sa.String(), nullable=True),
270 sa.Column(
271 "flow_runner_config",
272 prefect.server.utilities.database.JSON(astext_type=Text()),
273 nullable=True,
274 ),
275 sa.Column("flow_id", prefect.server.utilities.database.UUID(), nullable=False),
276 sa.ForeignKeyConstraint(
277 ["flow_id"],
278 ["flow.id"],
279 name=op.f("fk_deployment__flow_id__flow"),
280 ondelete="CASCADE",
281 ),
282 sa.PrimaryKeyConstraint("id", name=op.f("pk_deployment")),
283 )
284 op.create_index(
285 op.f("ix_deployment__flow_id"), "deployment", ["flow_id"], unique=False
286 )
287 op.create_index(
288 op.f("ix_deployment__updated"), "deployment", ["updated"], unique=False
289 )
290 op.create_index(
291 "uq_deployment__flow_id_name", "deployment", ["flow_id", "name"], unique=True
292 )
293 op.create_table(
294 "flow_run",
295 sa.Column(
296 "id",
297 prefect.server.utilities.database.UUID(),
298 server_default=sa.text("(GEN_RANDOM_UUID())"),
299 nullable=False,
300 ),
301 sa.Column(
302 "created",
303 prefect.server.utilities.database.Timestamp(timezone=True),
304 server_default=sa.text("CURRENT_TIMESTAMP"),
305 nullable=False,
306 ),
307 sa.Column(
308 "updated",
309 prefect.server.utilities.database.Timestamp(timezone=True),
310 server_default=sa.text("CURRENT_TIMESTAMP"),
311 nullable=False,
312 ),
313 sa.Column("name", sa.String(), nullable=False),
314 sa.Column(
315 "state_type",
316 sa.Enum(
317 "SCHEDULED",
318 "PENDING",
319 "RUNNING",
320 "COMPLETED",
321 "FAILED",
322 "CANCELLED",
323 name="state_type",
324 ),
325 nullable=True,
326 ),
327 sa.Column("run_count", sa.Integer(), server_default="0", nullable=False),
328 sa.Column(
329 "expected_start_time",
330 prefect.server.utilities.database.Timestamp(timezone=True),
331 nullable=True,
332 ),
333 sa.Column(
334 "next_scheduled_start_time",
335 prefect.server.utilities.database.Timestamp(timezone=True),
336 nullable=True,
337 ),
338 sa.Column(
339 "start_time",
340 prefect.server.utilities.database.Timestamp(timezone=True),
341 nullable=True,
342 ),
343 sa.Column(
344 "end_time",
345 prefect.server.utilities.database.Timestamp(timezone=True),
346 nullable=True,
347 ),
348 sa.Column("total_run_time", sa.Interval(), server_default="0", nullable=False),
349 sa.Column("flow_version", sa.String(), nullable=True),
350 sa.Column(
351 "parameters",
352 prefect.server.utilities.database.JSON(astext_type=Text()),
353 server_default="{}",
354 nullable=False,
355 ),
356 sa.Column("idempotency_key", sa.String(), nullable=True),
357 sa.Column(
358 "context",
359 prefect.server.utilities.database.JSON(astext_type=Text()),
360 server_default="{}",
361 nullable=False,
362 ),
363 sa.Column(
364 "empirical_policy",
365 prefect.server.utilities.database.JSON(astext_type=Text()),
366 server_default="{}",
367 nullable=False,
368 ),
369 sa.Column(
370 "tags",
371 prefect.server.utilities.database.JSON(astext_type=Text()),
372 server_default="[]",
373 nullable=False,
374 ),
375 sa.Column("flow_runner_type", sa.String(), nullable=True),
376 sa.Column(
377 "flow_runner_config",
378 prefect.server.utilities.database.JSON(astext_type=Text()),
379 nullable=True,
380 ),
381 sa.Column(
382 "empirical_config",
383 prefect.server.utilities.database.JSON(astext_type=Text()),
384 server_default="{}",
385 nullable=False,
386 ),
387 sa.Column("auto_scheduled", sa.Boolean(), server_default="0", nullable=False),
388 sa.Column("flow_id", prefect.server.utilities.database.UUID(), nullable=False),
389 sa.Column(
390 "deployment_id", prefect.server.utilities.database.UUID(), nullable=True
391 ),
392 sa.Column(
393 "parent_task_run_id",
394 prefect.server.utilities.database.UUID(),
395 nullable=True,
396 ),
397 sa.Column("state_id", prefect.server.utilities.database.UUID(), nullable=True),
398 sa.ForeignKeyConstraint(
399 ["deployment_id"],
400 ["deployment.id"],
401 name=op.f("fk_flow_run__deployment_id__deployment"),
402 ondelete="set null",
403 ),
404 sa.ForeignKeyConstraint(
405 ["flow_id"],
406 ["flow.id"],
407 name=op.f("fk_flow_run__flow_id__flow"),
408 ondelete="cascade",
409 ),
410 sa.ForeignKeyConstraint(
411 ["parent_task_run_id"],
412 ["task_run.id"],
413 name=op.f("fk_flow_run__parent_task_run_id__task_run"),
414 ondelete="SET NULL",
415 use_alter=True,
416 ),
417 sa.ForeignKeyConstraint(
418 ["state_id"],
419 ["flow_run_state.id"],
420 name=op.f("fk_flow_run__state_id__flow_run_state"),
421 ondelete="SET NULL",
422 use_alter=True,
423 ),
424 sa.PrimaryKeyConstraint("id", name=op.f("pk_flow_run")),
425 )
426 op.create_index(
427 op.f("ix_flow_run__deployment_id"), "flow_run", ["deployment_id"], unique=False
428 )
429 op.create_index(op.f("ix_flow_run__flow_id"), "flow_run", ["flow_id"], unique=False)
430 op.create_index(
431 op.f("ix_flow_run__flow_version"), "flow_run", ["flow_version"], unique=False
432 )
433 op.create_index(op.f("ix_flow_run__name"), "flow_run", ["name"], unique=False)
434 op.create_index(
435 op.f("ix_flow_run__parent_task_run_id"),
436 "flow_run",
437 ["parent_task_run_id"],
438 unique=False,
439 )
440 op.create_index("ix_flow_run__start_time", "flow_run", ["start_time"], unique=False)
441 op.create_index(
442 op.f("ix_flow_run__state_id"), "flow_run", ["state_id"], unique=False
443 )
444 op.create_index("ix_flow_run__state_type", "flow_run", ["state_type"], unique=False)
445 op.create_index(op.f("ix_flow_run__updated"), "flow_run", ["updated"], unique=False)
446 op.create_index(
447 "uq_flow_run__flow_id_idempotency_key",
448 "flow_run",
449 ["flow_id", "idempotency_key"],
450 unique=True,
451 )
452 op.create_table(
453 "flow_run_state",
454 sa.Column(
455 "id",
456 prefect.server.utilities.database.UUID(),
457 server_default=sa.text("(GEN_RANDOM_UUID())"),
458 nullable=False,
459 ),
460 sa.Column(
461 "created",
462 prefect.server.utilities.database.Timestamp(timezone=True),
463 server_default=sa.text("CURRENT_TIMESTAMP"),
464 nullable=False,
465 ),
466 sa.Column(
467 "updated",
468 prefect.server.utilities.database.Timestamp(timezone=True),
469 server_default=sa.text("CURRENT_TIMESTAMP"),
470 nullable=False,
471 ),
472 sa.Column(
473 "type",
474 sa.Enum(
475 "SCHEDULED",
476 "PENDING",
477 "RUNNING",
478 "COMPLETED",
479 "FAILED",
480 "CANCELLED",
481 name="state_type",
482 ),
483 nullable=False,
484 ),
485 sa.Column(
486 "timestamp",
487 prefect.server.utilities.database.Timestamp(timezone=True),
488 server_default=sa.text("CURRENT_TIMESTAMP"),
489 nullable=False,
490 ),
491 sa.Column("name", sa.String(), nullable=False),
492 sa.Column("message", sa.String(), nullable=True),
493 sa.Column(
494 "state_details",
495 prefect.server.utilities.database.Pydantic(
496 prefect.server.schemas.states.StateDetails
497 ),
498 server_default="{}",
499 nullable=False,
500 ),
501 sa.Column(
502 "data",
503 prefect.server.utilities.database.Pydantic(DataDocument),
504 nullable=True,
505 ),
506 sa.Column(
507 "flow_run_id", prefect.server.utilities.database.UUID(), nullable=False
508 ),
509 sa.ForeignKeyConstraint(
510 ["flow_run_id"],
511 ["flow_run.id"],
512 name=op.f("fk_flow_run_state__flow_run_id__flow_run"),
513 ondelete="cascade",
514 ),
515 sa.PrimaryKeyConstraint("id", name=op.f("pk_flow_run_state")),
516 )
517 op.create_index(
518 op.f("ix_flow_run_state__name"), "flow_run_state", ["name"], unique=False
519 )
520 op.create_index(
521 op.f("ix_flow_run_state__type"), "flow_run_state", ["type"], unique=False
522 )
523 op.create_index(
524 op.f("ix_flow_run_state__updated"), "flow_run_state", ["updated"], unique=False
525 )
526 op.create_table(
527 "task_run",
528 sa.Column(
529 "id",
530 prefect.server.utilities.database.UUID(),
531 server_default=sa.text("(GEN_RANDOM_UUID())"),
532 nullable=False,
533 ),
534 sa.Column(
535 "created",
536 prefect.server.utilities.database.Timestamp(timezone=True),
537 server_default=sa.text("CURRENT_TIMESTAMP"),
538 nullable=False,
539 ),
540 sa.Column(
541 "updated",
542 prefect.server.utilities.database.Timestamp(timezone=True),
543 server_default=sa.text("CURRENT_TIMESTAMP"),
544 nullable=False,
545 ),
546 sa.Column("name", sa.String(), nullable=False),
547 sa.Column(
548 "state_type",
549 sa.Enum(
550 "SCHEDULED",
551 "PENDING",
552 "RUNNING",
553 "COMPLETED",
554 "FAILED",
555 "CANCELLED",
556 name="state_type",
557 ),
558 nullable=True,
559 ),
560 sa.Column("run_count", sa.Integer(), server_default="0", nullable=False),
561 sa.Column(
562 "expected_start_time",
563 prefect.server.utilities.database.Timestamp(timezone=True),
564 nullable=True,
565 ),
566 sa.Column(
567 "next_scheduled_start_time",
568 prefect.server.utilities.database.Timestamp(timezone=True),
569 nullable=True,
570 ),
571 sa.Column(
572 "start_time",
573 prefect.server.utilities.database.Timestamp(timezone=True),
574 nullable=True,
575 ),
576 sa.Column(
577 "end_time",
578 prefect.server.utilities.database.Timestamp(timezone=True),
579 nullable=True,
580 ),
581 sa.Column("total_run_time", sa.Interval(), server_default="0", nullable=False),
582 sa.Column("task_key", sa.String(), nullable=False),
583 sa.Column("dynamic_key", sa.String(), nullable=False),
584 sa.Column("cache_key", sa.String(), nullable=True),
585 sa.Column(
586 "cache_expiration",
587 prefect.server.utilities.database.Timestamp(timezone=True),
588 nullable=True,
589 ),
590 sa.Column("task_version", sa.String(), nullable=True),
591 sa.Column(
592 "empirical_policy",
593 prefect.server.utilities.database.Pydantic(
594 prefect.server.schemas.core.TaskRunPolicy
595 ),
596 server_default="{}",
597 nullable=False,
598 ),
599 sa.Column(
600 "task_inputs",
601 prefect.server.utilities.database.Pydantic(
602 Dict[
603 str,
604 List[
605 Union[
606 prefect.server.schemas.core.TaskRunResult,
607 prefect.server.schemas.core.Parameter,
608 prefect.server.schemas.core.Constant,
609 ]
610 ],
611 ]
612 ),
613 server_default="{}",
614 nullable=False,
615 ),
616 sa.Column(
617 "tags",
618 prefect.server.utilities.database.JSON(astext_type=Text()),
619 server_default="[]",
620 nullable=False,
621 ),
622 sa.Column(
623 "flow_run_id", prefect.server.utilities.database.UUID(), nullable=False
624 ),
625 sa.Column("state_id", prefect.server.utilities.database.UUID(), nullable=True),
626 sa.ForeignKeyConstraint(
627 ["flow_run_id"],
628 ["flow_run.id"],
629 name=op.f("fk_task_run__flow_run_id__flow_run"),
630 ondelete="cascade",
631 ),
632 sa.ForeignKeyConstraint(
633 ["state_id"],
634 ["task_run_state.id"],
635 name=op.f("fk_task_run__state_id__task_run_state"),
636 ondelete="SET NULL",
637 use_alter=True,
638 ),
639 sa.PrimaryKeyConstraint("id", name=op.f("pk_task_run")),
640 )
641 op.create_index(
642 op.f("ix_task_run__flow_run_id"), "task_run", ["flow_run_id"], unique=False
643 )
644 op.create_index(op.f("ix_task_run__name"), "task_run", ["name"], unique=False)
645 op.create_index("ix_task_run__start_time", "task_run", ["start_time"], unique=False)
646 op.create_index(
647 op.f("ix_task_run__state_id"), "task_run", ["state_id"], unique=False
648 )
649 op.create_index("ix_task_run__state_type", "task_run", ["state_type"], unique=False)
650 op.create_index(op.f("ix_task_run__updated"), "task_run", ["updated"], unique=False)
651 op.create_index(
652 "uq_task_run__flow_run_id_task_key_dynamic_key",
653 "task_run",
654 ["flow_run_id", "task_key", "dynamic_key"],
655 unique=True,
656 )
657 op.create_table(
658 "task_run_state",
659 sa.Column(
660 "id",
661 prefect.server.utilities.database.UUID(),
662 server_default=sa.text("(GEN_RANDOM_UUID())"),
663 nullable=False,
664 ),
665 sa.Column(
666 "created",
667 prefect.server.utilities.database.Timestamp(timezone=True),
668 server_default=sa.text("CURRENT_TIMESTAMP"),
669 nullable=False,
670 ),
671 sa.Column(
672 "updated",
673 prefect.server.utilities.database.Timestamp(timezone=True),
674 server_default=sa.text("CURRENT_TIMESTAMP"),
675 nullable=False,
676 ),
677 sa.Column(
678 "type",
679 sa.Enum(
680 "SCHEDULED",
681 "PENDING",
682 "RUNNING",
683 "COMPLETED",
684 "FAILED",
685 "CANCELLED",
686 name="state_type",
687 ),
688 nullable=False,
689 ),
690 sa.Column(
691 "timestamp",
692 prefect.server.utilities.database.Timestamp(timezone=True),
693 server_default=sa.text("CURRENT_TIMESTAMP"),
694 nullable=False,
695 ),
696 sa.Column("name", sa.String(), nullable=False),
697 sa.Column("message", sa.String(), nullable=True),
698 sa.Column(
699 "state_details",
700 prefect.server.utilities.database.Pydantic(
701 prefect.server.schemas.states.StateDetails
702 ),
703 server_default="{}",
704 nullable=False,
705 ),
706 sa.Column(
707 "data",
708 prefect.server.utilities.database.Pydantic(DataDocument),
709 nullable=True,
710 ),
711 sa.Column(
712 "task_run_id", prefect.server.utilities.database.UUID(), nullable=False
713 ),
714 sa.ForeignKeyConstraint(
715 ["task_run_id"],
716 ["task_run.id"],
717 name=op.f("fk_task_run_state__task_run_id__task_run"),
718 ondelete="cascade",
719 ),
720 sa.PrimaryKeyConstraint("id", name=op.f("pk_task_run_state")),
721 )
722 op.create_index(
723 op.f("ix_task_run_state__name"), "task_run_state", ["name"], unique=False
724 )
725 op.create_index(
726 op.f("ix_task_run_state__type"), "task_run_state", ["type"], unique=False
727 )
728 op.create_index(
729 op.f("ix_task_run_state__updated"), "task_run_state", ["updated"], unique=False
730 )
731
732 # Foreign Keys do not automatically get created when `use_alter` is set, need to create manually
733 op.create_foreign_key(
734 constraint_name="fk_flow_run__parent_task_run_id__task_run",
735 source_table="flow_run",
736 referent_table="task_run",
737 local_cols=["parent_task_run_id"],
738 remote_cols=["id"],
739 ondelete="SET NULL",
740 )
741 op.create_foreign_key(
742 constraint_name="fk_flow_run__state_id__flow_run_state",
743 source_table="flow_run",
744 referent_table="flow_run_state",
745 local_cols=["state_id"],
746 remote_cols=["id"],
747 ondelete="SET NULL",
748 )
749 op.create_foreign_key(
750 constraint_name="fk_task_run__state_id__task_run_state",
751 source_table="task_run",
752 referent_table="task_run_state",
753 local_cols=["state_id"],
754 remote_cols=["id"],
755 ondelete="SET NULL",
756 )
757
758 # functional ordered indexes are skipped by auto generation and need to be created manually
759 op.create_index(
760 "ix_flow_run__end_time_desc",
761 "flow_run",
762 [sa.text("end_time DESC")],
763 unique=False,
764 )
765 op.create_index(
766 "ix_flow_run__expected_start_time_desc",
767 "flow_run",
768 [sa.text("expected_start_time DESC")],
769 unique=False,
770 )
771 op.create_index(
772 "ix_flow_run__next_scheduled_start_time_asc",
773 "flow_run",
774 [sa.text("next_scheduled_start_time ASC")],
775 unique=False,
776 )
777 op.create_index(
778 "uq_flow_run_state__flow_run_id_timestamp_desc",
779 "flow_run_state",
780 ["flow_run_id", sa.text("timestamp DESC")],
781 unique=True,
782 )
783 op.create_index(
784 "ix_task_run__expected_start_time_desc",
785 "task_run",
786 [sa.text("expected_start_time DESC")],
787 unique=False,
788 )
789 op.create_index(
790 "ix_task_run__next_scheduled_start_time_asc",
791 "task_run",
792 [sa.text("next_scheduled_start_time ASC")],
793 unique=False,
794 )
795 op.create_index(
796 "ix_task_run__end_time_desc",
797 "task_run",
798 [sa.text("end_time DESC")],
799 unique=False,
800 )
801 op.create_index(
802 "uq_task_run_state__task_run_id_timestamp_desc",
803 "task_run_state",
804 ["task_run_id", sa.text("timestamp DESC")],
805 unique=True,
806 )
807 op.create_index(
808 "ix_task_run_state_cache__cache_key_created_desc",
809 "task_run_state_cache",
810 ["cache_key", sa.text("created DESC")],
811 unique=False,
812 )
813
814
815def downgrade():
816 # functional ordered indexes are skipped by auto generation and need to be dropped manually
817 op.drop_index(
818 "ix_task_run_state_cache__cache_key_created_desc",
819 table_name="task_run_state_cache",
820 )
821 op.drop_index(
822 "uq_task_run_state__task_run_id_timestamp_desc", table_name="task_run_state"
823 )
824 op.drop_index("ix_task_run__next_scheduled_start_time_asc", table_name="task_run")
825 op.drop_index("ix_task_run__expected_start_time_desc", table_name="task_run")
826 op.drop_index("ix_task_run__end_time_desc", table_name="task_run")
827 op.drop_index(
828 "uq_flow_run_state__flow_run_id_timestamp_desc", table_name="flow_run_state"
829 )
830 op.drop_index("ix_flow_run__next_scheduled_start_time_asc", table_name="flow_run")
831 op.drop_index("ix_flow_run__expected_start_time_desc", table_name="flow_run")
832 op.drop_index("ix_flow_run__end_time_desc", table_name="flow_run")
833
834 # Foreign Keys do not automatically get created when `use_alter` is set, need to drop manually created keys
835 op.drop_constraint(
836 constraint_name="fk_flow_run__parent_task_run_id__task_run",
837 table_name="flow_run",
838 )
839 op.drop_constraint(
840 constraint_name="fk_flow_run__state_id__flow_run_state", table_name="flow_run"
841 )
842 op.drop_constraint(
843 constraint_name="fk_task_run__state_id__task_run_state", table_name="task_run"
844 )
845
846 # Drop tables
847 op.drop_index(op.f("ix_task_run_state__updated"), table_name="task_run_state")
848 op.drop_index(op.f("ix_task_run_state__type"), table_name="task_run_state")
849 op.drop_index(op.f("ix_task_run_state__name"), table_name="task_run_state")
850 op.drop_table("task_run_state")
851 op.drop_index(
852 "uq_task_run__flow_run_id_task_key_dynamic_key", table_name="task_run"
853 )
854 op.drop_index(op.f("ix_task_run__updated"), table_name="task_run")
855 op.drop_index("ix_task_run__state_type", table_name="task_run")
856 op.drop_index(op.f("ix_task_run__state_id"), table_name="task_run")
857 op.drop_index("ix_task_run__start_time", table_name="task_run")
858 op.drop_index(op.f("ix_task_run__name"), table_name="task_run")
859 op.drop_index(op.f("ix_task_run__flow_run_id"), table_name="task_run")
860 op.drop_table("task_run")
861 op.drop_index(op.f("ix_flow_run_state__updated"), table_name="flow_run_state")
862 op.drop_index(op.f("ix_flow_run_state__type"), table_name="flow_run_state")
863 op.drop_index(op.f("ix_flow_run_state__name"), table_name="flow_run_state")
864 op.drop_table("flow_run_state")
865 op.drop_index("uq_flow_run__flow_id_idempotency_key", table_name="flow_run")
866 op.drop_index(op.f("ix_flow_run__updated"), table_name="flow_run")
867 op.drop_index("ix_flow_run__state_type", table_name="flow_run")
868 op.drop_index(op.f("ix_flow_run__state_id"), table_name="flow_run")
869 op.drop_index("ix_flow_run__start_time", table_name="flow_run")
870 op.drop_index(op.f("ix_flow_run__parent_task_run_id"), table_name="flow_run")
871 op.drop_index(op.f("ix_flow_run__name"), table_name="flow_run")
872 op.drop_index(op.f("ix_flow_run__flow_version"), table_name="flow_run")
873 op.drop_index(op.f("ix_flow_run__flow_id"), table_name="flow_run")
874 op.drop_index(op.f("ix_flow_run__deployment_id"), table_name="flow_run")
875 op.drop_table("flow_run")
876 op.drop_index("uq_deployment__flow_id_name", table_name="deployment")
877 op.drop_index(op.f("ix_deployment__updated"), table_name="deployment")
878 op.drop_index(op.f("ix_deployment__flow_id"), table_name="deployment")
879 op.drop_table("deployment")
880 op.drop_index(
881 op.f("ix_task_run_state_cache__updated"), table_name="task_run_state_cache"
882 )
883 op.drop_table("task_run_state_cache")
884 op.drop_index(op.f("ix_concurrency_limit__updated"), table_name="concurrency_limit")
885 op.drop_index(op.f("ix_concurrency_limit__tag"), table_name="concurrency_limit")
886 op.drop_table("concurrency_limit")
887 op.drop_index(op.f("ix_saved_search__updated"), table_name="saved_search")
888 op.drop_table("saved_search")
889 op.drop_index(op.f("ix_log__updated"), table_name="log")
890 op.drop_index(op.f("ix_log__timestamp"), table_name="log")
891 op.drop_index(op.f("ix_log__task_run_id"), table_name="log")
892 op.drop_index(op.f("ix_log__level"), table_name="log")
893 op.drop_index(op.f("ix_log__flow_run_id"), table_name="log")
894 op.drop_table("log")
895 op.drop_index(op.f("ix_flow__updated"), table_name="flow")
896 op.drop_table("flow")
897
898 # Enum Type is not dropped by default
899 op.execute("DROP TYPE state_type")