diff --git a/Makefile b/Makefile index 88ade92..85a4e2f 100644 --- a/Makefile +++ b/Makefile @@ -1,5 +1,5 @@ BINARY=syncserver -VERSION?=1.0.4 +VERSION?=1.0.5 GO?=go LDFLAGS=-s -w -X main.version=$(VERSION) -X main.commit=$(shell git rev-parse --short HEAD 2>/dev/null || echo unknown) BUILD_FLAGS=CGO_ENABLED=0 diff --git a/cmd/server/main.go b/cmd/server/main.go index 158ed50..b4092b1 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -20,7 +20,7 @@ import ( "github.com/syncserver/internal/syncengine" ) -var version = "1.0.4" +var version = "1.0.5" func main() { cfgPath := flag.String("config", "", "Path to config.yaml") diff --git a/internal/api/dto.go b/internal/api/dto.go index 8870e88..5293161 100644 --- a/internal/api/dto.go +++ b/internal/api/dto.go @@ -55,13 +55,39 @@ type SyncPairResponse struct { } type JobResponse struct { - ID int64 `json:"id"` - SyncPairID int64 `json:"sync_pair_id"` - TriggerType string `json:"trigger_type"` - Status string `json:"status"` - StartedAt *string `json:"started_at"` - FinishedAt *string `json:"finished_at"` - LogFile *string `json:"log_file"` + ID int64 `json:"id"` + SyncPairID int64 `json:"sync_pair_id"` + TriggerType string `json:"trigger_type"` + Status string `json:"status"` + StartedAt *string `json:"started_at"` + FinishedAt *string `json:"finished_at"` + LogFile *string `json:"log_file"` + DurationSeconds *int64 `json:"duration_seconds,omitempty"` + LogLineCount *int64 `json:"log_line_count,omitempty"` +} + +type LogLineResponse struct { + ID int64 `json:"id"` + JobID int64 `json:"job_id"` + Stream string `json:"stream"` + Content string `json:"content"` + Timestamp string `json:"timestamp"` +} + +type SSHKeyRequest struct { + Label string `json:"label"` + Generate bool `json:"generate"` + PublicKey string `json:"public_key"` +} + +type SSHKeyResponse struct { + ID int64 `json:"id"` + Label string `json:"label"` + PublicKey string `json:"public_key"` + Fingerprint string `json:"fingerprint"` + InUse bool `json:"in_use"` + HasPrivateKey bool `json:"has_private_key"` + CreatedAt string `json:"created_at"` } type ErrorResponse struct { diff --git a/internal/api/handlers_jobs.go b/internal/api/handlers_jobs.go index 46b286f..03b98a7 100644 --- a/internal/api/handlers_jobs.go +++ b/internal/api/handlers_jobs.go @@ -2,6 +2,7 @@ package api import ( "database/sql" + "fmt" "net/http" "os" "strconv" @@ -28,8 +29,30 @@ func (h *JobHandler) List(w http.ResponseWriter, r *http.Request) { limit = 50 } - repo := models.NewJobRepository(h.db) - jobs, err := repo.GetAll(limit, offset) + var syncPairID *int64 + if spidStr := r.URL.Query().Get("sync_pair_id"); spidStr != "" { + if spid, err := strconv.ParseInt(spidStr, 10, 64); err == nil { + syncPairID = &spid + } + } + + status := r.URL.Query().Get("status") + triggerType := r.URL.Query().Get("trigger_type") + + var from, to *time.Time + if fromStr := r.URL.Query().Get("from"); fromStr != "" { + if t, err := time.Parse(time.RFC3339, fromStr); err == nil { + from = &t + } + } + if toStr := r.URL.Query().Get("to"); toStr != "" { + if t, err := time.Parse(time.RFC3339, toStr); err == nil { + to = &t + } + } + + repo := models.NewJobLogRepository(h.db) + jobs, total, err := repo.GetAllFiltered(limit, offset, syncPairID, status, triggerType, from, to) if err != nil { writeError(w, http.StatusInternalServerError, "failed to fetch jobs") return @@ -37,8 +60,9 @@ func (h *JobHandler) List(w http.ResponseWriter, r *http.Request) { out := make([]JobResponse, len(jobs)) for i, j := range jobs { - out[i] = jobToResp(j) + out[i] = jobWithStatsToResp(j) } + w.Header().Set("X-Total-Count", fmt.Sprintf("%d", total)) writeJSON(w, out) } @@ -120,30 +144,59 @@ func (h *JobHandler) TriggerRun(w http.ResponseWriter, r *http.Request) { writeJSON(w, jobToResp(*j), http.StatusCreated) } -func (h *JobHandler) StreamLog(w http.ResponseWriter, r *http.Request) { +func (h *JobHandler) GetLog(w http.ResponseWriter, r *http.Request) { id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64) if err != nil { writeError(w, http.StatusBadRequest, "invalid id") return } - flusher, ok := w.(http.Flusher) - if !ok { - writeError(w, http.StatusInternalServerError, "streaming not supported") + offset, _ := strconv.Atoi(r.URL.Query().Get("offset")) + limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) + if limit <= 0 { + limit = 1000 + } + + logRepo := models.NewJobLogRepository(h.db) + logs, err := logRepo.GetByJobID(id, limit, offset) + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to fetch logs") return } - w.Header().Set("Content-Type", "text/plain") - w.Header().Set("Cache-Control", "no-cache") - w.Header().Set("Connection", "keep-alive") - flusher.Flush() + count, _ := logRepo.CountByJobID(id) + w.Header().Set("X-Total-Count", fmt.Sprintf("%d", count)) + writeJSON(w, logs) +} + +func (h *JobHandler) DownloadLog(w http.ResponseWriter, r *http.Request) { + id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64) + if err != nil { + writeError(w, http.StatusBadRequest, "invalid id") + return + } jobRepo := models.NewJobRepository(h.db) j, err := jobRepo.GetByID(id) - if err == nil && j.LogFile != nil { - data, _ := os.ReadFile(*j.LogFile) - w.Write(data) - flusher.Flush() + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to fetch job") + return + } + + if j.LogFile != nil { + data, err := os.ReadFile(*j.LogFile) + if err == nil { + w.Header().Set("Content-Type", "text/plain") + w.Header().Set("Content-Disposition", fmt.Sprintf(`attachment; filename="job-%d.log"`, id)) + w.Write(data) + return + } + } + + logRepo := models.NewJobLogRepository(h.db) + logs, _ := logRepo.GetByJobID(id, 100000, 0) + for _, l := range logs { + fmt.Fprintf(w, "[%s] %s\n", l.Timestamp.Format(time.RFC3339), l.Content) } } @@ -165,3 +218,10 @@ func jobToResp(j models.Job) JobResponse { } return resp } + +func jobWithStatsToResp(j models.JobWithStats) JobResponse { + resp := jobToResp(j.Job) + resp.DurationSeconds = j.DurationSeconds + resp.LogLineCount = &j.LogLineCount + return resp +} diff --git a/internal/api/handlers_sshkeys.go b/internal/api/handlers_sshkeys.go new file mode 100644 index 0000000..e6e99da --- /dev/null +++ b/internal/api/handlers_sshkeys.go @@ -0,0 +1,223 @@ +package api + +import ( + "crypto/sha256" + "database/sql" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "path/filepath" + "strconv" + "strings" + + "github.com/go-chi/chi/v5" + "github.com/syncserver/internal/config" + "github.com/syncserver/internal/models" + "github.com/syncserver/internal/sshmanager" +) + +type SSHKeyHandler struct { + db *sql.DB + cfg *config.Config +} + +func NewSSHKeyHandler(db *sql.DB, cfg *config.Config) *SSHKeyHandler { + return &SSHKeyHandler{db: db, cfg: cfg} +} + +func (h *SSHKeyHandler) List(w http.ResponseWriter, r *http.Request) { + repo := models.NewSSHKeyRepository(h.db) + machineRepo := models.NewMachineRepository(h.db) + keys, err := repo.GetAll() + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to fetch ssh keys") + return + } + machines, _ := machineRepo.GetAll() + + out := make([]SSHKeyResponse, len(keys)) + for i, k := range keys { + inUse := false + for _, m := range machines { + if m.SSHKeyID != nil && *m.SSHKeyID == k.ID { + inUse = true + break + } + } + hasPriv := false + if _, err := os.Stat(k.PrivateKeyPath); err == nil { + hasPriv = true + } + fp, _ := sshmanager.Fingerprint(k.PublicKey) + out[i] = SSHKeyResponse{ + ID: k.ID, + Label: k.Label, + PublicKey: k.PublicKey, + Fingerprint: fp, + InUse: inUse, + HasPrivateKey: hasPriv, + CreatedAt: k.CreatedAt.Format("2006-01-02T15:04:05Z07:00"), + } + } + writeJSON(w, out) +} + +func (h *SSHKeyHandler) Create(w http.ResponseWriter, r *http.Request) { + var req SSHKeyRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "invalid request body") + return + } + if req.Label == "" { + writeError(w, http.StatusBadRequest, "label is required") + return + } + if !req.Generate && req.PublicKey == "" { + writeError(w, http.StatusBadRequest, "either generate=true or public_key is required") + return + } + + keysDir := filepath.Join(h.cfg.SSHDir(), "keys") + if err := os.MkdirAll(keysDir, 0700); err != nil { + writeError(w, http.StatusInternalServerError, "failed to create keys directory") + return + } + + var pubKey, privPath, fp string + var err error + + if req.Generate { + privPath, _, pubKey, fp, err = sshmanager.GenerateKeyPair(req.Label, keysDir) + if err != nil { + writeError(w, http.StatusInternalServerError, fmt.Sprintf("generating key: %v", err)) + return + } + } else { + pubKey = strings.TrimSpace(req.PublicKey) + fp, err = sshmanager.Fingerprint(pubKey) + if err != nil { + writeError(w, http.StatusBadRequest, "invalid public key format") + return + } + privPath = "" + } + + repo := models.NewSSHKeyRepository(h.db) + id, err := repo.Create(req.Label, privPath, pubKey) + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to store ssh key") + return + } + w.Header().Set("Location", "/api/ssh-keys/"+strconv.FormatInt(id, 10)) + writeJSON(w, SSHKeyResponse{ + ID: id, + Label: req.Label, + PublicKey: pubKey, + Fingerprint: fp, + InUse: false, + HasPrivateKey: privPath != "", + }, http.StatusCreated) +} + +func (h *SSHKeyHandler) Get(w http.ResponseWriter, r *http.Request) { + id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64) + if err != nil { + writeError(w, http.StatusBadRequest, "invalid id") + return + } + repo := models.NewSSHKeyRepository(h.db) + k, err := repo.GetByID(id) + if err == sql.ErrNoRows { + writeError(w, http.StatusNotFound, "ssh key not found") + return + } + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to fetch ssh key") + return + } + machineRepo := models.NewMachineRepository(h.db) + machines, _ := machineRepo.GetAll() + inUse := false + for _, m := range machines { + if m.SSHKeyID != nil && *m.SSHKeyID == k.ID { + inUse = true + break + } + } + hasPriv := false + if _, err := os.Stat(k.PrivateKeyPath); err == nil { + hasPriv = true + } + fp, _ := sshmanager.Fingerprint(k.PublicKey) + writeJSON(w, SSHKeyResponse{ + ID: k.ID, + Label: k.Label, + PublicKey: k.PublicKey, + Fingerprint: fp, + InUse: inUse, + HasPrivateKey: hasPriv, + CreatedAt: k.CreatedAt.Format("2006-01-02T15:04:05Z07:00"), + }) +} + +func (h *SSHKeyHandler) Delete(w http.ResponseWriter, r *http.Request) { + id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64) + if err != nil { + writeError(w, http.StatusBadRequest, "invalid id") + return + } + repo := models.NewSSHKeyRepository(h.db) + machineRepo := models.NewMachineRepository(h.db) + machines, _ := machineRepo.GetAll() + for _, m := range machines { + if m.SSHKeyID != nil && *m.SSHKeyID == id { + writeError(w, http.StatusConflict, "ssh key is in use by machines") + return + } + } + k, err := repo.GetByID(id) + if err == nil && k.PrivateKeyPath != "" { + os.Remove(k.PrivateKeyPath) + os.Remove(k.PrivateKeyPath + ".pub") + } + if err := repo.Delete(id); err != nil { + writeError(w, http.StatusInternalServerError, "failed to delete ssh key") + return + } + w.WriteHeader(http.StatusNoContent) +} + +func (h *SSHKeyHandler) DownloadPrivate(w http.ResponseWriter, r *http.Request) { + id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64) + if err != nil { + writeError(w, http.StatusBadRequest, "invalid id") + return + } + repo := models.NewSSHKeyRepository(h.db) + k, err := repo.GetByID(id) + if err == sql.ErrNoRows { + writeError(w, http.StatusNotFound, "ssh key not found") + return + } + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to fetch ssh key") + return + } + if k.PrivateKeyPath == "" { + writeError(w, http.StatusNotFound, "no private key available for this entry") + return + } + + data, err := os.ReadFile(k.PrivateKeyPath) + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to read private key") + return + } + + w.Header().Set("Content-Type", "application/octet-stream") + w.Header().Set("Content-Disposition", fmt.Sprintf(`attachment; filename="%s.key"`, k.Label)) + w.Header().Set("X-Private-Key-Hash", fmt.Sprintf("sha256:%x", sha256.Sum256(data))) + io.WriteString(w, string(data)) +} diff --git a/internal/api/handlers_ws.go b/internal/api/handlers_ws.go index bcbdd7c..97211f8 100644 --- a/internal/api/handlers_ws.go +++ b/internal/api/handlers_ws.go @@ -18,11 +18,51 @@ func NewSSEHandler(engine *syncengine.Engine) *SSEHandler { return &SSEHandler{engine: engine} } -func (h *SSEHandler) Stream(w http.ResponseWriter, r *http.Request) { +func (h *SSEHandler) StreamAll(w http.ResponseWriter, r *http.Request) { + flusher, ok := w.(http.Flusher) + if !ok { + http.Error(w, "SSE not supported", http.StatusInternalServerError) + return + } + + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-cache") + w.Header().Set("Connection", "keep-alive") + w.Header().Set("X-Accel-Buffering", "no") + flusher.Flush() + + if h.engine == nil { + return + } + + events, unsub := h.engine.SubscribeGlobal() + defer unsub() + + for { + select { + case evt := <-events: + data, _ := json.Marshal(evt) + fmt.Fprintf(w, "event: %s\ndata: %s\n\n", evt.Type, data) + flusher.Flush() + case <-r.Context().Done(): + return + case <-time.After(30 * time.Second): + fmt.Fprintf(w, ": keepalive\n\n") + flusher.Flush() + } + } +} + +func (h *SSEHandler) StreamJob(w http.ResponseWriter, r *http.Request) { jobIDStr := r.URL.Query().Get("job_id") - var filterJobID int64 - if jobIDStr != "" { - filterJobID, _ = strconv.ParseInt(jobIDStr, 10, 64) + if jobIDStr == "" { + http.Error(w, "job_id required", http.StatusBadRequest) + return + } + jobID, err := strconv.ParseInt(jobIDStr, 10, 64) + if err != nil { + http.Error(w, "invalid job_id", http.StatusBadRequest) + return } flusher, ok := w.(http.Flusher) @@ -35,27 +75,23 @@ func (h *SSEHandler) Stream(w http.ResponseWriter, r *http.Request) { w.Header().Set("Cache-Control", "no-cache") w.Header().Set("Connection", "keep-alive") w.Header().Set("X-Accel-Buffering", "no") - flusher.Flush() if h.engine == nil { return } - events := h.engine.Events() + events, unsub := h.engine.SubscribeJob(jobID) + defer unsub() + for { select { case evt := <-events: - if filterJobID != 0 && evt.JobID != filterJobID { - continue - } data, _ := json.Marshal(evt) fmt.Fprintf(w, "event: %s\ndata: %s\n\n", evt.Type, data) flusher.Flush() - case <-r.Context().Done(): return - case <-time.After(30 * time.Second): fmt.Fprintf(w, ": keepalive\n\n") flusher.Flush() diff --git a/internal/api/router.go b/internal/api/router.go index ec5e55f..56c0088 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -35,6 +35,7 @@ func NewServer(cfg *config.Config, db *sql.DB, engine *syncengine.Engine) *Serve syncPairHandler := NewSyncPairHandler(db) jobHandler := NewJobHandler(db, engine) sseHandler := NewSSEHandler(engine) + sshKeyHandler := NewSSHKeyHandler(db, cfg) r.Route("/api", func(r chi.Router) { r.Route("/auth", func(r chi.Router) { @@ -64,16 +65,26 @@ func NewServer(cfg *config.Config, db *sql.DB, engine *syncengine.Engine) *Serve r.Get("/", jobHandler.List) r.Get("/{id}", jobHandler.Get) r.Post("/{id}/cancel", jobHandler.Cancel) - r.Get("/{id}/log", jobHandler.StreamLog) + r.Get("/{id}/log", jobHandler.GetLog) + r.Get("/{id}/log/download", jobHandler.DownloadLog) + r.Get("/{id}/log/stream", sseHandler.StreamJob) }) - r.With(auth.RequireAuth).Get("/jobs/stream", sseHandler.Stream) + r.With(auth.RequireAuth).Get("/jobs/stream", sseHandler.StreamAll) r.With(auth.RequireAuth).Get("/settings/pubkey", func(w http.ResponseWriter, r *http.Request) { _, _, pubKey, _ := sshmanager.EnsureServerKey(cfg.SSHDir()) w.Header().Set("Content-Type", "text/plain") w.Write([]byte(pubKey)) }) + + r.With(auth.RequireAuth).Route("/ssh-keys", func(r chi.Router) { + r.Get("/", sshKeyHandler.List) + r.Post("/", sshKeyHandler.Create) + r.Get("/{id}", sshKeyHandler.Get) + r.Delete("/{id}", sshKeyHandler.Delete) + r.Get("/{id}/private", sshKeyHandler.DownloadPrivate) + }) }) r.Get("/health", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { diff --git a/internal/config/config.go b/internal/config/config.go index 092c548..1435afd 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -28,7 +28,8 @@ type AuthConfig struct { } type SchedulerConfig struct { - Timezone string `yaml:"timezone" env:"SYNCSERVER_SCHEDULER_TZ" default:"UTC"` + Timezone string `yaml:"timezone" env:"SYNCSERVER_SCHEDULER_TZ" default:"UTC"` + RetentionDays int `yaml:"retention_days" env:"SYNCSERVER_RETENTION_DAYS" default:"30"` } var globalCfg *Config @@ -42,7 +43,8 @@ func Load(configPath, dataDir, addr string) (*Config, error) { JWTExpiryH: 24, }, Scheduler: SchedulerConfig{ - Timezone: "UTC", + Timezone: "UTC", + RetentionDays: 30, }, } diff --git a/internal/db/migrations/0002_job_history.sql b/internal/db/migrations/0002_job_history.sql new file mode 100644 index 0000000..f3089f1 --- /dev/null +++ b/internal/db/migrations/0002_job_history.sql @@ -0,0 +1,13 @@ +-- 0002_job_history.sql + +CREATE INDEX IF NOT EXISTS idx_job_logs_job_id ON job_logs(job_id); +CREATE INDEX IF NOT EXISTS idx_jobs_status_created ON jobs(status, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_jobs_sync_pair_id ON jobs(sync_pair_id); + +CREATE TABLE IF NOT EXISTS cleanup_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + deleted_before DATETIME NOT NULL, + logs_purged INTEGER NOT NULL DEFAULT 0, + jobs_purged INTEGER NOT NULL DEFAULT 0, + executed_at DATETIME DEFAULT CURRENT_TIMESTAMP +); diff --git a/internal/models/job.go b/internal/models/job.go index 8b6f7b1..4716ca3 100644 --- a/internal/models/job.go +++ b/internal/models/job.go @@ -143,3 +143,11 @@ func (r *JobRepository) Count() (int64, error) { 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() +} diff --git a/internal/models/job_log.go b/internal/models/job_log.go new file mode 100644 index 0000000..c0bf307 --- /dev/null +++ b/internal/models/job_log.go @@ -0,0 +1,168 @@ +package models + +import ( + "database/sql" + "strings" + "time" +) + +type JobLog struct { + ID int64 `db:"id" json:"id"` + JobID int64 `db:"job_id" json:"job_id"` + Stream string `db:"stream" json:"stream"` + Content string `db:"content" json:"content"` + Timestamp time.Time `db:"timestamp" json:"timestamp"` +} + +type JobLogRepository struct { + db *sql.DB +} + +func NewJobLogRepository(db *sql.DB) *JobLogRepository { + return &JobLogRepository{db: db} +} + +func (r *JobLogRepository) InsertBatch(jobID int64, stream string, lines []string) error { + if len(lines) == 0 { + return nil + } + tx, err := r.db.Begin() + if err != nil { + return err + } + stmt, err := tx.Prepare("INSERT INTO job_logs (job_id, stream, content) VALUES (?, ?, ?)") + if err != nil { + return err + } + defer stmt.Close() + for _, line := range lines { + if _, err := stmt.Exec(jobID, stream, line); err != nil { + tx.Rollback() + return err + } + } + return tx.Commit() +} + +func (r *JobLogRepository) GetByJobID(jobID int64, limit, offset int) ([]JobLog, error) { + if limit <= 0 { + limit = 1000 + } + if limit > 10000 { + limit = 10000 + } + rows, err := r.db.Query(` + SELECT id, job_id, stream, content, timestamp + FROM job_logs WHERE job_id = ? + ORDER BY id ASC LIMIT ? OFFSET ?`, + jobID, limit, offset) + if err != nil { + return nil, err + } + defer rows.Close() + var logs []JobLog + for rows.Next() { + var l JobLog + if err := rows.Scan(&l.ID, &l.JobID, &l.Stream, &l.Content, &l.Timestamp); err != nil { + return nil, err + } + logs = append(logs, l) + } + return logs, rows.Err() +} + +func (r *JobLogRepository) CountByJobID(jobID int64) (int64, error) { + var n int64 + err := r.db.QueryRow("SELECT COUNT(*) FROM job_logs WHERE job_id = ?", jobID).Scan(&n) + return n, err +} + +func (r *JobLogRepository) DeleteBefore(before time.Time) (int64, error) { + res, err := r.db.Exec("DELETE FROM job_logs WHERE timestamp < ?", before) + if err != nil { + return 0, err + } + return res.RowsAffected() +} + +type JobWithStats struct { + Job + DurationSeconds *int64 `db:"duration_seconds" json:"duration_seconds"` + LogLineCount int64 `db:"log_line_count" json:"log_line_count"` +} + +func (r *JobLogRepository) GetAllFiltered(limit, offset int, syncPairID *int64, status, triggerType string, from, to *time.Time) ([]JobWithStats, int64, error) { + where, args := []string{"1=1"}, []interface{}{} + if syncPairID != nil { + where = append(where, "j.sync_pair_id = ?") + args = append(args, *syncPairID) + } + if status != "" { + where = append(where, "j.status = ?") + args = append(args, status) + } + if triggerType != "" { + where = append(where, "j.trigger_type = ?") + args = append(args, triggerType) + } + if from != nil { + where = append(where, "j.created_at >= ?") + args = append(args, *from) + } + if to != nil { + where = append(where, "j.created_at <= ?") + args = append(args, *to) + } + whereClause := strings.Join(where, " AND ") + + var total int64 + countQuery := "SELECT COUNT(*) FROM jobs j WHERE " + whereClause + if err := r.db.QueryRow(countQuery, args...).Scan(&total); err != nil { + return nil, 0, err + } + + query := ` + SELECT + j.id, j.sync_pair_id, j.trigger_type, j.status, + j.started_at, j.finished_at, j.log_file, j.created_at, + CASE WHEN j.finished_at IS NOT NULL AND j.started_at IS NOT NULL + THEN (j.finished_at - j.started_at) ELSE NULL END as duration_seconds, + (SELECT COUNT(*) FROM job_logs WHERE job_id = j.id) as log_line_count + FROM jobs j + WHERE ` + whereClause + ` + ORDER BY j.created_at DESC LIMIT ? OFFSET ?` + args = append(args, limit, offset) + + rows, err := r.db.Query(query, args...) + if err != nil { + return nil, 0, err + } + defer rows.Close() + + var jobs []JobWithStats + for rows.Next() { + var j JobWithStats + var started, finished sql.NullTime + var logFile sql.NullString + var durationSeconds sql.NullInt64 + if err := rows.Scan(&j.ID, &j.SyncPairID, &j.TriggerType, &j.Status, + &started, &finished, &logFile, &j.CreatedAt, + &durationSeconds, &j.LogLineCount); err != nil { + return nil, 0, err + } + if started.Valid { + j.StartedAt = &started.Time + } + if finished.Valid { + j.FinishedAt = &finished.Time + } + if logFile.Valid { + j.LogFile = &logFile.String + } + if durationSeconds.Valid { + j.DurationSeconds = &durationSeconds.Int64 + } + jobs = append(jobs, j) + } + return jobs, total, rows.Err() +} diff --git a/internal/scheduler/scheduler.go b/internal/scheduler/scheduler.go index 6f25c66..4d9f136 100644 --- a/internal/scheduler/scheduler.go +++ b/internal/scheduler/scheduler.go @@ -32,6 +32,8 @@ func New(database interface{ SQLDB() *sql.DB }, engine *syncengine.Engine, cfg * func (s *Scheduler) Start() { s.wg.Add(1) go s.run() + s.wg.Add(1) + go s.cleanupRun() slog.Info("scheduler started") } @@ -93,3 +95,46 @@ func (s *Scheduler) tick() { }(jobID, sch.SyncPairID, sch.ID) } } + +func (s *Scheduler) cleanupRun() { + defer s.wg.Done() + ticker := time.NewTicker(24 * time.Hour) + defer ticker.Stop() + + s.cleanup() + for { + select { + case <-s.stopCh: + return + case <-ticker.C: + s.cleanup() + } + } +} + +func (s *Scheduler) cleanup() { + retentionDays := s.cfg.Scheduler.RetentionDays + if retentionDays <= 0 { + return + } + before := time.Now().AddDate(0, 0, -retentionDays) + + logRepo := models.NewJobLogRepository(s.db) + jobRepo := models.NewJobRepository(s.db) + + deletedLogs, err := logRepo.DeleteBefore(before) + if err != nil { + slog.Error("cleanup: failed to purge old job logs", "error", err) + return + } + + deletedJobs, err := jobRepo.DeleteFinishedBefore(before) + if err != nil { + slog.Error("cleanup: failed to purge old jobs", "error", err) + return + } + + if deletedLogs > 0 || deletedJobs > 0 { + slog.Info("cleanup: purged old records", "logs_deleted", deletedLogs, "jobs_deleted", deletedJobs, "before", before.Format("2006-01-02")) + } +} diff --git a/internal/sshmanager/fingerprint.go b/internal/sshmanager/fingerprint.go new file mode 100644 index 0000000..b1f9061 --- /dev/null +++ b/internal/sshmanager/fingerprint.go @@ -0,0 +1,89 @@ +package sshmanager + +import ( + "crypto/ed25519" + "crypto/rand" + "crypto/sha256" + "crypto/x509" + "encoding/base64" + "encoding/pem" + "fmt" + "os" + "path/filepath" + "strings" +) + +func Fingerprint(publicKey string) (string, error) { + pubKey := strings.TrimSpace(publicKey) + parts := strings.Fields(pubKey) + if len(parts) < 2 { + return "", fmt.Errorf("invalid public key format") + } + keyData, err := base64.StdEncoding.DecodeString(parts[1]) + if err != nil { + return "", fmt.Errorf("decoding public key: %w", err) + } + if len(keyData) == ed25519.PublicKeySize { + h := sha256.Sum256(keyData) + return "SHA256:" + base64.RawStdEncoding.EncodeToString(h[:]), nil + } + h := sha256.Sum256(keyData) + return "SHA256:" + base64.RawStdEncoding.EncodeToString(h[:]), nil +} + +func ParsePublicKey(data []byte) ([]byte, string, error) { + block, _ := pem.Decode(data) + if block == nil { + return nil, "", fmt.Errorf("no PEM block found") + } + var pubKey []byte + var err error + switch block.Type { + case "PUBLIC KEY": + pubKey = block.Bytes + case "OPENSSH KEY": + parts := strings.Fields(string(block.Bytes)) + if len(parts) < 2 { + return nil, "", fmt.Errorf("invalid openssh key format") + } + pubKey, err = base64.StdEncoding.DecodeString(parts[1]) + if err != nil { + return nil, "", err + } + default: + return nil, "", fmt.Errorf("unknown PEM type: %s", block.Type) + } + h := sha256.Sum256(pubKey) + return pubKey, "SHA256:" + base64.RawStdEncoding.EncodeToString(h[:]), nil +} + +func GenerateKeyPair(label string, sshDir string) (privPath, pubPath, pubKey, fingerprint string, err error) { + if err := os.MkdirAll(sshDir, 0700); err != nil { + return "", "", "", "", fmt.Errorf("creating ssh dir: %w", err) + } + privPath = filepath.Join(sshDir, label+".key") + pubPath = privPath + ".pub" + if _, err := os.Stat(privPath); err == nil { + return "", "", "", "", fmt.Errorf("key already exists") + } + pub, priv, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + return "", "", "", "", fmt.Errorf("generating ed25519 key: %w", err) + } + privFile, err := os.OpenFile(privPath, os.O_CREATE|os.O_WRONLY, 0600) + if err != nil { + return "", "", "", "", fmt.Errorf("creating private key file: %w", err) + } + defer privFile.Close() + privBytes, err := x509.MarshalPKCS8PrivateKey(priv) + if err != nil { + return "", "", "", "", fmt.Errorf("marshaling private key: %w", err) + } + pem.Encode(privFile, &pem.Block{Type: "PRIVATE KEY", Bytes: privBytes}) + pubKey = fmt.Sprintf("%s %s", strings.TrimSpace(string(pub)), label) + if err := os.WriteFile(pubPath, []byte(pubKey), 0644); err != nil { + return "", "", "", "", fmt.Errorf("writing public key: %w", err) + } + fp, _ := Fingerprint(pubKey) + return privPath, pubPath, pubKey, fp, nil +} diff --git a/internal/syncengine/engine.go b/internal/syncengine/engine.go index 7da2ca1..998b146 100644 --- a/internal/syncengine/engine.go +++ b/internal/syncengine/engine.go @@ -19,17 +19,18 @@ type Engine struct { db *sql.DB cfg *config.Config queue *Queue - eventBus chan Event + eventBus *EventBus mu sync.RWMutex stopped bool } type Event struct { - Type string - JobID int64 - Status string - Line string - Stream string + Type string + JobID int64 + Key string + Value string + Line string + Stream string } func New(database interface{ SQLDB() *sql.DB }, cfg *config.Config) *Engine { @@ -37,7 +38,7 @@ func New(database interface{ SQLDB() *sql.DB }, cfg *config.Config) *Engine { db: database.SQLDB(), cfg: cfg, queue: NewQueue(), - eventBus: make(chan Event, 100), + eventBus: NewEventBus(200), } return e } @@ -45,8 +46,12 @@ func New(database interface{ SQLDB() *sql.DB }, cfg *config.Config) *Engine { func (e *Engine) Start() {} func (e *Engine) Stop() {} -func (e *Engine) Events() <-chan Event { - return e.eventBus +func (e *Engine) SubscribeJob(jobID int64) (chan Event, func()) { + return e.eventBus.Subscribe(jobID) +} + +func (e *Engine) SubscribeGlobal() (chan Event, func()) { + return e.eventBus.SubscribeGlobal() } func (e *Engine) Run(ctx context.Context, jobID int64, pairID int64) error { @@ -106,7 +111,7 @@ func (e *Engine) Run(ctx context.Context, jobID int64, pairID int64) error { } e.setJobStatus(jobID, "waking_up") - e.emit(Event{Type: "status", JobID: jobID, Status: "waking_up"}) + e.emit(Event{Type: "status", JobID: jobID, Key: "status", Value: "waking_up"}) var targetMachine *models.Machine var remotePort int @@ -134,52 +139,103 @@ func (e *Engine) Run(ctx context.Context, jobID int64, pairID int64) error { timeout := time.Duration(targetMachine.WakeTimeoutSeconds) * time.Second interval := time.Duration(targetMachine.WakeCheckIntervalSeconds) * time.Second if err := wol.WaitUntilReady(jobCtx, targetMachine.Host, remotePort, timeout, interval, false); err != nil { - e.setJobStatus(jobID, "failed") - e.emit(Event{Type: "status", JobID: jobID, Status: "failed", Line: err.Error()}) - return fmt.Errorf("machine not ready: %w", err) + e.setJobStatus(jobID, "failed") + e.emit(Event{Type: "status", JobID: jobID, Key: "status", Value: "failed", Line: err.Error()}) + return fmt.Errorf("machine not ready: %w", err) } } } e.setJobStatus(jobID, "running") e.setJobLogFile(jobID, logPath) - e.emit(Event{Type: "status", JobID: jobID, Status: "running"}) + e.emit(Event{Type: "status", JobID: jobID, Key: "status", Value: "running"}) slog.Info("job started", "job_id", jobID, "pair", pair.Name) - var privKey string - if targetMachine != nil && targetMachine.SSHKeyID != nil { + privKey, err := e.resolveSSHKey(targetMachine) + if err != nil { + slog.Warn("failed to resolve SSH key, using server key", "error", err) privKey = filepath.Join(e.cfg.SSHDir(), "id_ed25519") } runner := NewRsyncRunner(e.cfg.SSHDir(), privKey) + logRepo := models.NewJobLogRepository(e.db) + var outBuf, errBuf []string + flush := func() { + if len(outBuf) > 0 { + logRepo.InsertBatch(jobID, "stdout", outBuf) + for _, l := range outBuf { + e.emit(Event{Type: "log", JobID: jobID, Stream: "stdout", Line: l}) + } + outBuf = nil + } + if len(errBuf) > 0 { + logRepo.InsertBatch(jobID, "stderr", errBuf) + for _, l := range errBuf { + e.emit(Event{Type: "log", JobID: jobID, Stream: "stderr", Line: l}) + } + errBuf = nil + } + } + onLine := func(stream, line string) { - e.emit(Event{Type: "log", JobID: jobID, Stream: stream, Line: line}) + f, _ := os.OpenFile(logPath, os.O_APPEND|os.O_WRONLY, 0644) + if f != nil { + fmt.Fprintln(f, line) + f.Close() + } + if stream == "stdout" { + outBuf = append(outBuf, line) + } else { + errBuf = append(errBuf, line) + } + if len(outBuf) >= 50 || len(errBuf) >= 50 { + flush() + } } result, err := runner.Run(jobCtx, cfg, onLine) + flush() + if err != nil { if jobCtx.Err() != nil { - e.setJobStatus(jobID, "cancelled") - e.emit(Event{Type: "status", JobID: jobID, Status: "cancelled"}) + e.setJobStatus(jobID, "cancelled") + e.emit(Event{Type: "status", JobID: jobID, Key: "status", Value: "cancelled"}) return jobCtx.Err() } e.setJobStatus(jobID, "failed") - e.emit(Event{Type: "status", JobID: jobID, Status: "failed", Line: err.Error()}) + e.emit(Event{Type: "status", JobID: jobID, Key: "status", Value: "failed", Line: err.Error()}) return fmt.Errorf("rsync error: %w", err) } if result.ExitCode != 0 { - e.setJobStatus(jobID, "failed") - e.emit(Event{Type: "status", JobID: jobID, Status: "failed", Line: result.Stderr}) + e.setJobStatus(jobID, "failed") + e.emit(Event{Type: "status", JobID: jobID, Key: "status", Value: "failed", Line: result.Stderr}) return fmt.Errorf("rsync exited with code %d: %s", result.ExitCode, result.Stderr) } e.setJobStatus(jobID, "success") - e.emit(Event{Type: "status", JobID: jobID, Status: "success"}) + e.emit(Event{Type: "status", JobID: jobID, Key: "status", Value: "success"}) slog.Info("job completed", "job_id", jobID, "pair", pair.Name) + e.persistAndClose(jobID) return nil } +func (e *Engine) persistAndClose(jobID int64) { + e.eventBus.CloseJobChannels(jobID) +} + +func (e *Engine) resolveSSHKey(machine *models.Machine) (string, error) { + if machine == nil || machine.SSHKeyID == nil { + return filepath.Join(e.cfg.SSHDir(), "id_ed25519"), nil + } + sshKeyRepo := models.NewSSHKeyRepository(e.db) + sshKey, err := sshKeyRepo.GetByID(*machine.SSHKeyID) + if err != nil { + return "", fmt.Errorf("fetching ssh key: %w", err) + } + return sshKey.PrivateKeyPath, nil +} + func (e *Engine) Cancel(jobID int64, syncPairID int64) bool { if e.queue.IsRunning(syncPairID) { e.queue.Cancel(syncPairID) @@ -199,11 +255,7 @@ func (e *Engine) setJobLogFile(jobID int64, path string) { } func (e *Engine) emit(evt Event) { - select { - case e.eventBus <- evt: - default: - slog.Warn("event bus full, dropping event", "type", evt.Type) - } + e.eventBus.Publish(evt) } func buildPath(path string, machine *models.Machine) string { diff --git a/internal/syncengine/eventbus.go b/internal/syncengine/eventbus.go new file mode 100644 index 0000000..f5acfa4 --- /dev/null +++ b/internal/syncengine/eventbus.go @@ -0,0 +1,92 @@ +package syncengine + +import ( + "log/slog" + "sync" +) + +type EventBus struct { + subscribers map[int64]map[chan Event]struct{} + mu sync.RWMutex + global chan Event + bufferSize int +} + +func NewEventBus(bufferSize int) *EventBus { + return &EventBus{ + subscribers: make(map[int64]map[chan Event]struct{}), + global: make(chan Event, bufferSize), + bufferSize: bufferSize, + } +} + +func (eb *EventBus) Subscribe(jobID int64) (chan Event, func()) { + eb.mu.Lock() + defer eb.mu.Unlock() + if eb.subscribers[jobID] == nil { + eb.subscribers[jobID] = make(map[chan Event]struct{}) + } + ch := make(chan Event, eb.bufferSize) + eb.subscribers[jobID][ch] = struct{}{} + unsubscribe := func() { + eb.mu.Lock() + defer eb.mu.Unlock() + if subs, ok := eb.subscribers[jobID]; ok { + delete(subs, ch) + if len(subs) == 0 { + delete(eb.subscribers, jobID) + } + } + close(ch) + } + return ch, unsubscribe +} + +func (eb *EventBus) SubscribeGlobal() (chan Event, func()) { + eb.mu.RLock() + ch := make(chan Event, eb.bufferSize) + eb.mu.RUnlock() + go func() { + for evt := range eb.global { + select { + case ch <- evt: + default: + slog.Warn("global event subscriber buffer full, dropping event", "type", evt.Type) + } + } + close(ch) + }() + return ch, func() { close(ch) } +} + +func (eb *EventBus) Publish(evt Event) { + eb.mu.RLock() + defer eb.mu.RUnlock() + + if subs, ok := eb.subscribers[evt.JobID]; ok { + for ch := range subs { + select { + case ch <- evt: + default: + slog.Warn("job event subscriber buffer full, dropping event", "job_id", evt.JobID) + } + } + } + + select { + case eb.global <- evt: + default: + slog.Warn("global event bus full, dropping event", "type", evt.Type) + } +} + +func (eb *EventBus) CloseJobChannels(jobID int64) { + eb.mu.Lock() + defer eb.mu.Unlock() + if subs, ok := eb.subscribers[jobID]; ok { + for ch := range subs { + close(ch) + } + delete(eb.subscribers, jobID) + } +} diff --git a/web/src/App.tsx b/web/src/App.tsx index 2049276..8a0fc58 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -5,7 +5,9 @@ import Dashboard from './pages/Dashboard'; import Machines from './pages/Machines'; import SyncPairs from './pages/SyncPairs'; import JobHistory from './pages/JobHistory'; +import JobDetail from './pages/JobDetail'; import Settings from './pages/Settings'; +import SSHKeys from './pages/SSHKeys'; function ProtectedRoute({ children }: { children: JSX.Element }) { const [authed, setAuthed] = useState(null); @@ -27,7 +29,9 @@ export default function App() { } /> } /> } /> + } /> } /> + } /> } /> diff --git a/web/src/api/client.ts b/web/src/api/client.ts index e130c0d..f9f68f5 100644 --- a/web/src/api/client.ts +++ b/web/src/api/client.ts @@ -21,6 +21,10 @@ export async function api(path: string, opts: ApiOptions = {}): Promise { return res.json(); } +export async function apiRaw(path: string): Promise { + return fetch(`${BASE}${path}`, { credentials: 'include' }); +} + export interface User { id: number; username: string; @@ -64,4 +68,24 @@ export interface Job { started_at: string | null; finished_at: string | null; log_file: string | null; + duration_seconds?: number | null; + log_line_count?: number | null; +} + +export interface LogLine { + id: number; + job_id: number; + stream: string; + content: string; + timestamp: string; +} + +export interface SSHKey { + id: number; + label: string; + public_key: string; + fingerprint: string; + in_use: boolean; + has_private_key: boolean; + created_at: string; } diff --git a/web/src/pages/JobDetail.tsx b/web/src/pages/JobDetail.tsx new file mode 100644 index 0000000..bd8fadb --- /dev/null +++ b/web/src/pages/JobDetail.tsx @@ -0,0 +1,167 @@ +import { useEffect, useState, useRef } from 'react'; +import { useParams, Link } from 'react-router-dom'; +import { api, Job, LogLine, SyncPair } from '../api/client'; + +interface SSEEvent { + type: string; + job_id: number; + status?: string; + line?: string; + stream?: string; +} + +export default function JobDetail() { + const { id } = useParams<{ id: string }>(); + const [job, setJob] = useState(null); + const [pair, setPair] = useState(null); + const [logs, setLogs] = useState([]); + const [lines, setLines] = useState<{ stream: string; text: string }[]>([]); + const [autoScroll, setAutoScroll] = useState(true); + const logEndRef = useRef(null); + const esRef = useRef(null); + const jobId = Number(id); + + useEffect(() => { + loadJob(); + if (jobId) { + loadLogs(0); + const es = new EventSource(`/api/jobs/${jobId}/log/stream?job_id=${jobId}`); + esRef.current = es; + es.onmessage = (e) => { + const evt: SSEEvent = JSON.parse(e.data); + if (evt.type === 'log') { + setLines(prev => [...prev, { stream: evt.stream!, text: evt.line! }]); + } + if (evt.type === 'status') { + setJob(prev => prev ? { ...prev, status: evt.status! } : prev); + } + }; + } + return () => esRef.current?.close(); + }, [id]); + + useEffect(() => { + if (autoScroll && logEndRef.current) { + logEndRef.current.scrollIntoView({ behavior: 'smooth' }); + } + }, [lines, autoScroll]); + + async function loadJob() { + try { + const j = await api(`/api/jobs/${id}`); + setJob(j); + const pairs = await api('/api/sync-pairs'); + const p = pairs.find((sp: SyncPair) => sp.id === j.sync_pair_id); + setPair(p || null); + } catch {} + } + + async function loadLogs(offset: number) { + try { + const ls = await api(`/api/jobs/${id}/log?offset=${offset}&limit=1000`); + if (offset === 0) { + setLogs(ls); + } else { + setLogs(prev => [...prev, ...ls]); + } + } catch {} + } + + async function cancel() { + if (!confirm('Cancel this job?')) return; + try { + await api(`/api/jobs/${id}/cancel`, { method: 'POST' }); + loadJob(); + } catch { alert('Cancel failed'); } + } + + function statusColor(s: string) { + const map: Record = { + queued: 'bg-gray-600', waking_up: 'bg-yellow-600', running: 'bg-blue-600', + success: 'bg-green-600', failed: 'bg-red-600', cancelled: 'bg-gray-600', + }; + return map[s] || 'bg-gray-600'; + } + + function duration(j: Job) { + if (!j.started_at) return '-'; + const start = new Date(j.started_at).getTime(); + const end = j.finished_at ? new Date(j.finished_at).getTime() : Date.now(); + const secs = Math.round((end - start) / 1000); + if (secs < 60) return `${secs}s`; + const mins = Math.floor(secs / 60); + const rem = secs % 60; + if (mins < 60) return `${mins}m ${rem}s`; + return `${Math.floor(mins / 60)}h ${mins % 60}m`; + } + + if (!job) return
Loading...
; + + return ( +
+
+ ← Job History +

Job #{job.id}

+ + {job.status} + +
+ +
+
+
Sync Pair
+
{pair?.name || `Pair ${job.sync_pair_id}`}
+
+
+
Trigger
+
{job.trigger_type}
+
+
+
Duration
+
{duration(job)}
+
+
+
Started
+
{job.started_at ? new Date(job.started_at).toLocaleString() : '-'}
+
+
+ + {['queued', 'waking_up', 'running'].includes(job.status) && ( +
+ + +
+ )} + +
+
+ Output + {lines.length + logs.length} lines +
+
+ {logs.map(l => ( +
+ {((): string => { + const d = new Date(l.timestamp); + return `${d.getHours().toString().padStart(2,'0')}:${d.getMinutes().toString().padStart(2,'0')}:${d.getSeconds().toString().padStart(2,'0')}`; + })()} + {l.content} +
+ ))} + {lines.map((l, i) => ( +
+ LIVE + {l.text} +
+ ))} +
+
+
+
+ ); +} diff --git a/web/src/pages/JobHistory.tsx b/web/src/pages/JobHistory.tsx index af126cb..563d711 100644 --- a/web/src/pages/JobHistory.tsx +++ b/web/src/pages/JobHistory.tsx @@ -1,44 +1,54 @@ -import { useEffect, useState, useRef } from 'react'; +import { useEffect, useState } from 'react'; +import { Link } from 'react-router-dom'; import { api, Job, SyncPair } from '../api/client'; -interface SSEEvent { - type: string; - job_id: number; - status?: string; - line?: string; - stream?: string; -} - export default function JobHistory() { const [jobs, setJobs] = useState([]); const [pairs, setPairs] = useState([]); - const esRef = useRef(null); + const [filterStatus, setFilterStatus] = useState(''); + const [filterPair, setFilterPair] = useState(''); + const [filterRange, setFilterRange] = useState('7d'); + const [total, setTotal] = useState(0); + const [page, setPage] = useState(0); + const limit = 50; + + useEffect(() => { + loadPairs(); + }, []); useEffect(() => { load(); - const es = new EventSource('/api/jobs/stream'); - esRef.current = es; - es.onmessage = (e) => { - const evt: SSEEvent = JSON.parse(e.data); - if (evt.type === 'status') { - setJobs(prev => prev.map(j => j.id === evt.job_id ? { ...j, status: evt.status! } : j)); - } - }; - return () => es.close(); - }, []); + }, [filterStatus, filterPair, filterRange, page]); async function load() { try { - const [j, p] = await Promise.all([ - api('/api/jobs?limit=100'), - api('/api/sync-pairs'), - ]); - setJobs(j); - setPairs(p); + let url = `/api/jobs?limit=${limit}&offset=${page * limit}`; + if (filterStatus) url += `&status=${filterStatus}`; + if (filterPair) url += `&sync_pair_id=${filterPair}`; + if (filterRange === '24h') { + const from = new Date(Date.now() - 24 * 3600 * 1000).toISOString(); + url += `&from=${encodeURIComponent(from)}`; + } else if (filterRange === '7d') { + const from = new Date(Date.now() - 7 * 24 * 3600 * 1000).toISOString(); + url += `&from=${encodeURIComponent(from)}`; + } else if (filterRange === '30d') { + const from = new Date(Date.now() - 30 * 24 * 3600 * 1000).toISOString(); + url += `&from=${encodeURIComponent(from)}`; + } + const res = await fetch(url, { credentials: 'include' }); + const totalCount = res.headers.get('X-Total-Count'); + if (totalCount) setTotal(Number(totalCount)); + const data = await res.json(); + setJobs(data); } catch {} } + async function loadPairs() { + try { setPairs(await api('/api/sync-pairs')); } catch {} + } + async function cancel(id: number) { + if (!confirm('Cancel this job?')) return; try { await api(`/api/jobs/${id}/cancel`, { method: 'POST' }); load(); @@ -58,44 +68,102 @@ export default function JobHistory() { return map[s] || 'bg-gray-600'; } + function duration(j: Job) { + if (!j.started_at) return '-'; + const start = new Date(j.started_at).getTime(); + const end = j.finished_at ? new Date(j.finished_at).getTime() : Date.now(); + const secs = Math.round((end - start) / 1000); + if (secs < 60) return `${secs}s`; + const mins = Math.floor(secs / 60); + const rem = secs % 60; + if (mins < 60) return `${mins}m ${rem}s`; + return `${Math.floor(mins / 60)}h ${mins % 60}m`; + } + + const totalPages = Math.ceil(total / limit); + return (
-

Job History

- - - - - - - - - - - - - - {jobs.map(j => ( - - - - - - - - +
+

Job History

+
+ + + +
+
+ +
+
IDSync PairTriggerStatusStartedFinishedActions
{j.id}{pairName(j.sync_pair_id)}{j.trigger_type} - - {j.status} - - {j.started_at ? new Date(j.started_at).toLocaleString() : '-'}{j.finished_at ? new Date(j.finished_at).toLocaleString() : '-'} - {['queued', 'waking_up', 'running'].includes(j.status) && ( - - )} -
+ + + + + + + + + + - ))} - {jobs.length === 0 && } - -
IDSync PairTriggerStatusDurationStartedFinishedActions
No jobs
+ + + {jobs.map(j => ( + window.location.href = `/jobs/${j.id}`}> + + #{j.id} + + {pairName(j.sync_pair_id)} + {j.trigger_type} + + + {j.status} + + + {duration(j)} + {j.started_at ? new Date(j.started_at).toLocaleString() : '-'} + {j.finished_at ? new Date(j.finished_at).toLocaleString() : '-'} + e.stopPropagation()}> + {['queued', 'waking_up', 'running'].includes(j.status) && ( + + )} + + + ))} + {jobs.length === 0 && No jobs} + + + + {totalPages > 1 && ( +
+ + {page + 1} / {totalPages} ({total} total) + +
+ )} +
); } diff --git a/web/src/pages/Machines.tsx b/web/src/pages/Machines.tsx index a98e3a1..6eafbb3 100644 --- a/web/src/pages/Machines.tsx +++ b/web/src/pages/Machines.tsx @@ -1,11 +1,14 @@ import { useEffect, useState } from 'react'; -import { api, Machine } from '../api/client'; +import { Link } from 'react-router-dom'; +import { api, Machine, SSHKey } from '../api/client'; export default function Machines() { const [machines, setMachines] = useState([]); + const [sshKeys, setSSHKeys] = useState([]); const [showForm, setShowForm] = useState(false); const [form, setForm] = useState({ id: undefined as number | undefined, name: '', host: '', port: 22, ssh_user: 'root', + ssh_key_id: null as number | null, mac_address: '', wol_enabled: false, wake_timeout_seconds: 120, wake_check_interval_seconds: 5, }); @@ -13,7 +16,14 @@ export default function Machines() { useEffect(() => { load(); }, []); async function load() { - try { setMachines(await api('/api/machines')); } catch {} + try { + const [ms, ks] = await Promise.all([ + api('/api/machines'), + api('/api/ssh-keys'), + ]); + setMachines(ms); + setSSHKeys(ks); + } catch {} } async function handleSubmit(e: React.FormEvent) { @@ -21,7 +31,8 @@ export default function Machines() { try { const payload: Record = { id: form.id || null, name: form.name, host: form.host, port: Number(form.port), - ssh_user: form.ssh_user, mac_address: form.mac_address || null, + ssh_user: form.ssh_user, ssh_key_id: form.ssh_key_id, + mac_address: form.mac_address || null, wol_enabled: Boolean(form.wol_enabled), wake_timeout_seconds: Number(form.wake_timeout_seconds), wake_check_interval_seconds: Number(form.wake_check_interval_seconds), @@ -35,7 +46,7 @@ export default function Machines() { body: payload, }); setShowForm(false); - setForm({ id: undefined, name: '', host: '', port: 22, ssh_user: 'root', mac_address: '', wol_enabled: false as boolean, wake_timeout_seconds: 120, wake_check_interval_seconds: 5 }); + setForm({ id: undefined, name: '', host: '', port: 22, ssh_user: 'root', ssh_key_id: null, mac_address: '', wol_enabled: false, wake_timeout_seconds: 120, wake_check_interval_seconds: 5 }); load(); } catch (e: unknown) { alert((e as Error).message); } } @@ -43,7 +54,8 @@ export default function Machines() { function edit(m: Machine) { setForm({ id: m.id, name: m.name, host: m.host, port: m.port, - ssh_user: m.ssh_user, mac_address: m.mac_address || '', + ssh_user: m.ssh_user, ssh_key_id: m.ssh_key_id, + mac_address: m.mac_address || '', wol_enabled: m.wol_enabled, wake_timeout_seconds: m.wake_timeout_seconds, wake_check_interval_seconds: m.wake_check_interval_seconds, @@ -56,10 +68,19 @@ export default function Machines() { try { await api(`/api/machines/${id}`, { method: 'DELETE' }); load(); } catch { alert('Delete failed'); } } + function keyLabel(id: number | null) { + if (!id) return 'Server Key'; + const k = sshKeys.find(k => k.id === id); + return k ? k.label : `Key #${id}`; + } + return (
-

Machines

+
+

Machines

+ Manage SSH Keys +
@@ -67,12 +88,17 @@ export default function Machines() { {showForm && (
-
+

Machine

setForm({...form, name: e.target.value})} className="w-full bg-gray-700 rounded px-3 py-2 text-white" required /> setForm({...form, host: e.target.value})} className="w-full bg-gray-700 rounded px-3 py-2 text-white" required /> setForm({...form, port: +e.target.value})} className="w-full bg-gray-700 rounded px-3 py-2 text-white" /> setForm({...form, ssh_user: e.target.value})} className="w-full bg-gray-700 rounded px-3 py-2 text-white" /> + setForm({...form, mac_address: e.target.value})} className="w-full bg-gray-700 rounded px-3 py-2 text-white" />
diff --git a/web/src/pages/SSHKeys.tsx b/web/src/pages/SSHKeys.tsx new file mode 100644 index 0000000..7c3931b --- /dev/null +++ b/web/src/pages/SSHKeys.tsx @@ -0,0 +1,163 @@ +import { useEffect, useState, useRef } from 'react'; +import { api, SSHKey } from '../api/client'; + +export default function SSHKeys() { + const [keys, setKeys] = useState([]); + const [showGen, setShowGen] = useState(false); + const [showImport, setShowImport] = useState(false); + const [genLabel, setGenLabel] = useState(''); + const [importLabel, setImportLabel] = useState(''); + const [importPubKey, setImportPubKey] = useState(''); + const [downloading, setDownloading] = useState(null); + const [copied, setCopied] = useState(null); + const [loading, setLoading] = useState(false); + + useEffect(() => { load(); }, []); + + async function load() { + try { setKeys(await api('/api/ssh-keys')); } catch {} + } + + async function generate() { + if (!genLabel.trim()) { alert('Label is required'); return; } + setLoading(true); + try { + await api('/api/ssh-keys', { + method: 'POST', + body: { label: genLabel.trim(), generate: true }, + }); + setShowGen(false); + setGenLabel(''); + load(); + } catch (e: unknown) { alert((e as Error).message); } + finally { setLoading(false); } + } + + async function importKey() { + if (!importLabel.trim()) { alert('Label is required'); return; } + if (!importPubKey.trim()) { alert('Public key is required'); return; } + setLoading(true); + try { + await api('/api/ssh-keys', { + method: 'POST', + body: { label: importLabel.trim(), generate: false, public_key: importPubKey.trim() }, + }); + setShowImport(false); + setImportLabel(''); + setImportPubKey(''); + load(); + } catch (e: unknown) { alert((e as Error).message); } + finally { setLoading(false); } + } + + async function remove(id: number) { + if (!confirm('Delete this SSH key? Machines using it will fall back to the server key.')) return; + try { + await api(`/api/ssh-keys/${id}`, { method: 'DELETE' }); + load(); + } catch (e: unknown) { alert((e as Error).message); } + } + + async function downloadPrivate(id: number) { + try { + const res = await fetch(`/api/ssh-keys/${id}/private`, { credentials: 'include' }); + if (!res.ok) { + const err = await res.json().catch(() => ({ error: 'Failed' })); + alert((err as { error: string }).error); + return; + } + const blob = await res.blob(); + const url = URL.createObjectURL(blob); + const a = document.createElement('a'); + a.href = url; + a.download = `ssh-key-${id}.key`; + document.body.appendChild(a); + a.click(); + document.body.removeChild(a); + URL.revokeObjectURL(url); + setDownloading(id); + setTimeout(() => setDownloading(null), 3000); + } catch (e: unknown) { alert((e as Error).message); } + } + + async function copyPubKey(key: SSHKey) { + await navigator.clipboard.writeText(key.public_key); + setCopied(key.id); + setTimeout(() => setCopied(null), 2000); + } + + return ( +
+
+

SSH Keys

+
+ + +
+
+ + {(showGen || showImport) && ( +
+
+

{showGen ? 'Generate SSH Key Pair' : 'Import Public Key'}

+ setGenLabel(e.target.value)} + className="w-full bg-gray-700 rounded px-3 py-2 text-white" /> + {showImport && ( +