Compute and task dispatch¶
The client uses Compute for both a daemon and an embedded control plane. A Compute can receive one provider descriptor or several Spec alternatives. The alternatives are evaluated against the daemon's cached provider offers.
import skyward as sky
@sky.function
def train(batch):
return model(batch)
with sky.Compute(
sky.Spec(sky.AWS(region="us-east-1"), accelerator="A100"),
sky.Spec(sky.VastAI(), accelerator="A100", max_hourly_cost=2.0),
nodes=2,
allocation="spot_if_available",
selection="cheapest",
executor=sky.Executor(type="process", concurrency=2),
) as compute:
result = train(batch) >> compute
provider=... is shorthand for one Spec. Pass either provider or positional specs, not both. Spec describes hardware requirements; node count, allocation, selection, image, executor, options, ports, volumes, and lifecycle settings belong to Compute.
Lifecycle¶
Entering a with Compute(...) block registers provider accounts, creates or attaches the compute resource, waits for readiness, and starts the client-side lease. Leaving the block deletes the compute by default. Set delete_on_exit=False to keep it alive, then reconnect with Compute.attached(ref).
With no url, the client uses the daemon at http://127.0.0.1:17590 and starts one there if none is running; that daemon stays up after the block ends. With url or SKYWARD_URL, it uses the daemon it was given and never starts one. With database=, it runs the control plane in this process over that file. Every path uses the same control-plane API.
Dispatch¶
@sky.function creates an inert Pending call. It runs only when dispatched:
| Expression | Result |
|---|---|
call() >> compute |
Run on one node and return the value |
call() @ compute |
Run on every node and return a list |
call() > compute |
Start asynchronously and return a Future |
a() & b() >> compute |
Run a Group and return results in submission order |
sky.gather(a(), b(), stream=True) >> compute |
Return an iterator as results arrive |
@sky.stream and stream_call() >> compute |
Yield items from a remote generator |
Inside a with Compute(...) block, >> sky uses the active compute. The explicit compute target is required outside that context.
Compute.map(fn, items) submits one pending call per item and returns results in input order. Compute.current_nodes() reports the number of ready nodes. Compute.resize(nodes) asks for a different size — 4, (2, 8), or a Nodes — and returns as soon as the intent is recorded, so current_nodes() is what says how much of it is real. A compute running a collective plugin cannot be resized.
Compute.update(image=...) asks for a new image on a compute created with Image(mutable=True). It follows the same rules as resize: the intent is recorded and the call returns, and a task submitted meanwhile waits for a ready node. Only the mutable fields take effect; the others are refused (see Core concepts). Compute(name=...) naming a mutable compute that exists adopts it and reconciles its nodes and image; a different provider or accelerator is refused. Naming a compute whose image is fixed still fails with name_taken. Compute.attached is unchanged.
Watching the compute¶
callbacks= takes an iterable of Callable[[Event, ComputeView], None]; each callable sees every event on the compute's stream together with the whole compute folded into an immutable ComputeView — state, nodes, bootstrap phases, tasks, cost, and a bounded window of errors. The stream replays from the compute's creation, so a callback registered at construction still sees the provisioning it was not around for. Compute.events() is the primitive underneath: an iterator of decoded events, replayed and then followed, reconnecting on transport drops without repeating an event. Compute.attached accepts callbacks= too. See Watching a compute.
Specifications and runtime options¶
Spec accepts provider, accelerator, cpus, memory_gb, region, disk_gb, architecture, and max_hourly_cost.
Options accepts provisioning and worker timeouts, retry settings, health checks, autoscaling settings, and the cluster capability flag. ready_timeout and shutdown_timeout control how long the current client waits for its compute, and strict_version refuses a daemon on another version of Skyward instead of warning about it.
Executor supports thread, process, and loky. concurrency sets the number of task slots per node, and buffer sets how many additional tasks can be admitted ahead of those slots. reuse=False is valid only for the process executor.
Reference¶
skyward.Compute
¶
A pool of machines, for as long as the with block lasts.
url decides where the control plane is, and nothing else changes: given
one, the pool talks to that daemon; given none, it talks to the one at the
default address, starting it if nobody has. A daemon it started is left
running — the machines it bought outlive this block, and something has to be
reconciling them. database is the exception: it runs the control plane in
this process, over that file.
retry is the pool's answer to an attempt that did not answer, a
(reason, attempt) -> bool — see :func:function, whose retry overrides
it per call. The default tries once more after a loss and never after an
exception; None retries nothing.
A name the daemon already holds for a compute that is up is adopted when that
compute's image is mutable, and not refused: the pool attaches to it and sends it
the nodes and image given here. The machines must be the same kind — a name
is not a way to reach machines of another provider or accelerator, and that is an
error. A compute whose image was not built mutable keeps the contract it always
had, and the same name is refused by the daemon.
attach is :meth:attached's: the compute that already exists, in place of
a definition. It is a constructor argument only so that the class has one
constructor.
id
property
¶
loop
property
¶
client
property
¶
__init__(*specs, provider=None, accelerator=None, cpus=None, memory_gb=None, region=None, nodes=1, allocation='spot_if_available', selection='cheapest', image=DEFAULT_IMAGE, plugins=(), executor=DEFAULT_EXECUTOR, options=DEFAULT_OPTIONS, ports=(), volumes=(), ttl=600, retry=retry.default, name=None, url=None, database=None, delete_on_exit=True, console=True, callbacks=(), attach=None)
¶
attached(ref, url=None, database=None, console=True, callbacks=(), delete_on_exit=False)
classmethod
¶
The compute that is already there, by name or by id.
with sky.Compute(provider=sky.AWS(), nodes=8, name="training", delete_on_exit=False) as pool:
...
with sky.Compute.attached("training") as pool: # tomorrow, another process
more(data) >> pool
The machines outlive the process that asked for them, which is the whole reason the control plane is a daemon and not a library. This is how a second process says so — it takes no spec, because the compute it is joining already has one, and a spec here could only disagree with it.
It does not delete on exit by default. A pool somebody else is using is not a pool to take down on the way out.
__enter__()
¶
Buy the machines and hand back the pool, or leave nothing behind trying.
An entry that raises never runs the matching exit, and by then the compute
exists and its machines are billing. So a failure here tears down what it
was in the middle of building, on the same terms the exit would — the
delete_on_exit the caller asked for — and the error that caused it is
what the caller sees, not whatever the teardown ran into.
__exit__(*_)
¶
run(pending)
¶
start(pending)
¶
broadcast(pending)
¶
gather(group)
¶
gather_stream(group)
¶
Each answer as it lands, rather than all of them at the end.
Every task is submitted up front, so they overlap; the yielding is what
differs. ordered walks the futures as submitted and blocks on the next
one due — a slow first call holds back the rest; the unordered path hands
over whichever finishes first and never waits on a straggler out of turn.
stream(pending)
¶
The items, as the machine produces them.
The task is submitted here and dispatched by the request that reads it — the loop below pulls one frame at a time, and the pull reaches all the way to the generator on the node. A consumer that stops consuming stops it.
The failure comes back as the last frame rather than as a status, because by the time a generator raises, the caller already has the items it yielded before it, and there is no other way to say so.
map(fn, items)
¶
One task per item, spread over the nodes, answers in the order asked.
events()
¶
The compute's event log, replayed from the start and then followed.
The primitive under callbacks=: one decoded :class:Event at a
time, for as long as the caller keeps reading. The stream has no end
while the compute exists, so the consumer is what decides when to stop
— a break closes it. An event this build does not know is skipped,
not raised: a daemon may be newer than the client watching it.
current_nodes()
¶
How many machines are ready — counted off the compute, which carries its nodes.
resize(nodes)
¶
Ask for a different number of machines, in the spelling nodes= takes.
A size is the one part of a definition that changes without replacing
anything: the machines already up are kept, and the difference is bought or
drained. Like every other write to the compute this records an intent and
returns — :meth:current_nodes is what says how much of it is real.
A compute running a collective is refused, because a resize cannot reach it: the process group is formed on the first task and never formed again, so a rank that arrives afterwards blocks in a rendezvous nobody else will attend.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
nodes
|
NodeSpec
|
|
required |
update(image)
¶
Change the image a mutable compute runs from here on.
Only the packages and the code change: pip, pip_indexes and
includes. The base, the interpreter, the rest of the image and the
machines themselves cannot. Every ready node goes through bootstrapping
again and comes back ready with the new packages; a task submitted in the
meantime waits for one.
A compute whose image was not created with mutable=True refuses this
with image_fixed. Like every other write it records an intent and
returns — the nodes are bootstrapping until they are not.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
image
|
Image
|
The image to run from now on, as the whole image: a field left out is a field removed, not kept. |
required |
skyward.Compute.attached(ref, url=None, database=None, console=True, callbacks=(), delete_on_exit=False)
classmethod
¶
The compute that is already there, by name or by id.
with sky.Compute(provider=sky.AWS(), nodes=8, name="training", delete_on_exit=False) as pool:
...
with sky.Compute.attached("training") as pool: # tomorrow, another process
more(data) >> pool
The machines outlive the process that asked for them, which is the whole reason the control plane is a daemon and not a library. This is how a second process says so — it takes no spec, because the compute it is joining already has one, and a spec here could only disagree with it.
It does not delete on exit by default. A pool somebody else is using is not a pool to take down on the way out.
skyward.Spec
dataclass
¶
provider
instance-attribute
¶
accelerator = None
class-attribute
instance-attribute
¶
cpus = None
class-attribute
instance-attribute
¶
memory_gb = None
class-attribute
instance-attribute
¶
region = None
class-attribute
instance-attribute
¶
disk_gb = None
class-attribute
instance-attribute
¶
architecture = None
class-attribute
instance-attribute
¶
max_hourly_cost = None
class-attribute
instance-attribute
¶
__init__(provider, accelerator=None, cpus=None, memory_gb=None, region=None, disk_gb=None, architecture=None, max_hourly_cost=None)
¶
skyward.Options
dataclass
¶
Operational tuning for a compute — timeouts, retries, autoscaling.
Sensible defaults reproduce the runtime's built-in behavior, so most pools never
construct one. The daemon-side knobs are carried to the control plane on the
spec; the two session timeouts (ready_timeout, shutdown_timeout) and
strict_version stay in this process, because they govern how this client
waits for and reaches its own pool and never leave it.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
ssh_timeout
|
float
|
Seconds to keep dialing a machine before giving up on reaching it. |
240.0
|
provision_retry_delay
|
float
|
Seconds between connection attempts to a machine still coming up. |
2.0
|
max_provision_attempts
|
int
|
How many times a dropped connection is redialed before the node is lost. |
30
|
provision_timeout
|
float
|
Seconds a bought machine has to get closer to an address before it is
given up on and replaced. Counted from the last progress the provider
reported, so a machine still pulling its image is waited for. |
300.0
|
worker_timeout
|
float
|
Seconds to wait for the worker to come up once bootstrap has finished. |
180.0
|
autoscale_idle_timeout
|
float
|
Seconds a node must sit idle, counted from the moment it became ready, before an elastic pool may reclaim it. |
120.0
|
autoscale_cooldown
|
float
|
Seconds between autoscaling decisions. |
0.0
|
task_queue_timeout
|
float
|
Seconds an attempt may wait for a machine before it starts, for a function that
names no |
0.0
|
task_run_timeout
|
float
|
Seconds an attempt may run once started, for a function that names no
|
0.0
|
health_command
|
str | None
|
A shell command run on each node to ask whether the machine is still usable.
|
None
|
health_interval
|
float
|
Seconds between health probes. |
30.0
|
health_failures
|
int
|
How many consecutive failed probes make a node lost, and so replaced. |
3
|
ready_timeout
|
float
|
Seconds to wait for the pool to become ready before giving up. |
300.0
|
shutdown_timeout
|
float
|
Seconds to wait for the pool to finish deleting on exit. |
60.0
|
strict_version
|
bool
|
Refuse a daemon that runs another version of skyward, before a machine is
bought. Left |
False
|
Examples:
>>> with sky.Compute(provider=sky.AWS(), options=sky.Options(ssh_timeout=120)) as pool:
... train(data) >> pool
ssh_timeout = 240.0
class-attribute
instance-attribute
¶
provision_retry_delay = 2.0
class-attribute
instance-attribute
¶
max_provision_attempts = 30
class-attribute
instance-attribute
¶
provision_timeout = 300.0
class-attribute
instance-attribute
¶
worker_timeout = 180.0
class-attribute
instance-attribute
¶
autoscale_idle_timeout = 120.0
class-attribute
instance-attribute
¶
autoscale_cooldown = 0.0
class-attribute
instance-attribute
¶
task_queue_timeout = 0.0
class-attribute
instance-attribute
¶
task_run_timeout = 0.0
class-attribute
instance-attribute
¶
health_command = None
class-attribute
instance-attribute
¶
health_interval = 30.0
class-attribute
instance-attribute
¶
health_failures = 3
class-attribute
instance-attribute
¶
health_checker = None
class-attribute
instance-attribute
¶
cluster = None
class-attribute
instance-attribute
¶
ready_timeout = 300.0
class-attribute
instance-attribute
¶
shutdown_timeout = 60.0
class-attribute
instance-attribute
¶
strict_version = False
class-attribute
instance-attribute
¶
__init__(ssh_timeout=240.0, provision_retry_delay=2.0, max_provision_attempts=30, provision_timeout=300.0, worker_timeout=180.0, autoscale_idle_timeout=120.0, autoscale_cooldown=0.0, task_queue_timeout=0.0, task_run_timeout=0.0, health_command=None, health_interval=30.0, health_failures=3, health_checker=None, cluster=None, ready_timeout=300.0, shutdown_timeout=60.0, strict_version=False)
¶
skyward.Executor
dataclass
¶
How the tasks run on the machine: where, how many, and how far ahead.
thread runs tasks on a bounded thread pool — the default, and the only
one that shares the worker's own address space, so the distributed collections
reach the cluster with nothing in between. process and loky run each
task in a subprocess, which is what a task that holds the GIL or leaks state
wants; they reach the collections over a bridge back to the worker.
reuse is a process knob and nothing else: a process pool with
reuse=False spends one subprocess per task and throws it away, which is the
clean-slate every time. reuse=True keeps the subprocesses between tasks, and
loky is the reusable pool that also restarts a worker that died — so reuse
does not apply to it, nor to thread, whose threads are always reused.
concurrency is the pool's width — how many tasks run at once. buffer is
the slack above it: that many more tasks are admitted and their payloads made
ready, so a slot that frees finds the next one in hand rather than a round trip
away. It is also the depth the daemon reads as backpressure before it grows the
compute.
Attributes:
| Name | Type | Description |
|---|---|---|
type |
{'thread', 'process', 'loky'}
|
The backend the tasks run on. |
reuse |
bool
|
Whether subprocesses live between tasks. Only meaningful for |
concurrency |
int | None
|
How many tasks run at once. |
buffer |
int
|
How many more tasks to admit and keep ready above |
skyward.Nodes
¶
Bases: Struct
How many machines to open with, and how much of that is negotiable.
initial is the size the pool asks for once, when it starts. min is the
count it is willing to live at: what lets a job of eight begin on four, and the
only floor the pool is held to afterwards — a machine the opening request never
got is not asked for again. max is the ceiling autoscaling may reach, and
setting it is what makes the pool elastic at all. Both unset means the pool
opens at initial and stays there.
skyward.Image
¶
Bases: Struct
The environment a node builds before it runs anything.
The base, the interpreter, the packages and where they resolve from. What the
user shipped from their own machine is not here — includes is packed into
a blob client-side and only its hash travels, because a spec is written to the
compute row and served back by the API.
base = None
class-attribute
instance-attribute
¶
python = None
class-attribute
instance-attribute
¶
pip = ()
class-attribute
instance-attribute
¶
apt = ()
class-attribute
instance-attribute
¶
pip_indexes = ()
class-attribute
instance-attribute
¶
env = field(default_factory=dict)
class-attribute
instance-attribute
¶
shell_vars = field(default_factory=dict)
class-attribute
instance-attribute
¶
includes = ()
class-attribute
instance-attribute
¶
excludes = ()
class-attribute
instance-attribute
¶
includes_sha256 = None
class-attribute
instance-attribute
¶
The user-code tarball, once the client has built it and put it in the blob
store. includes/excludes are the client's inputs; this is what the node
reads.
metrics = None
class-attribute
instance-attribute
¶
What the node measures about itself: :data:READINGS when None, exactly the list otherwise.
A :data:Reading is served by the node's own collector, a :class:MetricSpec by a
loop of its own. A name may appear once, whichever kind it is.
bootstrap_timeout = 900
class-attribute
instance-attribute
¶
skyward = 'auto'
class-attribute
instance-attribute
¶
warm = False
class-attribute
instance-attribute
¶
Whether a machine that finished bootstrapping is kept as a boot image.
Off because what it creates is never removed: an AMI holds a snapshot that bills
for its storage until it is deregistered, and nothing here deregisters it. Turning
it on is taking that on. What is created carries :meth:content_hash as a tag, on
the image and on the snapshot behind it, so it can be found again and removed.
Only providers that can snapshot a running machine honor it.
mutable = False
class-attribute
instance-attribute
¶
Whether the image of a compute that is up may be changed, best-effort.
Only :data:MUTABLE may change, and a node that is ready with another image is
sent through bootstrapping again and comes back ready. A package removed from
pip is not promised to leave the machine. A compute built without it keeps the
contract it always had: a different image is a different compute.
MUTABLE = frozenset({'pip', 'pip_indexes', 'includes', 'excludes', 'includes_sha256'})
class-attribute
¶
The fields a PATCH may change on a mutable image.
__post_init__()
¶
content_hash(source)
¶
Name the environment a bootstrapped machine ends up in.
Covers what the bootstrap installs — the base, the interpreter, the packages
and the indexes they are resolved from — together with source, which is
what stands in for a skyward version now that a node installs whatever the
daemon is running.
Left out is everything the bootstrap re-applies on every boot: the exports, the shell vars, the metric commands, and the user code, which is synced per run. Folding those in would split the images over changes that cost nothing to redo.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source
|
str
|
:attr: |
required |
Returns:
| Type | Description |
|---|---|
str
|
Twelve hex characters — long enough to name an image, short enough to read in one. |
digest()
¶
Name exactly this image, every field included.
A node reports the digest of the image it materialized, and the reconciler compares
it with the spec's to tell when the spec asks for another one. Unlike
:meth:content_hash it is not about warm images, and it includes what the
bootstrap re-applies: the exports, the shell vars, the metrics and the includes.
Returns:
| Type | Description |
|---|---|
str
|
The sha256 of the image's deterministic JSON encoding, in full. |
skyward.Port
dataclass
¶
Expose a node port on a fixed local port.
Each connection to 127.0.0.1:<local> is bridged to remote on a ready
node, chosen by route, over that node's existing SSH connection. The local
listener binds loopback only.
Attributes:
| Name | Type | Description |
|---|---|---|
remote |
int
|
The port the service listens on inside the node. |
local |
int
|
The local port to bind. |
route |
Route
|
How connections are spread across the ready nodes. |
Examples:
skyward.Volume
dataclass
¶
A bucket every node mounts, at a path you choose.
The nodes see a directory: ordinary open and os.listdir work, and a
dataset too big to ship as user code is read straight off object storage
instead. It is a network filesystem wearing a directory's clothes, so it is
read-heavy by design — read_only is the default for that reason.
Where the credentials come from is what storage decides. Left None,
the daemon resolves them from the provider the compute was bought on, and for
a bucket in that same account nothing secret is created or transmitted at all.
Set, they are yours — an R2 or Backblaze bucket the daemon has no account for —
and they travel to the daemon as a blob rather than on the spec, so they are
never served back by the compute API.
A compute takes one kind or the other, never a mix: two sets of credentials under one set of mounts is a compute nobody can say who is paying for.
Attributes:
| Name | Type | Description |
|---|---|---|
bucket |
str
|
The bucket to mount. On RunPod, which attaches storage instead of mounting it, the id or name of a network volume. |
mount |
str
|
The absolute path the bucket appears at on every node. |
prefix |
str
|
A subdirectory of the bucket to mount, rather than its root. |
read_only |
bool
|
Whether writes are refused. Buckets shared by several volumes are mounted writable if any of them asked to write. |
storage |
Storage | None
|
Credentials for buckets the provider cannot reach on its own. |
Examples:
skyward.function
¶
What the user writes, and what it does not do.
@function builds nothing but a description of a call. No pickling, no HTTP,
no compute: a Pending is inert until an operator hands it to a pool, which
is what lets the same call be dispatched to one node, to all of them, or not at
all.
Target = Pool | _Sky | ModuleType
¶
Where a call goes: a named pool, the sky stand-in, or the skyward
module itself (what import skyward as sky binds sky to).
Pool
¶
Pending
dataclass
¶
One call, described and not yet made.
timeout is how long each attempt may run once started, and queue_timeout
how long it may wait for a machine before it starts: None takes the pool's
task_run_timeout and task_queue_timeout, and 0 is no limit whatever
the pool says. retry is the call's retry decision: unset takes the pool's,
None turns retrying off for this call, and a function decides — see
:func:function.
fn
instance-attribute
¶
args
instance-attribute
¶
kwargs
instance-attribute
¶
timeout = None
class-attribute
instance-attribute
¶
retry = UNSET
class-attribute
instance-attribute
¶
queue_timeout = None
class-attribute
instance-attribute
¶
with_timeout(timeout)
¶
with_queue_timeout(queue_timeout)
¶
with_retry(retry)
¶
__rshift__(target)
¶
__matmul__(target)
¶
__gt__(target)
¶
__and__(other)
¶
__init__(fn, args, kwargs, timeout=None, retry=UNSET, queue_timeout=None)
¶
Group
dataclass
¶
Calls that go together.
Typed by what they return in common: a & b where both give an int is
a Group[int]. Mixing return types is allowed and lands on object —
the group is honest about what it can promise rather than pretending to know
which slot holds which type.
stream changes what >> gives back: a list once every call is in, or an
iterator that hands over each result the moment it is ready. ordered picks
between the two ways to be early — submission order, blocking only on the next
one due, or completion order, whichever finishes first.
Streaming
dataclass
¶
A call whose answer arrives in pieces.
timeout is how long the stream may run once started — None takes the
pool's task_run_timeout, and 0 is no limit.
What a generator function becomes. It is a separate type from Pending and
not a flag on it, because it is a separate promise: >> gives back an
iterator here, and the difference is worth knowing before the code runs rather
than after — a generator dispatched as an ordinary call would pickle the
generator object and fail on the machine.
gather(*pendings, stream=False, ordered=True)
¶
The same thing & builds, for when there are more than a few.
stream turns >> from a list into an iterator that yields each result as
it lands; ordered keeps that iterator in submission order, waiting on the
next one due, rather than in completion order. Both are inert until dispatched.
function(fn=None, *, timeout=None, queue_timeout=None, retry=UNSET)
¶
Turn a function into one that describes a call instead of making it.
Bare (@function) or with defaults (@function(timeout=600)), which any
single call can override with .with_timeout, .with_queue_timeout and
.with_retry.
timeout is how long each attempt may run, counted from when it starts on a
machine — not from the submission, so a call behind a long queue keeps all of
it, and a retry starts the count again. One that runs past it is stopped on the
machine and ends timed_out. queue_timeout is how long each attempt may
wait for a machine before it starts. None takes the pool's
task_run_timeout and task_queue_timeout; 0 is no limit.
retry is a (reason, attempt) -> bool asked when an attempt does not
answer. reason is the exception the function raised, or a :class:sky.Lost
when the attempt was lost — the process died under it, the worker restarted, the
machine went away — and attempt is the one that just failed, from one. An
exception is asked about on the node, where it was raised; a loss on the daemon.
Left unset, the call takes the pool's decision, whose default tries once more
after a loss and never after an exception: a function that raised is retried
only if the decision says so. None turns retrying off for this function.
stream(fn=None, *, timeout=None)
¶
Turn a generator into one that describes a stream instead of making it.
@sky.stream
def tokens(prompt: str) -> Iterator[str]:
yield from model.generate(prompt)
for token in tokens("hi") >> pool:
print(token)
A separate decorator rather than a flag on @function, because it is a
separate promise: >> hands back an iterator, and the items arrive as the
machine produces them. Worth knowing where the function is defined rather than
where it is called.
skyward.stream(fn=None, *, timeout=None)
¶
Turn a generator into one that describes a stream instead of making it.
@sky.stream
def tokens(prompt: str) -> Iterator[str]:
yield from model.generate(prompt)
for token in tokens("hi") >> pool:
print(token)
A separate decorator rather than a flag on @function, because it is a
separate promise: >> hands back an iterator, and the items arrive as the
machine produces them. Worth knowing where the function is defined rather than
where it is called.
skyward.Pending
dataclass
¶
One call, described and not yet made.
timeout is how long each attempt may run once started, and queue_timeout
how long it may wait for a machine before it starts: None takes the pool's
task_run_timeout and task_queue_timeout, and 0 is no limit whatever
the pool says. retry is the call's retry decision: unset takes the pool's,
None turns retrying off for this call, and a function decides — see
:func:function.
fn
instance-attribute
¶
args
instance-attribute
¶
kwargs
instance-attribute
¶
timeout = None
class-attribute
instance-attribute
¶
retry = UNSET
class-attribute
instance-attribute
¶
queue_timeout = None
class-attribute
instance-attribute
¶
with_timeout(timeout)
¶
with_queue_timeout(queue_timeout)
¶
with_retry(retry)
¶
__rshift__(target)
¶
__matmul__(target)
¶
__gt__(target)
¶
__and__(other)
¶
__init__(fn, args, kwargs, timeout=None, retry=UNSET, queue_timeout=None)
¶
skyward.Group
dataclass
¶
Calls that go together.
Typed by what they return in common: a & b where both give an int is
a Group[int]. Mixing return types is allowed and lands on object —
the group is honest about what it can promise rather than pretending to know
which slot holds which type.
stream changes what >> gives back: a list once every call is in, or an
iterator that hands over each result the moment it is ready. ordered picks
between the two ways to be early — submission order, blocking only on the next
one due, or completion order, whichever finishes first.
skyward.Streaming
dataclass
¶
A call whose answer arrives in pieces.
timeout is how long the stream may run once started — None takes the
pool's task_run_timeout, and 0 is no limit.
What a generator function becomes. It is a separate type from Pending and
not a flag on it, because it is a separate promise: >> gives back an
iterator here, and the difference is worth knowing before the code runs rather
than after — a generator dispatched as an ordinary call would pickle the
generator object and fail on the machine.
skyward.gather(*pendings, stream=False, ordered=True)
¶
The same thing & builds, for when there are more than a few.
stream turns >> from a list into an iterator that yields each result as
it lands; ordered keeps that iterator in submission order, waiting on the
next one due, rather than in completion order. Both are inert until dispatched.