The runner read unconsumed_signals (a DB query) on every frame including each
ephemeral token, so stop polling scaled with the token stream. A cooperative
stop is only honored at vertex/milestone boundaries, so poll only on durable
frames; the post-loop check still catches a STOP that lands after the last
frame.
The v2 stop handler called update_job_status(CANCELLED) unconditionally, so
stopping an already-COMPLETED/FAILED/TIMED_OUT job flipped it to CANCELLED while
its result/error blob stayed put. Return the existing terminal state instead of
transitioning when the row is already finished.
stop() cancelled in-flight tasks but returned without awaiting them, so a job's
shielded terminal reconcile (a DB write) could race a closing engine and leave a
'Task was destroyed but it is pending' warning. Cancel and gather the in-flight
job tasks first, then tear down the workers. Cancelling the job tasks directly
(rather than relying on cancellation propagating through the worker's await)
avoids the worker sitting in cancellation limbo behind a job that swallows its
cancel on the user-stop path, which otherwise stalls teardown.
unconsumed_signals filters consumed_at IS NULL but nothing stamped it, so STOP
rows lingered: the table grew and a re-enqueued job self-cancelled off a stale
STOP. Add JobService.consume_signals and stamp the STOP in the runner's
_reconcile_stop, the single point where a stop is finalized.
Two uvicorn workers booting against the same DB both re-enqueued the same QUEUED
row, so a non-idempotent flow executed twice. Add JobService.claim_queued_job, a
conditional UPDATE ... WHERE status='QUEUED' that only one racer can win
(rowcount==1), and gate sweep re-enqueue on it. Works on SQLite and Postgres.
The background_job_timeout setting was read by nobody, so a runaway run never
timed out. Bound the runner's drive with asyncio.wait_for(timeout=...) when the
setting is configured; execute_with_status already maps asyncio.TimeoutError to
TIMED_OUT, so the run ends TIMED_OUT with a terminal event. None keeps the prior
unbounded behaviour. The facade passes settings.background_job_timeout through.
events() decided 'is the job finished' from the process-local live bus, whose
_closed marker is empty after a restart. A reattach to an already-terminal job
replayed durable rows then blocked forever on the live tail. Key the decision off
the persisted job status instead: when terminal, replay the durable log and
return. Re-frame durable rows through the same format_sse_event formatter the
live path uses (with id=str(seq)) so replayed and live frames are byte-compatible
and Last-Event-ID resumes correctly.
The worker used asyncio.Task.cancelling() to tell a job-cancel apart from a
worker-cancel, but that API is Python 3.11+. On 3.10 the first job cancellation
raised AttributeError inside the cancel handler and permanently killed the
worker, hanging later jobs QUEUED. Replace it with an explicit _closed flag set
by stop(): a CancelledError while shutting down re-raises to exit the worker,
otherwise it is a per-job cancel that is swallowed so the pool keeps serving.
Give v2-workflows sync and the langflow stream protocol one parser. The
stream now emits a normalized "output" event per terminal output carrying
an OutputEvent (the ComponentOutput shape sync returns in outputs[id], plus
component_id). A shared build_component_output() backs both the sync
converter and the adapter, and the build loop ships authoritative vertex
metadata as an additive output_meta key on end_vertex (existing consumers
read build_data and ignore it).
This is access-pattern parity (one parser, same fields, same terminal set),
not byte-identical content: the stream reuses the v1 build path whose
display serialization differs from sync's run_graph output.
Let a sync caller name the output(s) they want via output_ids so
output.text resolves deterministically (reason=single) on multi-output
flows instead of going null. Selection is steer-only: it picks the
answer among the named outputs without filtering the outputs map.
Invalid ids are rejected with 422 before the flow runs (and before any
job row is created), so a typo costs no compute. Resolution considers
selected outputs that actually fired, so branching flows resolve to
whichever candidate ran.
Replace the flat output_text shortcut with an `output` object carrying the
resolved text answer plus a `reason` that explains why it resolved that way
(single/multiple/none/non_string/failed), so a null answer is always
diagnosable instead of silently None. `reason` follows the LLM-domain
finish_reason/stop_reason convention, distinct from the lifecycle status.
Also add `display_name` to each ComponentOutput (the stable component id
stays the dict key) and a computed `has_errors` flag derived from errors.
Pin the sync-response shortcuts on the v2 workflows endpoint:
- output_text surfaces the lone ChatOutput/TextOutput text and stays None for
non-output message nodes, data-only flows, and multi-text flows
- session_id echoes the resolved session; the error response exposes neither
- each outputs entry exposes only {type, status, content, metadata}, with the
component id carried by the dict key
Also drop the component_id kwarg the converter passed to ComponentOutput, which
has no such field and silently dropped it.
The synchronous /api/v2/workflows response keyed every result under its
component id, so reading the answer meant knowing an id you can't predict.
Surface two additive fields:
- output_text: the flow's single text answer (ChatOutput/TextOutput). None
when the flow has zero or multiple text outputs, so callers read outputs
rather than the shortcut guessing which channel is the answer.
- session_id: echoes the resolved session so chat/memory callers can
continue the same thread (v1 /run returned this; v2 had dropped it).
outputs is unchanged, so this is non-breaking.
Rebased onto release-1.10.0. The base independently rebuilt the v2
workflows backend (RBAC, body globals, share-aware fetch); keep our
forward design and conform its auth to that work:
1. Auth: keep get_current_user_for_workflow (session-or-API-key authN
that does not hold a DB connection during the inline run, avoiding
the SQLite lock contention api_key_security would cause) and enforce
the base's RBAC on top: ensure_flow_permission(EXECUTE) before run,
(READ) before status reconstruct, with widen_for_shares fetch.
2. Port the base's request-body globals onto the v2 WorkflowRunRequest.
The X-LANGFLOW-GLOBAL-VAR-* headers stay supported (the Responses API
passes globals that way); body globals win on conflict. Converters
echo the effective globals via effective_globals.
3. Public endpoint keeps the v1 build_public_tmp posture
(access_type==PUBLIC, run-as-owner); RBAC applies to the
authenticated endpoint only.
4. Preserve the base's post-build KB-cache invalidation in the AG-UI
build path.
The endpoint, AG-UI bridge, pluggable stream adapters, public endpoint,
and re-attach are unchanged.
* docs: remove 1.8 env vars
* docs: link out to deployment guide
* docs: redis queue
* fix-links-and-combine-env-vars-table
* docs: cleanup
* peer-review
* peer-review
* Apply suggestions from code review
Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>
* move-prereqs
---------
Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>
* fix: remove src/backend/base/uv.lock and Dockerfile references
- Deleted src/backend/base/uv.lock (monorepo should have only one uv.lock at root)
- Removed COPY ./src/backend/base/uv.lock lines from all Dockerfiles:
- docker/build_and_push.Dockerfile
- docker/build_and_push_base.Dockerfile
- docker/build_and_push_ep.Dockerfile
- docker/build_and_push_with_extras.Dockerfile
- docker/dev.Dockerfile
Fixes LE-1093
* fix: remove stale base uv.lock regeneration paths