iop/apps/node/internal/store/store.go
toki d764adb390 chore: agent-task 변경 사항 반영
- edge node store 및 console 업데이트
- node adapters(cli, ollama, vllm, mock) 리팩토링
- node router, run_manager, store 개선
- proto 파일(control.proto, runtime.proto) 및 생성 코드 업데이트
- config 업데이트
2026-05-10 20:41:33 +09:00

173 lines
4.3 KiB
Go

package store
import (
"context"
"database/sql"
"fmt"
"time"
_ "modernc.org/sqlite"
"go.uber.org/zap"
)
const schema = `
CREATE TABLE IF NOT EXISTS runs (
run_id TEXT PRIMARY KEY,
adapter TEXT NOT NULL,
target TEXT NOT NULL,
session_id TEXT NOT NULL DEFAULT 'default',
background INTEGER NOT NULL DEFAULT 0,
status TEXT NOT NULL DEFAULT 'running',
created_at DATETIME NOT NULL,
completed_at DATETIME,
error TEXT
);
`
// RunRecord holds the persisted state of an execution.
type RunRecord struct {
RunID string
Adapter string
Target string
SessionID string
Background bool
Status string
CreatedAt time.Time
CompletedAt *time.Time
Error string
}
// Store persists run history to SQLite.
type Store struct {
db *sql.DB
logger *zap.Logger
}
// New opens (or creates) the SQLite database and runs the schema migration.
func New(dsn string, logger *zap.Logger) (*Store, error) {
db, err := sql.Open("sqlite", dsn)
if err != nil {
return nil, fmt.Errorf("store: open %s: %w", dsn, err)
}
db.SetMaxOpenConns(1) // SQLite is single-writer
if _, err := db.Exec(schema); err != nil {
_ = db.Close()
return nil, fmt.Errorf("store: migrate: %w", err)
}
if err := migrateRunsTable(db); err != nil {
_ = db.Close()
return nil, fmt.Errorf("store: migrate runs: %w", err)
}
logger.Info("store ready", zap.String("dsn", dsn))
return &Store{db: db, logger: logger}, nil
}
func migrateRunsTable(db *sql.DB) error {
cols, err := tableColumns(db, "runs")
if err != nil {
return err
}
if !cols["session_id"] {
if _, err := db.Exec(`ALTER TABLE runs ADD COLUMN session_id TEXT NOT NULL DEFAULT 'default'`); err != nil {
return err
}
}
if !cols["background"] {
if _, err := db.Exec(`ALTER TABLE runs ADD COLUMN background INTEGER NOT NULL DEFAULT 0`); err != nil {
return err
}
}
if !cols["target"] {
if cols["model"] {
// Local DB rename: copy model into target.
if _, err := db.Exec(`ALTER TABLE runs RENAME COLUMN model TO target`); err != nil {
return err
}
} else {
if _, err := db.Exec(`ALTER TABLE runs ADD COLUMN target TEXT NOT NULL DEFAULT ''`); err != nil {
return err
}
}
}
return nil
}
func tableColumns(db *sql.DB, table string) (map[string]bool, error) {
rows, err := db.Query(fmt.Sprintf("PRAGMA table_info(%s)", table))
if err != nil {
return nil, err
}
defer rows.Close()
cols := make(map[string]bool)
for rows.Next() {
var (
cid int
name string
ctype string
notnull int
dfltValue sql.NullString
pk int
)
if err := rows.Scan(&cid, &name, &ctype, &notnull, &dfltValue, &pk); err != nil {
return nil, err
}
cols[name] = true
}
return cols, rows.Err()
}
// Close closes the underlying database connection.
func (s *Store) Close() error {
return s.db.Close()
}
// InsertRun records a new run.
func (s *Store) InsertRun(ctx context.Context, r RunRecord) error {
sessionID := r.SessionID
if sessionID == "" {
sessionID = "default"
}
bg := 0
if r.Background {
bg = 1
}
_, err := s.db.ExecContext(ctx,
`INSERT INTO runs (run_id, adapter, target, session_id, background, status, created_at) VALUES (?,?,?,?,?,?,?)`,
r.RunID, r.Adapter, r.Target, sessionID, bg, r.Status, r.CreatedAt.UTC(),
)
return err
}
// CompleteRun marks a run as completed, failed, or cancelled.
func (s *Store) CompleteRun(ctx context.Context, runID, status, errMsg string) error {
_, err := s.db.ExecContext(ctx,
`UPDATE runs SET status=?, completed_at=?, error=? WHERE run_id=?`,
status, time.Now().UTC(), errMsg, runID,
)
return err
}
// GetRun retrieves a run record by ID.
func (s *Store) GetRun(ctx context.Context, runID string) (*RunRecord, error) {
row := s.db.QueryRowContext(ctx,
`SELECT run_id, adapter, target, session_id, background, status, created_at, completed_at, error FROM runs WHERE run_id=?`,
runID,
)
var r RunRecord
var bg int
var completedAt sql.NullTime
var errMsg sql.NullString
if err := row.Scan(&r.RunID, &r.Adapter, &r.Target, &r.SessionID, &bg, &r.Status, &r.CreatedAt, &completedAt, &errMsg); err != nil {
if err == sql.ErrNoRows {
return nil, nil
}
return nil, err
}
r.Background = bg != 0
if completedAt.Valid {
r.CompletedAt = &completedAt.Time
}
r.Error = errMsg.String
return &r, nil
}