package main import ( "context" "errors" "flag" "fmt" "io" "strings" "time" "github.com/atomine-elektrine/tarakan-client/internal/agent" "github.com/atomine-elektrine/tarakan-client/internal/api" "github.com/atomine-elektrine/tarakan-client/internal/app" repoctx "github.com/atomine-elektrine/tarakan-client/internal/context" "github.com/atomine-elektrine/tarakan-client/internal/updatecheck" ) func runWorker(ctx context.Context, arguments []string, stdout, stderr io.Writer, cfg api.Config) int { flags := flag.NewFlagSet("worker", flag.ContinueOnError) flags.SetOutput(stderr) var agentName, model, statePath string var once bool var interval, runFor time.Duration var maxJobs int var jobsOnly bool var skipCritic bool var urlFlag, hostFlag, tokenFlag string var minStars int var language, kind string flags.StringVar(&agentName, "agent", "", "local review backend (required)") flags.StringVar(&model, "model", "", "override the model for HTTP backends") flags.BoolVar(&once, "once", false, "process the current queue once and exit") flags.DurationVar(&interval, "interval", 30*time.Second, "delay between queue polls") flags.IntVar(&maxJobs, "max-jobs", 100, "maximum Jobs and repositories per queue pass") flags.BoolVar(&jobsOnly, "jobs-only", false, "process explicit Jobs only; skip the unscanned repository queue") flags.BoolVar(&skipCritic, "skip-critic", false, "skip the second evidence-validation agent pass") flags.StringVar(&statePath, "state-file", "", "durable worker state path") flags.IntVar(&minStars, "min-stars", 0, "only repos/jobs with at least this many stars") flags.StringVar(&language, "language", "", "only repos with this primary language (e.g. Rust, Elixir)") flags.StringVar(&language, "lang", "", "alias for --language") flags.StringVar(&kind, "kind", "", "only jobs of this kind (e.g. code_review, verify_findings)") // Subscription quota is use-it-or-lose-it on a rolling window. A bounded // run turns "spend my budget on strangers' repos" into "salvage what I was // going to lose". The client cannot see a provider's reset time, so the // window is the operator's to state rather than something guessed here. flags.DurationVar(&runFor, "for", 0, "stop after this long (e.g. 45m); salvages idle quota without running indefinitely") addAPIFlags(flags, &urlFlag, &hostFlag, &tokenFlag) flags.Usage = func() { fmt.Fprintln(stderr, "Usage: tarakan worker --agent codex [--once] [--min-stars N] [--language Rust]") fmt.Fprintln(stderr, "Continuously completes agent Jobs against pinned snapshots: Reports, Checks, and patch proposals.") flags.PrintDefaults() } if err := flags.Parse(arguments); err != nil { return 2 } if flags.NArg() != 0 || strings.TrimSpace(agentName) == "" { flags.Usage() return 2 } var err error cfg, err = mergeFlagConfig(cfg, urlFlag, hostFlag, tokenFlag) if err != nil { fmt.Fprintln(stderr, err) return 2 } registry := agent.Detect() provider, ok := registry.Find(agentName) if !ok { fmt.Fprintf(stderr, "agent %q is not installed or configured\n", agentName) return 1 } provider = provider.WithModel(model) local, _ := repoctx.Current() updatecheck.MaybeNotify(stderr, version) if runFor > 0 { var stopAfter context.CancelFunc ctx, stopAfter = context.WithTimeout(ctx, runFor) defer stopAfter() fmt.Fprintf(stdout, "%s Running for %s, then stopping.\n", time.Now().Format(time.RFC3339), runFor) } err = app.RunWorker(ctx, app.WorkerOptions{ APIConfig: cfg, Provider: provider, Local: local, Once: once, Interval: interval, MaxJobs: maxJobs, ReviewUnscanned: !jobsOnly, SkipCritic: skipCritic, StatePath: statePath, Filter: api.QueueFilter{ MinStars: minStars, Language: language, Kind: kind, }, Progress: func(message string) { fmt.Fprintln(stdout, time.Now().Format(time.RFC3339), message) }, }) // Reaching the --for window, or being interrupted, is the expected way to // stop; neither is a failure. if errors.Is(err, context.DeadlineExceeded) { fmt.Fprintf(stdout, "%s Run window reached; stopping.\n", time.Now().Format(time.RFC3339)) return 0 } if err != nil && !errors.Is(err, context.Canceled) { fmt.Fprintf(stderr, "worker stopped: %v\n", err) return 1 } return 0 }