Documentation
¶
Overview ¶
Package cloud is CSF's record of the paid jobs it runs on other people's machines, and its hand on them: every job it launches, what each has spent against its cap, a stop the operator presses that is checked against the provider's own job list, and a watcher that cancels a job at its cap and every job at the day's cap even when the provider's own timeout was set wrong. Every start, stop and cap hit is a notice to the operator.
The record is one file under the harness state directory, LedgerFile, replaced whole on every change, which the Workbench's Cloud panel follows. It is a service: the binary grants the provider, the state directory and the notifier, and mounts it; one goroutine owns the record.
Index ¶
- Constants
- Variables
- func FormatUSD(amount float64) string
- type BurstCredentials
- type CloudJob
- type CloudJobs
- func (jobs *CloudJobs) Jobs(ctx context.Context, input ListCloudJobsInput) (Ledger, error)
- func (jobs *CloudJobs) LaunchBurst(ctx context.Context, input LaunchBurstInput) (LaunchBurstOutput, error)
- func (jobs *CloudJobs) LaunchProbe(ctx context.Context, input LaunchProbeInput) (CloudJob, error)
- func (jobs *CloudJobs) LaunchSuiteShard(ctx context.Context, input LaunchSuiteShardInput) (LaunchBurstOutput, error)
- func (jobs *CloudJobs) Register(router gin.IRouter)
- func (jobs *CloudJobs) Start(scope *runtime.Scope) error
- func (jobs *CloudJobs) Stop(ctx context.Context, input StopCloudJobInput) (CloudJob, error)
- func (jobs *CloudJobs) Tools() []csf.Option
- type Collector
- type DryRunProvider
- func (provider *DryRunProvider) Cancel(ctx context.Context, id string) error
- func (provider *DryRunProvider) Hardware(ctx context.Context) ([]hfjobs.Hardware, error)
- func (provider *DryRunProvider) JobURL(id string) string
- func (provider *DryRunProvider) List(ctx context.Context) ([]hfjobs.Job, error)
- func (provider *DryRunProvider) Run(ctx context.Context, spec hfjobs.Spec) (hfjobs.Job, error)
- type IJobProvider
- type INotifier
- type LaunchBurstInput
- type LaunchBurstOutput
- type LaunchProbeInput
- type LaunchSuiteShardInput
- type Ledger
- type ListCloudJobsInput
- type Notice
- type NoticeKind
- type Option
- func WithBurstImage(image string) Option
- func WithBurstRepository(repository string) Option
- func WithClock(source clock.IClock) Option
- func WithDailyCap(capUSD float64) Option
- func WithLogger(logger *slog.Logger) Option
- func WithMergedCount(merged func(ctx context.Context, branches []string) (int, error)) Option
- func WithNotifier(notifier INotifier) Option
- func WithProvider(name string, provider IJobProvider) Option
- func WithSliceQueue(ready func(ctx context.Context, limit int) ([]dispatch.ReadySlice, error), ...) Option
- func WithStateDirectory(directory string) Option
- type StopCloudJobInput
- type StopReason
Constants ¶
const ( // LedgerFile is the record, under the harness state directory. LedgerFile = "cloud.json" // WatchInterval is the watcher's period. A job can overrun its cap by at // most one period of its own spend: at cpu-xl's measured price, 16,667 // micro-dollars a minute (/api/jobs/hardware, 2026-10-05), 30 s is $0.0083, // under 1% of a one-dollar cap, while one list call a period is nothing // against the API's limits. WatchInterval = 30 * time.Second // DefaultDailyCapUSD is the operator's hard cap on a day's fixer spend // (ruling of 2026-10-04: $100 a day), which cloud spend is part of. DefaultDailyCapUSD = 100.0 )
const ( // ProvidersFile holds the operator's provider credentials, under the // harness state directory, readable by its owner only. ProvidersFile = "providers.json" // The providers.json keys a burst reads, as a dry run names a missing one. KeyGitHubToken = "burst.github_token" KeyHuggingFaceToken = "burst.hf_token" KeyCorpusDataset = "burst.corpus_dataset" KeyExecutorEnvironment = "burst.executor_environment" )
A burst: N fixer sessions in one disposable cloud job, each a whole harness. The job runs burst.sh in the session image; this side picks the slices, builds the spec with the credentials as job secrets, launches it under its cap and holds the slices here so this host never runs them too.
const ( ListCloudJobsTool = "ListCloudJobs" StopCloudJobTool = "StopCloudJob" LaunchProbeTool = "LaunchCloudProbe" LaunchBurstTool = "LaunchBurstSessions" LaunchSuiteTool = "LaunchSuiteShard" JobsPath = "/api/cloud" StopPath = "/api/cloud/stop" ProbePath = "/api/cloud/probe" BurstPath = "/api/cloud/burst" SuitePath = "/api/cloud/suite" )
The cloud operations: each is one MCP tool on the CSF service and one HTTP route on the host's router; the Workbench's Cloud panel and the csf burst verb are clients of the same calls.
const (
// ProbeImage is busybox 1.37, pinned by its index digest (2026-10-05).
ProbeImage = "busybox:1.37@sha256:bdf57e528e45e4433820e045b29b4597825a1c9e38353532d90a01445013f82e"
)
A probe is the smallest job: a pinned busybox that sleeps until the provider's timeout stops it. It runs nothing of CSF and carries no secret, so it proves the panel, the Stop and both caps against the real provider before the first burst is paid for.
const ProviderDryRun = "dry-run"
ProviderDryRun names DryRunProvider in the record.
Variables ¶
var ( // ErrNoProvider reports a service built without a provider. ErrNoProvider = errors.New("cloud: a job provider is required") // ErrNoState reports a service built without the state directory. ErrNoState = errors.New("cloud: the harness state directory is required") // ErrNotStarted reports an operation before the service was mounted. ErrNotStarted = errors.New("cloud: the service is not running") // ErrUnknownJob reports a job the record does not hold. ErrUnknownJob = errors.New("cloud: no job in the record has that identifier") // ErrAlreadyStopped reports a stop of a job that is not running. ErrAlreadyStopped = errors.New("cloud: the job is not running") // ErrNotConfirmed reports a cancel the provider's job list does not show // yet; the watcher records the stop when it does. ErrNotConfirmed = errors.New("cloud: cancel requested, but the provider's job list does not show the job stopped yet") // ErrOverDailyCap reports a launch whose cap would take the day's spend // past the day's cap. ErrOverDailyCap = errors.New("cloud: the job's cap would pass the day's cap") // ErrInvalidCap reports a cap that buys less than a minute, or none. ErrInvalidCap = errors.New("cloud: the cap must buy at least one minute of the flavor") // ErrUnknownFlavor reports a flavor the provider does not price. ErrUnknownFlavor = errors.New("cloud: the provider has no such flavor") )
var ( // ErrNoImage reports a burst on a service granted no session image. ErrNoImage = errors.New("cloud: no session image is published for bursts; give serve -burst-image") // ErrNoRepository reports a burst on a service granted no repository. ErrNoRepository = errors.New("cloud: no repository is named for bursts; give serve -burst-repository") // ErrNoQueue reports a burst on a service granted no dispatcher. ErrNoQueue = errors.New("cloud: bursts need the slice dispatcher") // ErrNoReadySlice reports a burst with nothing to run. ErrNoReadySlice = errors.New("cloud: no slice is ready and uncontended") // ErrMissingCredentials reports a launch the operator's credentials do // not cover yet. ErrMissingCredentials = errors.New("cloud: providers.json lacks credentials a burst needs") // ErrInvalidSessions reports a session count out of range. ErrInvalidSessions = fmt.Errorf("cloud: a burst runs between 1 and %d sessions", maxSessions) )
Functions ¶
Types ¶
type BurstCredentials ¶
type BurstCredentials struct {
// GitHubToken is a fine-grained token for the repository, with contents
// and pull requests write.
GitHubToken string `json:"github_token"`
// HuggingFaceToken writes the run logs to the corpus dataset.
HuggingFaceToken string `json:"hf_token"`
// CorpusDataset is the private dataset the run logs go to, owner/name.
CorpusDataset string `json:"corpus_dataset"`
// ExecutorEnvironment is the executor's model credential, by the
// environment variable it reads, such as ANTHROPIC_API_KEY.
ExecutorEnvironment map[string]string `json:"executor_environment"`
}
BurstCredentials are the "burst" object of providers.json: what a burst job receives as job secrets. None is ever written into the image or the record.
type CloudJob ¶
type CloudJob struct {
ID string `json:"id"`
Provider string `json:"provider"`
Flavor string `json:"flavor"`
// Purpose is what the job is for: the burst and its slices.
Purpose string `json:"purpose"`
// URL is the job's own page at the provider.
URL string `json:"url,omitempty"`
CreatedAt time.Time `json:"created_at"`
FinishedAt *time.Time `json:"finished_at,omitempty"`
Stage hfjobs.Stage `json:"stage"`
CapUSD float64 `json:"cap_usd"`
// UnitCostMicroUSD is the flavor's price per minute when it launched.
UnitCostMicroUSD int64 `json:"unit_cost_micro_usd"`
// TimeoutSeconds is the provider's own stop: the minutes the cap buys.
TimeoutSeconds int64 `json:"timeout_seconds"`
// SpendUSD is the spend so far: the minutes since the job was created,
// until it finished, at the flavor's price. Billing from creation rather
// than from the start of running reads high, never low.
SpendUSD float64 `json:"spend_usd"`
StopReason StopReason `json:"stop_reason,omitempty"`
// Slices are the dispatcher's slices the job runs, held on this host.
Slices []string `json:"slices,omitempty"`
// Branches are the work branches its sessions push.
Branches []string `json:"branches,omitempty"`
// PullRequests is how many of those branches' pull requests merged,
// counted once the job ended; -1 until then.
PullRequests int `json:"pull_requests"`
}
CloudJob is one job in the record.
type CloudJobs ¶
type CloudJobs struct {
// contains filtered or unexported fields
}
CloudJobs is the record and the operations on it.
func NewCloudJobs ¶
NewCloudJobs validates the option set before building the service.
func (*CloudJobs) LaunchBurst ¶
func (jobs *CloudJobs) LaunchBurst(ctx context.Context, input LaunchBurstInput) (LaunchBurstOutput, error)
LaunchBurst picks the ready slices, builds the job and, unless it is a dry run, launches it under its cap and holds its slices on this host.
func (*CloudJobs) LaunchProbe ¶
LaunchProbe starts one probe under its cap.
func (*CloudJobs) LaunchSuiteShard ¶
func (jobs *CloudJobs) LaunchSuiteShard(ctx context.Context, input LaunchSuiteShardInput) (LaunchBurstOutput, error)
LaunchSuiteShard builds the job that replays a shard of the suite on a cloud node and, unless it is a dry run, launches it under its cap. The job runs burst.sh in suite mode: the build's released binary, replays at the suite's budget, pushes that never leave the job, no merge, and the replays uploaded beside the run logs for csf eval record.
func (*CloudJobs) Start ¶
Start starts the owner: it holds the record, runs every operation in turn and watches every running job once a WatchInterval.
type Collector ¶
type Collector struct {
// contains filtered or unexported fields
}
Collector measures the cloud record at each scrape.
func NewCloudCollector ¶
NewCloudCollector measures the record under the state directory.
func (*Collector) Collect ¶
func (collector *Collector) Collect(metrics chan<- prometheus.Metric)
Collect reads the record and sends every series.
func (*Collector) Describe ¶
func (collector *Collector) Describe(descriptions chan<- *prometheus.Desc)
Describe sends every series' description.
type DryRunProvider ¶
type DryRunProvider struct {
// contains filtered or unexported fields
}
DryRunProvider is a provider that runs nothing: a job it starts is a record that runs until it is cancelled, priced as the real flavor is. It is the fake job the Cloud panel is proved on before anything is paid for. Only the owner of CloudJobs calls it, one call at a time, so it holds its jobs without a lock.
func NewDryRunProvider ¶
func NewDryRunProvider(source clock.IClock) *DryRunProvider
NewDryRunProvider builds the provider over the clock its jobs are stamped by.
func (*DryRunProvider) Cancel ¶
func (provider *DryRunProvider) Cancel(ctx context.Context, id string) error
Cancel ends one running job.
func (*DryRunProvider) JobURL ¶
func (provider *DryRunProvider) JobURL(id string) string
JobURL is empty: a dry-run job has no page.
type IJobProvider ¶
type IJobProvider interface {
Run(ctx context.Context, spec hfjobs.Spec) (hfjobs.Job, error)
List(ctx context.Context) ([]hfjobs.Job, error)
Cancel(ctx context.Context, id string) error
Hardware(ctx context.Context) ([]hfjobs.Hardware, error)
JobURL(id string) string
}
IJobProvider is the provider capability: start, list and cancel jobs, and price a flavor. *hfjobs.HFJobsClient and *DryRunProvider satisfy it.
type LaunchBurstInput ¶
type LaunchBurstInput struct {
Sessions int `json:"sessions" jsonschema:"how many fixer sessions the job runs, one ready slice each"`
Flavor string `json:"flavor" jsonschema:"the provider's hardware flavor, such as cpu-xl"`
CapUSD float64 `` /* 128-byte string literal not displayed */
DryRun bool `json:"dry_run,omitempty" jsonschema:"build and show the job spec, secrets redacted, and launch nothing"`
}
LaunchBurstInput is one burst.
type LaunchBurstOutput ¶
type LaunchBurstOutput struct {
DryRun bool `json:"dry_run"`
Spec hfjobs.Spec `json:"spec"`
Slices []string `json:"slices"`
Missing []string `json:"missing"`
Job *CloudJob `json:"job,omitempty"`
}
LaunchBurstOutput is the burst: its spec (secrets redacted), the slices it runs, what providers.json still lacks, and the job once launched.
type LaunchProbeInput ¶
type LaunchProbeInput struct {
Flavor string `json:"flavor" jsonschema:"the provider's hardware flavor, such as cpu-basic"`
Purpose string `json:"purpose" jsonschema:"what the probe is for, as the Cloud panel shows it"`
CapUSD float64 `json:"cap_usd" jsonschema:"its spend cap in dollars; the provider's timeout is the minutes it buys"`
}
LaunchProbeInput is one probe job.
type LaunchSuiteShardInput ¶
type LaunchSuiteShardInput struct {
Build string `json:"build"`
Suite json.RawMessage `json:"suite"`
Jobs json.RawMessage `json:"jobs"`
Tickets []string `json:"tickets"`
Flavor string `json:"flavor"`
CapUSD float64 `json:"cap_usd"`
DryRun bool `json:"dry_run,omitempty"`
}
LaunchSuiteShardInput is one shard of an evaluation suite's replays run as a burst job (#416): the replay jobs and the suite as the evaluate package encodes them, the build they replay on, and the job's cap.
type Ledger ¶
type Ledger struct {
DailyCapUSD float64 `json:"daily_cap_usd"`
Jobs []CloudJob `json:"jobs"`
Notices []Notice `json:"notices"`
UpdatedAt time.Time `json:"updated_at"`
}
Ledger is the record: every job, the latest notices, and the day's cap.
func ReadLedger ¶
ReadLedger decodes a record; an empty or unreadable one is an empty record.
func (Ledger) RunningJobs ¶
RunningJobs is how many jobs may still be spending.
type Notice ¶
type Notice struct {
At time.Time `json:"at"`
Kind NoticeKind `json:"kind"`
JobID string `json:"job_id"`
Text string `json:"text"`
}
Notice is one message to the operator.
type NoticeKind ¶
type NoticeKind string
NoticeKind is what a notice reports.
const ( NoticeStart NoticeKind = "start" NoticeStop NoticeKind = "stop" NoticeCap NoticeKind = "cap" )
type Option ¶
Option configures CloudJobs.
func WithBurstImage ¶
WithBurstImage grants the published session image burst jobs run.
func WithBurstRepository ¶
WithBurstRepository names the repository burst jobs clone and open their pull requests on, owner/name.
func WithDailyCap ¶
WithDailyCap replaces DefaultDailyCapUSD.
func WithMergedCount ¶
WithMergedCount grants the count of a set of branches' merged pull requests, which a burst's output is measured by once it ends.
func WithNotifier ¶
WithNotifier grants the path notices reach the operator's phone by. Without it a notice is on the Workbench only.
func WithProvider ¶
func WithProvider(name string, provider IJobProvider) Option
WithProvider grants the provider, recorded on each job by name. Required.
func WithSliceQueue ¶
func WithSliceQueue(ready func(ctx context.Context, limit int) ([]dispatch.ReadySlice, error), hold func(ctx context.Context, slice string, reason string) error) Option
WithSliceQueue grants the dispatcher: ready lists the slices that may run away from this host, hold holds one here with a reason.
func WithStateDirectory ¶
WithStateDirectory grants the harness state directory the record lives in and providers.json is read from. Required.
type StopCloudJobInput ¶
type StopCloudJobInput struct {
ID string `json:"id" jsonschema:"the job's identifier, as the record and the Cloud panel show it"`
}
StopCloudJobInput names the job to stop.
type StopReason ¶
type StopReason string
StopReason is why a job stopped.
const ( // StopOperator is the operator's Stop. StopOperator StopReason = "operator" // StopCap is the job's own cap, reached. StopCap StopReason = "cap" // StopDailyCap is the day's cap, reached. StopDailyCap StopReason = "daily cap" // StopEnded is a job the provider ended: it finished, failed, or its // timeout fired. StopEnded StopReason = "ended" )