feat(store): tasks, accounts, runs, dedup journal
This commit is contained in:
@@ -0,0 +1,74 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
)
|
||||
|
||||
type Account struct {
|
||||
ID int64
|
||||
TaskID int64
|
||||
SrcLogin string
|
||||
SrcPassEnc string
|
||||
DstLogin string
|
||||
DstPassEnc string
|
||||
TestSrcStatus string
|
||||
TestDstStatus string
|
||||
Status string
|
||||
Copied int64
|
||||
Skipped int64
|
||||
Errors int64
|
||||
}
|
||||
|
||||
func (s *Store) CreateAccount(ctx context.Context, a Account) (int64, error) {
|
||||
var id int64
|
||||
err := s.Pool.QueryRow(ctx,
|
||||
`INSERT INTO accounts (task_id, src_login, src_pass_enc, dst_login, dst_pass_enc)
|
||||
VALUES ($1,$2,$3,$4,$5) RETURNING id`,
|
||||
a.TaskID, a.SrcLogin, a.SrcPassEnc, a.DstLogin, a.DstPassEnc).Scan(&id)
|
||||
return id, err
|
||||
}
|
||||
|
||||
func (s *Store) ListAccountsByTask(ctx context.Context, taskID int64) ([]Account, error) {
|
||||
rows, err := s.Pool.Query(ctx,
|
||||
`SELECT id, task_id, src_login, src_pass_enc, dst_login, dst_pass_enc,
|
||||
test_src_status, test_dst_status, status, copied_count, skipped_count, error_count
|
||||
FROM accounts WHERE task_id=$1 ORDER BY id`, taskID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []Account
|
||||
for rows.Next() {
|
||||
var a Account
|
||||
if err := rows.Scan(&a.ID, &a.TaskID, &a.SrcLogin, &a.SrcPassEnc, &a.DstLogin, &a.DstPassEnc,
|
||||
&a.TestSrcStatus, &a.TestDstStatus, &a.Status, &a.Copied, &a.Skipped, &a.Errors); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, a)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// side = "src" | "dst"
|
||||
func (s *Store) SetAccountTestStatus(ctx context.Context, id int64, side, status string) error {
|
||||
col := "test_src_status"
|
||||
if side == "dst" {
|
||||
col = "test_dst_status"
|
||||
}
|
||||
_, err := s.Pool.Exec(ctx, fmt.Sprintf(`UPDATE accounts SET %s=$2 WHERE id=$1`, col), id, status)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *Store) SetAccountStatus(ctx context.Context, id int64, status string) error {
|
||||
_, err := s.Pool.Exec(ctx, `UPDATE accounts SET status=$2 WHERE id=$1`, id, status)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *Store) IncAccountCounters(ctx context.Context, id, copied, skipped, errs int64) error {
|
||||
_, err := s.Pool.Exec(ctx,
|
||||
`UPDATE accounts SET copied_count=copied_count+$2,
|
||||
skipped_count=skipped_count+$3, error_count=error_count+$4 WHERE id=$1`,
|
||||
id, copied, skipped, errs)
|
||||
return err
|
||||
}
|
||||
Reference in New Issue
Block a user