Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/database/_migrations/versions/postgresql/2021_01_20_122127_25f4b90a7a42_initial_migration.py: 55%

132 statements  

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

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")