Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Weft The book

The ctx

ctx is everything your node’s code can reach: its inputs, and every service weft provides. The runtime hands it to run, and to a trigger’s setup_trigger.

Every call here returns WeftResult, and your node never names a weft error type. Wrap an outside error with .node_err("what you were doing"), and fail on something you detected yourself with node_bail!.

Six rules that save you a debugging session

Calls on a connection go through ctx.open or ctx.client. That is what signs the request and records what it cost. ctx.http() is for everything else. A reqwest client you built yourself skips both.

ctx.await_signal is refused once your node has emitted anything, because a resume replays the body from the top and would emit again. It is also refused on a node with a Generator input, because a stream cannot be replayed.

ctx.run keys on call order, not on the name you gave it. The name is for reading logs. If your body can call it a different number of times on a replay, wrap the thing that varies.

put with no keep policy is run-scoped, and the file is swept shortly after the run ends. Right for scratch, wrong for anything a person will open later.

ctx.stop_tagged returns when the stop is queued, not when the other runs are gone. Do not write a node whose next step assumes they stopped.

ctx.register_signal is once per node per setup. A second call is a loud error rather than a trigger that quietly never fires.

Which firing am I

FieldTypeWhat it is
colorColor (a UUID)This run’s id
project_idUuidThe project it belongs to
node_idStringThis node’s id in the graph
node_typeStringIts catalog type, such as ExecPython
node_labelOption<String>The title shown on the box, if it has one
framesLoopFramesWhich loop iterations this firing sits inside
member()Option<&MemberId>Who this run is for, when it is for one member of the program (see programs with members)

Reading inputs

ctx.inputs holds this firing’s inputs, whether they arrived on a wire, were written into the body, or came from a declared default. ctx.wake holds the event payload, and only on the trigger that actually fired. Both are a ValueBag with the same methods.

CallWhat you getIt fails when
get::<T>("name")A required valueIt is absent, or it does not fit T
opt::<T>("name")Option<T>; absent or null gives NoneA value is there and does not fit T
get_or::<T>("name", default)The value, or your defaultA value is there and does not fit T
list::<T>("name")Vec<T>; absent gives [], and one value gives a one-item listAny item does not fit T
access("name")The connection picked on an access inputNobody picked one and the service requires it
raw("name")Option<&Value>, the JSON as it arrivednever
object() / record()The whole bag as one recordOn a wake bag whose event was missing or not an object
nested("name")The object under name, as its own bagThe value is there and is not an object
iter()Every named valuenever
custom()Only the ports this instance added, not the type’s own settingsnever
declared()Only the type’s own declared settingsnever
in_order()Every input that arrived, in source orderOn a wake or nested bag
declares("name")Whether the node declares this input at allnever
for_holes(holes, what, spell, declare)Values matched to named holes in a text, such as $user_id in SQLA hole names no port, or a wired port the text never reads

Never write get(...).unwrap_or(...). It turns a real type error into your default and you will not find out. get_or is the same shape and fails honestly.

Emitting

CallWhat it doesIt fails when
pulse_downstream(output).awaitSends values on. The only way a node emitsA port is mentioned twice, is not declared, or gets a value its type does not allow
yield_downstream(output).awaitThe same, but waits until the value was takenSame, plus a delivery that can never happen, and an empty output
close_port("port").awaitCloses one output. On a stream port this is the end of the streamThe port was already mentioned this firing
set_max_buffered_items("port", n)Raises the un-taken cap on a stream port above 4096The port is not a stream output, or n is 0
fan_declared(&value)Builds an output by matching an object’s keys to your declared ports, skipping the restnever
output_type("port")The resolved type of one outputnever
declared_outputs() / declared_inputs()Every port this instance declares, with its typenever

An ordinary port emits at most once per firing. Whatever you never mention is closed for you when the body returns. A Generator[T] port emits as often as you like until you close it.

NodeOutput is the payload you hand those calls:

CallWhat it does
NodeOutput::new()An empty one
.set("port", value)Sets one port
.extend_from_object(&value)Fans every top-level key onto a same-named port. Last write wins
NodeOutput::stored_file(stored)The four ports a stored file travels as: file, filename, mimeType, sizeBytes
.get("port")Reads back what you already set

Waiting, and surviving a restart

CallWhat it doesIt fails when
await_signal(kind).awaitParks this firing until the signal fires, releasing the worker. Returns the payloadThe body already emitted, or the node has a stream input, or the body’s call order changed between replays
register_signal(kind).awaitSets up a trigger that starts a new run on every fire. Called during setupCalled twice for one node in one setup
run("name", closure).awaitRuns the closure and writes down its result, or gives back what was written down last timeThe closure itself fails

ctx.run is how you stop a replay from doing an expensive thing twice. It cannot stop it from doing it twice when the worker dies between the action and the write, so a service you call more than once has to cope with that itself.

Connections

CallWhat you getIt fails when
open(&access).awaitThe connection, leased for this firing, with its token refreshedIt needs reconnecting, the service does not match, or the door refuses
open_within(&access, window).awaitThe same, with a window you choose instead of the default 15 minutesSame
client(access).awaitStraight to the signed-in HTTP client. Takes None and gives you a plain oneSame, when it is Some
http()The shared plain HTTP client, no credentialsnever
publish_access(values).awaitPublishes a connection to something this node runs itself, such as a database it provisionedA field the service does not declare, or a required one missing
published_access().awaitWhat this node published before, or NoneThe broker cannot answer
endpoint("name").awaitA handle on one of this node’s declared infra endpoints, once something answers thereThe endpoint is not declared, or the infra is not running

On an opened connection: client(), credential(), value(name), identity(), owner(). On an endpoint handle: url(), host_and_port(), and call(method, path, body).

Storage

ctx.storage(scope) gives you a handle. The scope decides where new files go and what list sees; get, delete, keep and presign act on whatever scope the key itself belongs to, so a later node can read a file without knowing where it came from.

ScopeWhere it lives, and how long
Execution (the default)This run only, swept after it ends unless you keep it
ProjectOutlives runs, gone on weft clean or weft rm
Shared { name }Shared across projects that name the same space
AssetThe project’s published assets. Readable, and the worker refuses writes
CallWhat it does
identified(id)Names what the next put is a copy of, so storing it twice stores it once
put(bytes, mime, filename, keep)Stores bytes and gives back the stored-file value
put_stream(stream, mime, filename, keep)The same without holding the file in memory
put_response(resp, what, mime, filename, keep)Streams an HTTP response you already made into storage
put_from_url(url, filename, keep)Fetches a URL straight into storage
copy(&file, keep)Copies a file into this handle’s scope, leaving the original
get(&file)The file’s bytes, as a stream
get_range(&file, range)Part of the file
get_bytes(&file)The whole file in memory. Only for small ones
delete(&file)Deletes it. Stored files only, not URL-backed ones
list()Everything under this scope
keep(&file, ttl)Saves a run-scoped file from the sweep. There is no un-keep
presign(&file, ttl)A temporary link. None means the default, about 15 minutes
public_link(&file, ttl)A link the open internet can fetch, or None if this install serves none
caller_link(&file, ttl)A link a caller of this install can fetch
externalize(&value, &ty, policy)Turns every file inside a typed value into a link or inline bytes, for handing out
internalize(&value, &ty, keep)The reverse: pulls URLs and inline data into stored files

keep takes KeepTtl::Default (30 days, pushed back each time the file is read), Secs { secs }, or Never.

What comes out of externalize is no longer a value of that type, because the links expire. Hand it to whoever asked and do not store it.

Streams and buses

CallWhat it doesIt fails when
create_bus(opts)Makes a bus and gives you a handle and a marker to emitThe options do not make sense
bus(&marker)Turns a marker value back into a live handleThe bus is gone, or that is not a marker
open_bus("port", opts, "name").awaitThe whole producer move: make it, emit the marker, register your name. Closes on dropEmitting or registering fails
join_bus("port", "name")The consumer twin, also closing on dropThe input is not a live bus
bus_from_input("port")A handle that does not close the bus. For an observerSame

Read a stream input like any other value: ctx.inputs.get::<Generator<Row>>("rows")?.

CallWhat it does
next().awaitThe next item, waiting. None once, at a clean end
try_next()Item, Empty or Finished, without waiting
drain().awaitEverything, waiting for the end. A failed stream errors instead of returning a partial list
end()Whether the producer finished or failed, or None while it is open

Live callers

CallWhat you get
is_api_call() / is_websocket()Whether a caller of that kind is attached
caller_data_type()Whether this connection speaks bytes or JSON
caller()The caller as a handle, or None
http_caller().awaitThe HTTP caller, attached and connected, or a loud error
ws_caller().awaitThe same for a WebSocket
live_caller().awaitWhichever one is connected
caller_request()What the caller sent to open the exchange, without waiting for the connection

Logging, tags, cancelling

CallWhat it doesIt fails when
log(level, message).awaitWrites a line into the run’s log. Trace, Debug, Info, Warn, ErrorThe write fails
tag_execution(["user_7"]).awaitTags this run. Additive and safe to repeatThe list is empty, or a tag is not [A-Za-z0-9_-], 1 to 64 characters
stop_tagged("user_7", StopSelf::Keep).awaitQueues a stop for every live run of this project carrying that tagThe tag is invalid
is_cancelled()Whether this run was cancelled. Cheap enough to pollnever
cancellation()The flag itself, to select! against your own worknever

StopSelf::Keep only reaches runs that took the tag before this one did, so two runs racing to stop each other leave the later one alive. StopSelf::Include ends this run too.

Your project beyond this run

The runtime answers each of these calls on behalf of this run and writes the answer into the run’s journal, so a replay reads the answer back instead of asking again. Most of them take an optional .member(id), which picks that member’s copy instead of the shared one; values(), connections() and tokens().member(..) always name a member.

CallWhat it does
infra("bridge").start().awaitBrings the shared copy up. Returns once it runs, parking the run between looks; fails with the reason if it does not come up, or if somebody stops it while this waits
infra("bridge").stop(spec, stop_self).awaitScales the copy down, keeping its disk
infra("bridge").terminate(spec, stop_self).awaitDeletes the copy and its disk
infra("bridge").status().awaitThe copy’s state, None when there is none. A start or stop on its way reads provisioning or stopping at once, the same answer weft status gives
infra("bridge").copies().awaitEvery copy: the shared one and each member’s
trigger("receive").activate().awaitTurns one trigger on
triggers().deactivate(spec, stop_self).awaitTurns off every trigger of the program (or of one member, with .member(id)); if you want only some, name them with .only([..])
values().member(id).get().awaitWhat the member gave for the program’s @member_filled fields, by step and field
values().member(id).set(step, field, value).clear(step, field).apply().awaitGives and clears values in one change, each checked against its node’s rules; the member’s live triggers reading one are set up again, and their names come back
values().member(id).forget().awaitForgets everything the member gave
connections().member(id).list() / forget()Lists a member’s connections, or forgets all of them (the values naming them go too)
members().list().awaitEvery member weft holds anything for, one entry each: how many values, connections and live tokens, their infra copies, and their triggers, each with the events waiting on a field they have not filled and why
costs().member(id).service(s).since(t).list().awaitThe cost records, each saying whose credential paid, narrowed by member, node, service, run, paid_by or since (unix seconds)
runs().member(id).status(s).older_than(d).clean(running, stop_self).awaitDeletes runs; runs still going follow running
tokens().mint_for_member(id, expires_in).awaitA member token, value shown once
tokens().member(id).revoke().awaitRevokes a member’s tokens

For what spec and stop_self decide, go and read taking something down from a run.

What is not here

ContextHandle, the trait underneath, is the seam the runtime implements. Your node only ever sees the ExecutionContext wrappers on this page.