(ctx context.Context, engine storepb.Engine, db *sql.DB, appName string)
| 609 | } |
| 610 | |
| 611 | func getSession(ctx context.Context, engine storepb.Engine, db *sql.DB, appName string) (*v1pb.TaskRunSession, error) { |
| 612 | switch engine { |
| 613 | case storepb.Engine_POSTGRES, storepb.Engine_COCKROACHDB: |
| 614 | query := ` |
| 615 | WITH target_session AS ( |
| 616 | SELECT pid FROM pg_catalog.pg_stat_activity WHERE application_name = $1 LIMIT 1 |
| 617 | ) |
| 618 | SELECT |
| 619 | a.pid, |
| 620 | pg_blocking_pids(a.pid) AS blocked_by_pids, |
| 621 | a.query, |
| 622 | a.state, |
| 623 | a.wait_event_type, |
| 624 | a.wait_event, |
| 625 | a.datname, |
| 626 | a.usename, |
| 627 | a.application_name, |
| 628 | a.client_addr, |
| 629 | a.client_port, |
| 630 | a.backend_start, |
| 631 | a.xact_start, |
| 632 | a.query_start |
| 633 | FROM |
| 634 | pg_catalog.pg_stat_activity a |
| 635 | WHERE a.application_name = $1 |
| 636 | OR (SELECT pid FROM target_session) = ANY(pg_blocking_pids(a.pid)) |
| 637 | OR a.pid = ANY(pg_blocking_pids((SELECT pid FROM target_session))) |
| 638 | ORDER BY a.pid |
| 639 | ` |
| 640 | rows, err := db.QueryContext(ctx, query, appName) |
| 641 | if err != nil { |
| 642 | return nil, errors.Wrapf(err, "failed to query rows") |
| 643 | } |
| 644 | defer rows.Close() |
| 645 | |
| 646 | ss := &v1pb.TaskRunSession_Postgres{} |
| 647 | for rows.Next() { |
| 648 | var s v1pb.TaskRunSession_Postgres_Session |
| 649 | |
| 650 | var blockedByPids pgtype.TextArray |
| 651 | |
| 652 | var bs time.Time |
| 653 | var xs, qs *time.Time |
| 654 | if err := rows.Scan( |
| 655 | &s.Pid, |
| 656 | &blockedByPids, |
| 657 | &s.Query, |
| 658 | &s.State, |
| 659 | &s.WaitEventType, |
| 660 | &s.WaitEvent, |
| 661 | &s.Datname, |
| 662 | &s.Usename, |
| 663 | &s.ApplicationName, |
| 664 | &s.ClientAddr, |
| 665 | &s.ClientPort, |
| 666 | &bs, |
| 667 | &xs, |
| 668 | &qs, |
no test coverage detected