Watching a compute¶
Everything observable about a compute — machines moving through their lifecycle, bootstrap phases turning over, lines the nodes print, gauges, cost, task outcomes — is an event in one log. The pool's live console is a reader of that log, and nothing more. callbacks= gives your code the same seat: every event, as it happens, together with the whole compute folded into one value.
Callbacks¶
A callback is a plain callable, Callable[[Event, ComputeView], None]. You hand any number of them to the pool, and each one sees every event:
import skyward as sky
def observer(event: sky.Event, compute: sky.ComputeView) -> None:
match event:
case sky.NodeEvent(node=node, state="ready"):
print(f"{node} up — {compute.nodes_ready}/{compute.nodes_total}")
case sky.PhaseEvent(phase=phase, event="completed"):
print(f"bootstrap: {phase} done")
case sky.ConsoleEvent(node=node, content=line):
forward_to_logging(node, line)
case sky.CostEvent():
budget.observe(compute.cost)
case _:
pass
with sky.Compute(provider=sky.AWS(), nodes=4, callbacks=(observer,)) as pool:
train(data) >> pool
Event is a tagged union — one struct per compute fact (ComputeReady, ComputeProvisioning, ComputeDegraded, ComputeDeleting, ComputeDeleted, ComputeBound, ComputeAbandoned, …), plus NodeEvent, PhaseEvent, ConsoleEvent, ProgressEvent, MetricEvent, CostEvent, TaskEvent — so match is the whole dispatch mechanism. There is no registration API to learn: one callable, one match, and the cases you did not write fall through.
Three properties make callbacks safe to lean on:
- The stream replays. The event log starts at the compute's creation and the subscription reads it from the beginning, so a callback registered at construction still sees the provisioning and bootstrap it could not possibly have been early enough for. And every change of the compute's state is in it: the daemon writes the state and the event that moved it in one transaction, so there is no transition a callback can miss. This is why
callbacks=is a constructor argument: the pool provisions inside__enter__, and by the time your code runs inside thewithblock, the interesting part already happened. - Callbacks run off the event loop. They are dispatched from the pool's background thread, in registration order, and never sit between your script and its tasks. A slow callback delays the next event, not the training run.
- A callback that raises is reported and skipped. The exception is printed to stderr and the stream carries on. A broken observer must not take the run with it.
The view¶
The second argument is the fold: every event that has passed, accumulated into one immutable ComputeView. It answers the question an event alone cannot — a ConsoleEvent does not know whether the compute is ready; the view does.
compute.state # "requested" → "provisioning" ⇄ "ready" → ... one ComputeState word, moved by the same table the daemon uses
compute.cost # accrued so far, from the cost gauge
compute.errors # tuple[str, ...] — what has gone wrong, most recent last
compute.nodes # tuple[NodeView, ...]
compute.tasks # tuple[TaskView, ...]
compute.nodes_ready # counted from the rows
Each NodeView carries the machine's lifecycle state, its address, price_per_hour and market once the API has said them, the bootstrap checklist as phases (each PhaseView named and either underway, done, or broken — the pip installs your plugins asked for show up here), a short metrics history per gauge, the tail of what the node printed, and the progress line a machine reports while it is still short of an address. Each TaskView carries the task's state, the function's real name, and its timings.
The view is a frozen dataclass: a callback that keeps a reference keeps that moment, and two callbacks handed the same event are handed the same value. Its windows are bounded — the last 40 tail lines, 12 samples per gauge, 32 error messages, the latest 200 tasks — so a compute that stays up for days never grows the view with it. The full history is the event log itself, which sky log export hands you whole.
The view is fed from both directions. Events move what events can say; the fields only the API carries — a node's address, a task's submitted_at, the spec the compute was created from — are hydrated by reads made beside the stream rather than in its way. A node or task transition asks for a read, a burst of them shares one, no two reads start less than two seconds apart, and none more than about ten, because some of what the API carries — a machine given an address, a task queued behind a busy node — moves without an event. So a callback can be handed an event before the read it asked for has landed: a task's state is in the view at once, its timings in a later one. What a read never does is take back a state an event has already moved.
One stream, every watcher¶
The pool holds one SSE connection per compute, whoever is watching. The live console, the Rich dashboard, and all of your callbacks hang off the same consumer: the stream is read once, folded once, and each subscriber is handed the same view. Registering five callbacks costs five function calls per event, not five connections.
The connection is nursed. Every frame carries the log's global sequence, so when the transport drops — a daemon bounce, a flaky network — the client reconnects with Last-Event-ID and resumes exactly after the last event it delivered, never repeating one. What a gap can cost you are the published gauges (cost, metrics, progress) that fell inside it, which do not replay by design. A daemon that stays silent past the client's retry window is a different situation, and then the error is the answer.
The raw stream¶
Callbacks are sugar over an iterator. When you would rather own the loop — a dashboard's render cycle, an asyncio app bridging events into its own queue — pool.events() hands you the decoded stream directly:
with sky.Compute.attached("training") as pool:
for event in pool.events():
match event:
case sky.TaskEvent(task=task, state="failed"):
alert(task)
case sky.ComputeDeleted():
break
The iterator replays the log and then follows it, with the same reconnect behaviour, for as long as you keep reading — the stream has no end while the compute exists, so the consumer decides when to stop, and a break closes the connection. Each events() call is its own subscription with its own connection; the shared one belongs to the callbacks and the console.
Building on it¶
Compute.attached takes callbacks= too, which is what turns this into an integration surface rather than a logging convenience: a process that did not create the compute can join it and watch it.
A monitoring sidecar, a Slack bot announcing failed tasks, a cost widget on an internal dashboard — each is an attach, a callback, and a match over the cases it cares about. The daemon owns the compute either way; watchers come and go.
Next steps¶
- Events — the wire-level reference: every event's payload, recorded vs published, replay and cursors
- Compute and task dispatch — the pool's full constructor surface
- Core Concepts — the programming model the events narrate
- CLI —
sky monitorandsky log export, the same stream from a shell