Add coordinator admin console foundation (system, jobs, metrics)
This commit is contained in:
@@ -0,0 +1,198 @@
|
||||
package postgres
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
sq "github.com/Masterminds/squirrel"
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
|
||||
)
|
||||
|
||||
// AdminReadRepo backs the coordinator admin console: paginated jobs, status
|
||||
// counters, metrics buckets and storage figures. Read-only.
|
||||
type AdminReadRepo struct{ pool *pgxpool.Pool }
|
||||
|
||||
func NewAdminReadRepo(pool *pgxpool.Pool) *AdminReadRepo { return &AdminReadRepo{pool: pool} }
|
||||
|
||||
var _ usecase.AdminReadRepository = (*AdminReadRepo)(nil)
|
||||
|
||||
func (r *AdminReadRepo) ListJobsPaginated(ctx context.Context, status string, limit, offset int) ([]domain.Job, int, error) {
|
||||
if limit < 1 || limit > 100 || offset < 0 {
|
||||
return nil, 0, domain.ErrInvalidInput
|
||||
}
|
||||
countQ := psql.Select("COUNT(*)").From("jobs")
|
||||
listQ := psql.Select(jobColumns...).From("jobs")
|
||||
if status != "" {
|
||||
countQ = countQ.Where(sq.Eq{"status": status})
|
||||
listQ = listQ.Where(sq.Eq{"status": status})
|
||||
}
|
||||
countSQL, args, err := countQ.ToSql()
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
var total int
|
||||
if err := conn(ctx, r.pool).QueryRow(ctx, countSQL, args...).Scan(&total); err != nil {
|
||||
return nil, 0, fmt.Errorf("count jobs: %w", err)
|
||||
}
|
||||
listSQL, args, err := listQ.OrderBy("created_at DESC", "id DESC").Limit(uint64(limit)).Offset(uint64(offset)).ToSql()
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
rows, err := conn(ctx, r.pool).Query(ctx, listSQL, args...)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("list jobs paginated: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
jobs := make([]domain.Job, 0)
|
||||
for rows.Next() {
|
||||
var j domain.Job
|
||||
var statusRaw string
|
||||
if err := rows.Scan(
|
||||
&j.ID, &j.Workload, &j.InputURI, &j.Parameters, &statusRaw, &j.CreatedAt, &j.CompletedAt,
|
||||
&j.InputArtifactID, &j.ResultArtifactID, &j.ErrorCode, &j.ErrorMessage, &j.ReducerStartedAt,
|
||||
&j.OwnerID,
|
||||
); err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
j.Status = domain.JobStatus(statusRaw)
|
||||
jobs = append(jobs, j)
|
||||
}
|
||||
return jobs, total, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) CountJobsByStatus(ctx context.Context) (map[string]int, error) {
|
||||
rows, err := conn(ctx, r.pool).Query(ctx, "SELECT status, COUNT(*) FROM jobs GROUP BY status")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("count jobs by status: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
out := make(map[string]int)
|
||||
for rows.Next() {
|
||||
var status string
|
||||
var count int
|
||||
if err := rows.Scan(&status, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[status] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) TaskCountsByJobs(ctx context.Context, jobIDs []uuid.UUID) (map[uuid.UUID]map[string]int, error) {
|
||||
out := make(map[uuid.UUID]map[string]int, len(jobIDs))
|
||||
if len(jobIDs) == 0 {
|
||||
return out, nil
|
||||
}
|
||||
sql, args, err := psql.Select("job_id", "status", "COUNT(*)").From("tasks").
|
||||
Where(sq.Eq{"job_id": jobIDs}).GroupBy("job_id", "status").ToSql()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
rows, err := conn(ctx, r.pool).Query(ctx, sql, args...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("task counts by jobs: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
for rows.Next() {
|
||||
var jobID uuid.UUID
|
||||
var status string
|
||||
var count int
|
||||
if err := rows.Scan(&jobID, &status, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out[jobID] == nil {
|
||||
out[jobID] = make(map[string]int)
|
||||
}
|
||||
out[jobID][status] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) JobCountsByDay(ctx context.Context, since time.Time) (map[string]int, error) {
|
||||
rows, err := conn(ctx, r.pool).Query(ctx,
|
||||
"SELECT to_char(date_trunc('day', created_at AT TIME ZONE 'UTC'), 'YYYY-MM-DD') AS day, COUNT(*) FROM jobs WHERE created_at >= $1 GROUP BY 1",
|
||||
since.UTC())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("job counts by day: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
out := make(map[string]int)
|
||||
for rows.Next() {
|
||||
var day string
|
||||
var count int
|
||||
if err := rows.Scan(&day, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[day] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) JobCountsByWorkload(ctx context.Context) (map[string]int, error) {
|
||||
rows, err := conn(ctx, r.pool).Query(ctx, "SELECT workload, COUNT(*) FROM jobs GROUP BY workload")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("job counts by workload: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
out := make(map[string]int)
|
||||
for rows.Next() {
|
||||
var workload string
|
||||
var count int
|
||||
if err := rows.Scan(&workload, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[workload] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) TaskStats(ctx context.Context) (int64, int64, float64, error) {
|
||||
var completed, failed int64
|
||||
var avgSeconds *float64
|
||||
err := conn(ctx, r.pool).QueryRow(ctx, `
|
||||
SELECT
|
||||
COALESCE(SUM(CASE WHEN status = 'completed' THEN 1 ELSE 0 END), 0),
|
||||
COALESCE(SUM(CASE WHEN status = 'failed' THEN 1 ELSE 0 END), 0),
|
||||
AVG(CASE WHEN status = 'completed' AND started_at IS NOT NULL
|
||||
THEN EXTRACT(EPOCH FROM completed_at - started_at) END)::float8
|
||||
FROM tasks`).Scan(&completed, &failed, &avgSeconds)
|
||||
if err != nil {
|
||||
return 0, 0, 0, fmt.Errorf("task stats: %w", err)
|
||||
}
|
||||
var avg float64
|
||||
if avgSeconds != nil {
|
||||
avg = *avgSeconds
|
||||
}
|
||||
return completed, failed, avg, nil
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) ArtifactSizeByKind(ctx context.Context) (map[string]int64, error) {
|
||||
rows, err := conn(ctx, r.pool).Query(ctx, "SELECT kind, COALESCE(SUM(size_bytes), 0) FROM artifacts GROUP BY kind")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("artifact sizes: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
out := make(map[string]int64)
|
||||
for rows.Next() {
|
||||
var kind string
|
||||
var size int64
|
||||
if err := rows.Scan(&kind, &size); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[kind] = size
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) DatabaseSizeBytes(ctx context.Context) (int64, error) {
|
||||
var size int64
|
||||
if err := conn(ctx, r.pool).QueryRow(ctx, "SELECT pg_database_size(current_database())").Scan(&size); err != nil {
|
||||
return 0, fmt.Errorf("database size: %w", err)
|
||||
}
|
||||
return size, nil
|
||||
}
|
||||
@@ -0,0 +1,191 @@
|
||||
package sqlite
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
|
||||
)
|
||||
|
||||
// AdminReadRepo backs the coordinator admin console: paginated jobs, status
|
||||
// counters, metrics buckets and storage figures. Read-only.
|
||||
type AdminReadRepo struct{ db *sql.DB }
|
||||
|
||||
func NewAdminReadRepo(db *sql.DB) *AdminReadRepo { return &AdminReadRepo{db: db} }
|
||||
|
||||
var _ usecase.AdminReadRepository = (*AdminReadRepo)(nil)
|
||||
|
||||
func (r *AdminReadRepo) ListJobsPaginated(ctx context.Context, status string, limit, offset int) ([]domain.Job, int, error) {
|
||||
if limit < 1 || limit > 100 || offset < 0 {
|
||||
return nil, 0, domain.ErrInvalidInput
|
||||
}
|
||||
where := ""
|
||||
args := []any{}
|
||||
if status != "" {
|
||||
where = " WHERE status = ?"
|
||||
args = append(args, status)
|
||||
}
|
||||
var total int
|
||||
if err := conn(ctx, r.db).QueryRowContext(ctx, "SELECT COUNT(*) FROM jobs"+where, args...).Scan(&total); err != nil {
|
||||
return nil, 0, fmt.Errorf("count jobs: %w", err)
|
||||
}
|
||||
rows, err := conn(ctx, r.db).QueryContext(ctx,
|
||||
"SELECT "+jobColumns+" FROM jobs"+where+" ORDER BY created_at DESC, id DESC LIMIT ? OFFSET ?",
|
||||
append(args, limit, offset)...)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("list jobs paginated: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
jobs := make([]domain.Job, 0)
|
||||
for rows.Next() {
|
||||
job, err := scanJob(rows)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
jobs = append(jobs, *job)
|
||||
}
|
||||
return jobs, total, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) CountJobsByStatus(ctx context.Context) (map[string]int, error) {
|
||||
rows, err := conn(ctx, r.db).QueryContext(ctx, "SELECT status, COUNT(*) FROM jobs GROUP BY status")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("count jobs by status: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
out := make(map[string]int)
|
||||
for rows.Next() {
|
||||
var status string
|
||||
var count int
|
||||
if err := rows.Scan(&status, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[status] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) TaskCountsByJobs(ctx context.Context, jobIDs []uuid.UUID) (map[uuid.UUID]map[string]int, error) {
|
||||
out := make(map[uuid.UUID]map[string]int, len(jobIDs))
|
||||
if len(jobIDs) == 0 {
|
||||
return out, nil
|
||||
}
|
||||
placeholders := make([]string, 0, len(jobIDs))
|
||||
args := make([]any, 0, len(jobIDs))
|
||||
for _, id := range jobIDs {
|
||||
placeholders = append(placeholders, "?")
|
||||
args = append(args, id.String())
|
||||
}
|
||||
rows, err := conn(ctx, r.db).QueryContext(ctx,
|
||||
"SELECT job_id, status, COUNT(*) FROM tasks WHERE job_id IN ("+strings.Join(placeholders, ", ")+") GROUP BY job_id, status",
|
||||
args...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("task counts by jobs: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
for rows.Next() {
|
||||
var jobRaw, status string
|
||||
var count int
|
||||
if err := rows.Scan(&jobRaw, &status, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
jobID, err := uuid.Parse(jobRaw)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("task counts: parse job id: %w", err)
|
||||
}
|
||||
if out[jobID] == nil {
|
||||
out[jobID] = make(map[string]int)
|
||||
}
|
||||
out[jobID][status] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) JobCountsByDay(ctx context.Context, since time.Time) (map[string]int, error) {
|
||||
// created_at is unix nanos; the bucket is the UTC calendar day.
|
||||
rows, err := conn(ctx, r.db).QueryContext(ctx,
|
||||
"SELECT strftime('%Y-%m-%d', created_at / 1000000000, 'unixepoch') AS day, COUNT(*) FROM jobs WHERE created_at >= ? GROUP BY day",
|
||||
since.UTC().UnixNano())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("job counts by day: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
out := make(map[string]int)
|
||||
for rows.Next() {
|
||||
var day string
|
||||
var count int
|
||||
if err := rows.Scan(&day, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[day] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) JobCountsByWorkload(ctx context.Context) (map[string]int, error) {
|
||||
rows, err := conn(ctx, r.db).QueryContext(ctx, "SELECT workload, COUNT(*) FROM jobs GROUP BY workload")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("job counts by workload: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
out := make(map[string]int)
|
||||
for rows.Next() {
|
||||
var workload string
|
||||
var count int
|
||||
if err := rows.Scan(&workload, &count); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[workload] = count
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) TaskStats(ctx context.Context) (int64, int64, float64, error) {
|
||||
var completed, failed int64
|
||||
var avgNanos sql.NullFloat64
|
||||
err := conn(ctx, r.db).QueryRowContext(ctx, `
|
||||
SELECT
|
||||
COALESCE(SUM(CASE WHEN status = 'completed' THEN 1 ELSE 0 END), 0),
|
||||
COALESCE(SUM(CASE WHEN status = 'failed' THEN 1 ELSE 0 END), 0),
|
||||
AVG(CASE WHEN status = 'completed' AND started_at IS NOT NULL THEN completed_at - started_at END)
|
||||
FROM tasks`).Scan(&completed, &failed, &avgNanos)
|
||||
if err != nil {
|
||||
return 0, 0, 0, fmt.Errorf("task stats: %w", err)
|
||||
}
|
||||
return completed, failed, avgNanos.Float64 / 1e9, nil
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) ArtifactSizeByKind(ctx context.Context) (map[string]int64, error) {
|
||||
rows, err := conn(ctx, r.db).QueryContext(ctx, "SELECT kind, COALESCE(SUM(size_bytes), 0) FROM artifacts GROUP BY kind")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("artifact sizes: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
out := make(map[string]int64)
|
||||
for rows.Next() {
|
||||
var kind string
|
||||
var size int64
|
||||
if err := rows.Scan(&kind, &size); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[kind] = size
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *AdminReadRepo) DatabaseSizeBytes(ctx context.Context) (int64, error) {
|
||||
var pageCount, pageSize int64
|
||||
if err := conn(ctx, r.db).QueryRowContext(ctx, "PRAGMA page_count").Scan(&pageCount); err != nil {
|
||||
return 0, fmt.Errorf("page count: %w", err)
|
||||
}
|
||||
if err := conn(ctx, r.db).QueryRowContext(ctx, "PRAGMA page_size").Scan(&pageSize); err != nil {
|
||||
return 0, fmt.Errorf("page size: %w", err)
|
||||
}
|
||||
return pageCount * pageSize, nil
|
||||
}
|
||||
@@ -0,0 +1,171 @@
|
||||
package sqlite
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
|
||||
)
|
||||
|
||||
func TestAdminListJobsPaginatedAndCounts(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
ctx := context.Background()
|
||||
jobRepo := NewJobRepo(db)
|
||||
adminRepo := NewAdminReadRepo(db)
|
||||
|
||||
jobs := make([]*domain.Job, 5)
|
||||
for i := range jobs {
|
||||
jobs[i] = seedJob(t, db, 2)
|
||||
}
|
||||
// Two completed, two running, one pending.
|
||||
if err := jobRepo.UpdateStatus(ctx, jobs[0].ID, domain.JobCompleted, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := jobRepo.UpdateStatus(ctx, jobs[1].ID, domain.JobCompleted, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := jobRepo.UpdateStatus(ctx, jobs[2].ID, domain.JobRunning, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := jobRepo.UpdateStatus(ctx, jobs[3].ID, domain.JobRunning, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
all, total, err := adminRepo.ListJobsPaginated(ctx, "", 100, 0)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if total != 5 || len(all) != 5 {
|
||||
t.Errorf("all: total=%d len=%d, want 5/5", total, len(all))
|
||||
}
|
||||
completed, total, err := adminRepo.ListJobsPaginated(ctx, "completed", 100, 0)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if total != 2 || len(completed) != 2 {
|
||||
t.Errorf("completed: total=%d len=%d, want 2/2", total, len(completed))
|
||||
}
|
||||
page, total, err := adminRepo.ListJobsPaginated(ctx, "", 2, 2)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if total != 5 || len(page) != 2 {
|
||||
t.Errorf("page: total=%d len=%d, want 5/2", total, len(page))
|
||||
}
|
||||
|
||||
counts, err := adminRepo.CountJobsByStatus(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if counts["completed"] != 2 || counts["running"] != 2 || counts["pending"] != 1 {
|
||||
t.Errorf("counts = %v, want completed=2 running=2 pending=1", counts)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAdminTaskCountsByJobs(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
ctx := context.Background()
|
||||
job := seedJob(t, db, 3)
|
||||
if _, err := db.ExecContext(ctx, "UPDATE tasks SET status = 'completed', result_artifact_id = ? WHERE chunk_index = 0 AND job_id = ?", uuid.NewString(), job.ID.String()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := db.ExecContext(ctx, "UPDATE tasks SET status = 'failed' WHERE chunk_index = 1 AND job_id = ?", job.ID.String()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
counts, err := NewAdminReadRepo(db).TaskCountsByJobs(ctx, []uuid.UUID{job.ID})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got := counts[job.ID]
|
||||
if got["completed"] != 1 || got["failed"] != 1 || got["pending"] != 1 {
|
||||
t.Errorf("task counts = %v, want completed=1 failed=1 pending=1", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAdminJobCountsByDay(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
ctx := context.Background()
|
||||
job := seedJob(t, db, 1)
|
||||
// Move the seed job to two days ago; create two more today.
|
||||
old := fixedTime().Add(-48 * time.Hour)
|
||||
if _, err := db.ExecContext(ctx, "UPDATE jobs SET created_at = ? WHERE id = ?", old.UnixNano(), job.ID.String()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
seedJob(t, db, 1)
|
||||
seedJob(t, db, 1)
|
||||
|
||||
repo := NewAdminReadRepo(db)
|
||||
counts, err := repo.JobCountsByDay(ctx, fixedTime().Add(-6*24*time.Hour))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
today := fixedTime().UTC().Format("2006-01-02")
|
||||
oldDay := old.UTC().Format("2006-01-02")
|
||||
if counts[today] != 2 {
|
||||
t.Errorf("today count = %d, want 2 (got %v)", counts[today], counts)
|
||||
}
|
||||
if counts[oldDay] != 1 {
|
||||
t.Errorf("old day count = %d, want 1 (got %v)", counts[oldDay], counts)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAdminTaskStatsAndStorage(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
ctx := context.Background()
|
||||
repo := NewAdminReadRepo(db)
|
||||
|
||||
// One completed task with a known duration, one failed.
|
||||
job := seedJob(t, db, 2)
|
||||
start := fixedTime().Add(-2 * time.Minute)
|
||||
done := fixedTime().Add(-90 * time.Second)
|
||||
queries := []string{
|
||||
"UPDATE tasks SET status='completed', result_artifact_id=?, started_at=?, completed_at=? WHERE job_id=? AND chunk_index=0",
|
||||
"UPDATE tasks SET status='failed' WHERE job_id=? AND chunk_index=1",
|
||||
}
|
||||
for i, q := range queries {
|
||||
args := []any{uuid.NewString(), start.UnixNano(), done.UnixNano(), job.ID.String()}
|
||||
if i == 1 {
|
||||
args = []any{job.ID.String()}
|
||||
}
|
||||
if _, err := db.ExecContext(ctx, q, args...); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
completed, failed, avg, err := repo.TaskStats(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if completed != 1 || failed != 1 {
|
||||
t.Errorf("stats = completed %d failed %d, want 1/1", completed, failed)
|
||||
}
|
||||
if avg < 29 || avg > 31 {
|
||||
t.Errorf("avg duration = %.1fs, want ~30s", avg)
|
||||
}
|
||||
|
||||
// Artifact sizes by kind.
|
||||
for _, kind := range []string{"input", "shard", "final_result"} {
|
||||
if _, err := db.ExecContext(ctx, "INSERT INTO artifacts (id, job_id, kind, filename, storage_key, content_type, size_bytes, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
uuid.NewString(), job.ID.String(), kind, kind+".csv", "key-"+kind, "text/csv", int64(len(kind)*1000), fixedTime().UnixNano()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
sizes, err := repo.ArtifactSizeByKind(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if sizes["input"] != 5000 || sizes["shard"] != 5000 || sizes["final_result"] != 12000 {
|
||||
t.Errorf("sizes = %v", sizes)
|
||||
}
|
||||
dbBytes, err := repo.DatabaseSizeBytes(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if dbBytes <= 0 {
|
||||
t.Errorf("database size = %d, want > 0", dbBytes)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user