Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 88 additions & 1 deletion internal/cli/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (

"github.com/CrisSTEM/signalscope/internal/config"
fetchruntime "github.com/CrisSTEM/signalscope/internal/fetch"
scheduleruntime "github.com/CrisSTEM/signalscope/internal/schedule"
"github.com/CrisSTEM/signalscope/internal/storage"
)

Expand Down Expand Up @@ -81,6 +82,8 @@ func (app *App) Run(ctx context.Context, args []string) int {
return app.runSeedDemo(ctx, args[1:])
case "fetch":
return app.runFetch(ctx, args[1:])
case "schedule":
return app.runSchedule(ctx, args[1:])
default:
fmt.Fprintf(app.Stderr, "unknown command %q\n\n", args[0])
app.printUsage()
Expand Down Expand Up @@ -290,7 +293,88 @@ func (app *App) runFetch(ctx context.Context, args []string) int {
}

fmt.Fprint(app.Stdout, fetchruntime.FormatSummary(summary))
if summary.HasFailures() {
return 1
}

return 0
}

func (app *App) runSchedule(ctx context.Context, args []string) int {
flags := flag.NewFlagSet("schedule", flag.ContinueOnError)
flags.SetOutput(app.Stderr)

configPath := flags.String("config", "", "Path to the config pack directory")
dbPath := flags.String("db", "", "Path to the SQLite database file")
jobID := flags.String("job", "", "Optional schedule job ID to narrow execution to a single configured job")
flags.Usage = func() {
fmt.Fprintln(app.Stderr, "Usage: signalscope schedule --config <dir> --db <path> [--job <schedule-job-id>]")
fmt.Fprintln(app.Stderr)
flags.PrintDefaults()
}

if err := flags.Parse(args); err != nil {
if errors.Is(err, flag.ErrHelp) {
return 0
}
return exitUsage
}
if flags.NArg() != 0 {
fmt.Fprintf(app.Stderr, "schedule: unexpected arguments: %s\n\n", strings.Join(flags.Args(), " "))
flags.Usage()
return exitUsage
}
if strings.TrimSpace(*configPath) == "" {
fmt.Fprintln(app.Stderr, "schedule: --config is required")
fmt.Fprintln(app.Stderr)
flags.Usage()
return exitUsage
}
if strings.TrimSpace(*dbPath) == "" {
fmt.Fprintln(app.Stderr, "schedule: --db is required")
fmt.Fprintln(app.Stderr)
flags.Usage()
return exitUsage
}

pack, err := app.loadPack(*configPath)
if err != nil {
fmt.Fprintf(app.Stderr, "schedule: %v\n", err)
return 1
}

db, err := app.openDB(*dbPath)
if err != nil {
fmt.Fprintf(app.Stderr, "schedule: open database: %v\n", err)
return 1
}
defer db.Close()

if err := app.bootstrap(ctx, db); err != nil {
fmt.Fprintf(app.Stderr, "schedule: bootstrap database: %v\n", err)
return 1
}

if err := app.syncPack(ctx, db, pack); err != nil {
fmt.Fprintf(app.Stderr, "schedule: sync config pack: %v\n", err)
return 1
}

runtime := scheduleruntime.Runtime{
DB: db,
Registry: app.registry(),
Now: app.Now,
}

summary, err := runtime.Execute(ctx, pack, scheduleruntime.Options{
JobID: *jobID,
})
if err != nil {
fmt.Fprintf(app.Stderr, "schedule: %v\n", err)
return 1
}

fmt.Fprint(app.Stdout, scheduleruntime.FormatSummary(summary))
if summary.HasFailures() {
return 1
}
Expand All @@ -303,14 +387,16 @@ func (app *App) printUsage() {
fmt.Fprintln(app.Stderr, " signalscope check-config --config <dir>")
fmt.Fprintln(app.Stderr, " signalscope seed-demo --db <path> [--config <dir>]")
fmt.Fprintln(app.Stderr, " signalscope fetch --config <dir> --db <path> --source <source-kind> [--binding <binding-id>]")
fmt.Fprintln(app.Stderr, " signalscope schedule --config <dir> --db <path> [--job <schedule-job-id>]")
fmt.Fprintln(app.Stderr)
fmt.Fprintln(app.Stderr, "Commands:")
fmt.Fprintln(app.Stderr, " check-config Load and validate a config pack")
fmt.Fprintln(app.Stderr, " seed-demo Seed a SQLite database from the demo pack")
fmt.Fprintln(app.Stderr, " fetch Execute the ingestion runtime for one source kind")
fmt.Fprintln(app.Stderr, " schedule Execute one deterministic scheduler pass")
fmt.Fprintln(app.Stderr)
fmt.Fprintln(app.Stderr, "Note:")
fmt.Fprintln(app.Stderr, " Only source kinds registered in the runtime can be fetched.")
fmt.Fprintln(app.Stderr, " Only source kinds registered in the runtime can be fetched or scheduled.")
}

func (app *App) ensureDefaults() {
Expand Down Expand Up @@ -358,6 +444,7 @@ func (app *App) bootstrap(ctx context.Context, db *sql.DB) error {
func (app *App) seedPack(ctx context.Context, db *sql.DB, pack config.Pack) error {
return app.SeedPack(ctx, db, pack)
}

func (app *App) syncPack(ctx context.Context, db *sql.DB, pack config.Pack) error {
return app.SyncPack(ctx, db, pack)
}
Expand Down
Loading
Loading