client

package
v0.1.4 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 31, 2026 License: Apache-2.0 Imports: 19 Imported by: 0

Documentation

Overview

Package client is the typed HTTP client shared by every CLI verb.

Index

Constants

This section is empty.

Variables

View Source
var ErrDeviceHeldByAnother = errors.New("device is held by someone else")

ErrDeviceHeldByAnother means the device is leased to somebody else. A copy runs under a lease exactly like any other execution, so this is the whole of "there is no unauthenticated upload path".

View Source
var ErrNoDevice = errors.New("no device available")

ErrNoDevice mirrors the controller's 409: nothing free right now. A plain Submit queues instead of returning this; it only surfaces when the caller passes NoWait, opting out of the queue in favor of an immediate answer.

View Source
var ErrNotCancellable = errors.New("job not cancellable")

ErrNotCancellable mirrors the controller's 409 not_cancellable from POST /v1/jobs/{id}/kill: the job raced onto a different state (started running, or already finished) between the caller's last look and this kill attempt. Exposed as a sentinel, rather than left for callers to string-match the message, because rc run needs to tell this apart from an ordinary kill failure: it means the job may now actually be running unattended and deserves Stage 1's honest "STILL RUNNING" warning.

Functions

This section is empty.

Types

type Client

type Client struct {
	// contains filtered or unexported fields
}

func New

func New(baseURL, token string) *Client

func (*Client) AttachTTY

func (c *Client) AttachTTY(ctx context.Context, jobID string, term Terminal) error

AttachTTY joins a job's terminal: it copies the job's output to the local terminal and everything typed locally back to the job, until the job's output stream ends.

It may be called the moment a job has an ID, before it is scheduled. The relay holds both connect orders, and anything typed during a queue wait is held rather than dropped — which for a controller whose whole job is leasing exclusive GPUs is the ordinary case, not an edge one.

func (*Client) Clear

func (c *Client) Clear(ctx context.Context, deviceID string) error

Clear returns a quarantined device to the pool. Admin-only server-side.

The controller answers 409 when the device still has a live lease, and that is deliberately surfaced as an error rather than swallowed: this is the one manual override for a stuck GPU, and reporting success when nothing was cleared is worse than failing.

func (*Client) CopyFrom

func (c *Client) CopyFrom(ctx context.Context, o CopyOptions) error

CopyFrom streams remote off the device into local.

The mirror image: `tar -cf -` on the box writes the archive to its stdout, which is the relay's output direction, and it is extracted here. The remote end is not trusted with where those bytes land — see extractTar.

func (*Client) CopyTo

func (c *Client) CopyTo(ctx context.Context, o CopyOptions) error

CopyTo streams local onto the device, extracting it at remote.

The far side runs `tar -xf -`, reading the archive from the job's stdin, which is the relay's input direction. Nothing here builds the archive first: the tar writer writes straight into the socket as the walk produces entries.

func (*Client) Describe

func (c *Client) Describe(ctx context.Context, deviceID string) (*server.DescribeResponse, error)

Describe fetches everything `rc describe` shows about one device: its state and holder, every label with its provenance and age, the usage sheet and when it was last written, and recent job history.

func (*Client) Explain

func (c *Client) Explain(ctx context.Context, selector string) (*server.ExplainResponse, error)

Explain answers "if I submitted this selector right now, what would happen" without submitting anything: which devices match, which of those are free, and how deep the queue already is behind the ones that aren't.

func (*Client) Hold

func (c *Client) Hold(ctx context.Context, opts HoldOptions) (*model.Job, error)

Hold submits a hold: a job with kind "hold" and no command of its own. The worker chooses the sleeper it actually runs (internal/worker's execute) — never this client, never the caller — so a hold can never be used to run arbitrary code under a different label; see server.handleSubmit, which rejects a hold submission that carries a command. Hold goes through the exact same Submit/Enqueue/ScheduleOnce path an ordinary job does: the same allocation transaction, the same queue, the same wall-clock watchdog for expiry.

func (*Client) Job

func (c *Client) Job(ctx context.Context, id string) (*model.Job, error)

func (*Client) JobView

func (c *Client) JobView(ctx context.Context, id string) (*server.JobView, error)

JobView fetches a job together with its queue position.

func (*Client) Jobs

func (c *Client) Jobs(ctx context.Context, opts JobsOptions) ([]model.Job, error)

func (*Client) Kill

func (c *Client) Kill(ctx context.Context, id, submitter string) error

Kill cancels a queued job or terminates a running one.

func (*Client) Logs

func (c *Client) Logs(ctx context.Context, id string, out io.Writer, follow bool) error

Logs copies either a snapshot of the job's currently stored output or follows it until completion. Attached TTY and pipe jobs do not have stored logs.

func (*Client) Release

func (c *Client) Release(ctx context.Context, jobID, submitter string) error

Release ends a hold — or, for that matter, any job — early. It is a thin alias over Kill so ownership is checked identically: only the job's own submitter, or an admin token, may release what they didn't hold.

func (*Client) Retire

func (c *Client) Retire(ctx context.Context, deviceID string) (string, error)

Retire removes a device from the fleet for good. Admin-only. The note the controller returns is passed back so the caller can repeat it: a worker that still declares this device will recreate it on its next registration.

func (*Client) State

func (c *Client) State(ctx context.Context) (*server.StateResponse, error)

func (*Client) StreamLogs

func (c *Client) StreamLogs(ctx context.Context, id string, out io.Writer) error

StreamLogs copies the job's output to out until the job finishes. It remains a direct request so run and attach callers preserve their existing behavior.

func (*Client) Submit

func (c *Client) Submit(ctx context.Context, opts SubmitOptions) (*model.Job, error)

func (*Client) WaitScheduled

func (c *Client) WaitScheduled(ctx context.Context, id string, onPosition func(int)) (*model.Job, error)

WaitScheduled polls until the job leaves the queue, calling onPosition each time its position changes so the caller can show progress.

"Leaves the queue" is exactly what it returns on, which is NOT the same as "was scheduled": a queued job can also be killed or cancelled, and that too ends the wait and is returned with no error. Callers must look at the state they get back before announcing a job as running — see cli.NewRunCmd, which reports a job that left the queue without ever starting as what it is rather than printing "job X on gpu0" and then streaming a log that will never have anything in it.

Like WaitTerminal, it tolerates a bounded run of consecutive poll failures — a dropped connection, a controller restarting for a couple of seconds — with a fixed retry interval, resetting the count on every success, rather than giving up on the very first one (see the constants above for the resulting wall-clock budgets). This distinction matters more here than in WaitTerminal: a caller (rc run) that sees WaitScheduled give up may go on to cancel the job, since nothing is running yet and there is no lease to protect — but only if the job is actually still queued. A transient network blip must never be indistinguishable from "genuinely done waiting", or a brief controller hiccup can end up cancelling — or, worse, if the job was scheduled during the blip, flagging for kill — a job that was never in trouble.

The giveup error, when the retry budget above is exhausted, is NOT distinguished from an ordinary error by its type or content — deciding "the caller's own bound elapsed" by unwrapping context.DeadlineExceeded from this error is unsound: each poll already runs under its own scheduledPollTimeout-bounded context, so a hung controller produces exactly that same error, with no caller deadline in sight at all. Callers must instead check their own context's Err() directly (see cli.NewRunCmd's waitCtx) to tell a real, caller-requested deadline apart from this function simply giving up.

func (*Client) WaitTerminal

func (c *Client) WaitTerminal(ctx context.Context, id string, maxWait time.Duration) (*model.Job, error)

WaitTerminal polls until the job reaches a terminal state. maxWait bounds the whole call so a job stuck in "assigned" (no worker ever attached) or a controller that stops responding cannot hang it forever; pass 0 for no bound. On giveup — either the overall bound or too many consecutive poll failures — the returned error names the job and its last known state rather than leaving the caller with a bare timeout.

type CopyOptions

type CopyOptions struct {
	// DeviceID is the box, always — a copy has exactly one remote side.
	DeviceID string
	// Local is the path on this machine, Remote the path on the box.
	Local  string
	Remote string
	// Submitter is who the copy runs as, and therefore who holds the lease
	// it takes.
	Submitter string
	// Progress, if non-nil, receives one-line status notes.
	Progress io.Writer
}

CopyOptions is one transfer.

type HoldOptions

type HoldOptions struct {
	DeviceID  string
	Selector  string
	Submitter string
	// Reason is why the device is being held (e.g. "manual profiling"),
	// shown by rc devices and the dashboard.
	Reason string
	// TTL is required: unlike an ordinary job's MaxRuntime, a hold has no
	// "device default" to fall back on, since the whole point is a human
	// deciding how long they need the device. Capped by the device's
	// max_runtime exactly as a job's MaxRuntime is — rejected, never
	// clamped.
	TTL time.Duration
}

HoldOptions configures rc hold: taking a device for a human to use directly (a shell), not for a job. Give exactly one of DeviceID or Selector, matching SubmitOptions.

type JobsOptions

type JobsOptions struct {
	Limit     int
	DeviceID  string
	Submitter string
	State     model.JobState
}

type SubmitOptions

type SubmitOptions struct {
	DeviceID string
	// Selector picks a device by its labels instead of by exact ID. Give
	// exactly one of DeviceID or Selector.
	Selector       string
	Command        []string
	Cwd            string
	Env            map[string]string
	Submitter      string
	IdempotencyKey string
	// Priority orders queued jobs on the same device: higher runs sooner.
	Priority int
	// MaxRuntime and IdleTimeout are watchdog ceilings enforced by the
	// worker (task 5); zero means "use the device's configured default".
	MaxRuntime  time.Duration
	IdleTimeout time.Duration
	// NoWait opts out of stage 2's queue: a busy device fails fast with
	// ErrNoDevice instead of the job sitting queued behind it.
	NoWait bool
	// Kind is model.LeaseKindJob or model.LeaseKindHold; empty means job.
	// Callers submitting an ordinary job never set this — use Hold instead
	// of setting it here directly, since the controller rejects a hold
	// submission (Kind == model.LeaseKindHold) that also carries a Command.
	Kind string
	// Reason is why a hold was taken; meaningless for an ordinary job.
	Reason string
	// Stdio is model.StdioLogs (the default), StdioTTY or StdioPipe: where
	// the job's standard streams are wired. The two attached modes put the
	// process on the controller's relay instead of the log store, which is
	// what AttachTTY and CopyTo/CopyFrom then join.
	Stdio string
}

type Terminal

type Terminal struct {
	// In is where keystrokes come from, normally os.Stdin in raw mode.
	In io.Reader
	// Out is where the job's output is written, normally os.Stdout.
	Out io.Writer
	// Size reports the terminal's current size. Required: a session that
	// never announces a size leaves the far side at the PTY's 24x80 default,
	// which is why remote shells wrap at 80 columns forever.
	Size func() (rows, cols uint16, err error)
	// Resized, if non-nil, fires whenever the terminal changed size. Every
	// tick sends a fresh resize frame.
	Resized <-chan struct{}
}

Terminal is the local terminal AttachTTY drives, as an interface rather than an *os.File so the part that CAN be tested is tested. Raw mode and SIGWINCH need a real tty and belong to the caller (see cli.attachTerminal); what happens to the bytes afterwards does not, and that is what this covers.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL