Files
move-data-nas/internal/models/job.go
T
darroyo bfa006f4ab Add SSH key management, job history persistence, and live streaming
- SSH key management: generate ed25519 keypairs or import public keys
  from UI (/ssh-keys), per-machine key selection in Machines form,
  one-time private key download with hash verification
- Fix engine to use machine-specific SSH key (was hardcoded to server key)
- Job log persistence: write to job_logs table (DB) with batched inserts,
  buffer of 50 lines; GetAllFiltered with status/pair/date range filters
- EventBus refactor: per-job subscriber channels, global channel, non-blocking
- SSE endpoints: /jobs/stream (all), /jobs/:id/log/stream (per-job live)
- JobDetail page: live log streaming, auto-scroll, cancel, duration
- JobHistory: filters (pair, status, date range), pagination, link to detail
- Cleanup scheduler: daily purge of job_logs and finished jobs older than
  SYNCSERVER_RETENTION_DAYS (default 30)
- Migration 0002: indexes on job_logs(job_id), jobs(status,created_at),
  jobs(sync_pair_id)
2026-07-07 20:36:11 -04:00

154 lines
4.1 KiB
Go

package models
import (
"database/sql"
"time"
)
type Job struct {
ID int64 `db:"id" json:"id"`
SyncPairID int64 `db:"sync_pair_id" json:"sync_pair_id"`
TriggerType string `db:"trigger_type" json:"trigger_type"`
Status string `db:"status" json:"status"`
StartedAt *time.Time `db:"started_at" json:"started_at"`
FinishedAt *time.Time `db:"finished_at" json:"finished_at"`
LogFile *string `db:"log_file" json:"log_file"`
CreatedAt time.Time `db:"created_at" json:"created_at"`
}
type JobRepository struct {
db *sql.DB
}
func NewJobRepository(db *sql.DB) *JobRepository {
return &JobRepository{db: db}
}
func (r *JobRepository) Create(syncPairID int64, triggerType, status string) (int64, error) {
res, err := r.db.Exec(`
INSERT INTO jobs (sync_pair_id, trigger_type, status) VALUES (?, ?, ?)`,
syncPairID, triggerType, status,
)
if err != nil {
return 0, err
}
return res.LastInsertId()
}
func (r *JobRepository) GetByID(id int64) (*Job, error) {
var j Job
var started, finished sql.NullTime
var logFile sql.NullString
err := r.db.QueryRow(`
SELECT id, sync_pair_id, trigger_type, status, started_at, finished_at,
log_file, created_at FROM jobs WHERE id = ?`, id).Scan(
&j.ID, &j.SyncPairID, &j.TriggerType, &j.Status, &started, &finished,
&logFile, &j.CreatedAt)
if err != nil {
return nil, err
}
if started.Valid {
j.StartedAt = &started.Time
}
if finished.Valid {
j.FinishedAt = &finished.Time
}
if logFile.Valid {
j.LogFile = &logFile.String
}
return &j, nil
}
func (r *JobRepository) GetAll(limit, offset int) ([]Job, error) {
rows, err := r.db.Query(`
SELECT id, sync_pair_id, trigger_type, status, started_at, finished_at,
log_file, created_at FROM jobs ORDER BY created_at DESC LIMIT ? OFFSET ?`,
limit, offset)
if err != nil {
return nil, err
}
defer rows.Close()
var jobs []Job
for rows.Next() {
var j Job
var started, finished sql.NullTime
var logFile sql.NullString
if err := rows.Scan(&j.ID, &j.SyncPairID, &j.TriggerType, &j.Status,
&started, &finished, &logFile, &j.CreatedAt); err != nil {
return nil, err
}
if started.Valid {
j.StartedAt = &started.Time
}
if finished.Valid {
j.FinishedAt = &finished.Time
}
if logFile.Valid {
j.LogFile = &logFile.String
}
jobs = append(jobs, j)
}
return jobs, rows.Err()
}
func (r *JobRepository) UpdateStatus(id int64, status string) error {
var query string
var args []interface{}
switch status {
case "running", "waking_up":
query = "UPDATE jobs SET status = ?, started_at = COALESCE(started_at, CURRENT_TIMESTAMP) WHERE id = ?"
args = []interface{}{status, id}
case "success", "failed", "cancelled":
query = "UPDATE jobs SET status = ?, finished_at = CURRENT_TIMESTAMP WHERE id = ?"
args = []interface{}{status, id}
default:
query = "UPDATE jobs SET status = ? WHERE id = ?"
args = []interface{}{status, id}
}
_, err := r.db.Exec(query, args...)
return err
}
func (r *JobRepository) SetLogFile(id int64, path string) error {
_, err := r.db.Exec("UPDATE jobs SET log_file = ? WHERE id = ?", path, id)
return err
}
func (r *JobRepository) GetRunningBySyncPair(syncPairID int64) (*Job, error) {
var j Job
var started sql.NullTime
var logFile sql.NullString
err := r.db.QueryRow(`
SELECT id, sync_pair_id, trigger_type, status, started_at, finished_at,
log_file, created_at FROM jobs
WHERE sync_pair_id = ? AND status IN ('queued','waking_up','running')
ORDER BY created_at DESC LIMIT 1`, syncPairID).Scan(
&j.ID, &j.SyncPairID, &j.TriggerType, &j.Status, &started,
&j.FinishedAt, &logFile, &j.CreatedAt)
if err != nil {
return nil, err
}
if started.Valid {
j.StartedAt = &started.Time
}
if logFile.Valid {
j.LogFile = &logFile.String
}
return &j, nil
}
func (r *JobRepository) Count() (int64, error) {
var n int64
err := r.db.QueryRow("SELECT COUNT(*) FROM jobs").Scan(&n)
return n, err
}
func (r *JobRepository) DeleteFinishedBefore(before time.Time) (int64, error) {
res, err := r.db.Exec("DELETE FROM jobs WHERE finished_at IS NOT NULL AND finished_at < ?", before)
if err != nil {
return 0, err
}
return res.RowsAffected()
}