mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-07-27 08:28:00 +00:00
* fix(runs): cancel degrades to lease takeover for multi-worker Work item 4 of the multi-worker ownership epic (https://github.com/bytedance/deer-flow/issues/3948). Problem: POST /runs/{run_id}/cancel landing on a non-owning worker returns 409 — the cancel button silently fails under GATEWAY_WORKERS>1 with no sticky routing. cancel() required the current worker to hold the in-memory task/abort_event, which any non-owner pod cannot satisfy. Changes: - RunManager.cancel() returns CancelOutcome enum (cancelled / taken_over / lease_valid_elsewhere / not_active_locally / not_cancellable / unknown) instead of bool, so the router can map each outcome to the right HTTP response. - New store primitive claim_for_takeover(): a single atomic conditional UPDATE that marks a run as error only when status IN (pending, running) AND (lease IS NULL OR lease < now - grace). Closes the stale-read / concurrent-heartbeat race — if the owner renews between our read and write, the UPDATE matches 0 rows and we surface lease_valid_elsewhere. - HTTP cancel + stream-join endpoints route on CancelOutcome: cancelled -> 202 (or 204 with wait=true); taken_over -> 202 immediately (no SSE streaming — the run is terminal on another worker, streaming would hang); lease_valid_elsewhere -> 409 + Retry-After header computed from lease_expires_at + grace_seconds. - RunManager.grace_seconds exposed as a public property; the router no longer reaches into _run_ownership_config. - _is_lease_expired extracted to a module-level function, shared by RunManager.cancel() and MemoryRunStore.claim_for_takeover(). - GATEWAY_WORKERS=1 + heartbeat_enabled=false is zero-regression: the non-local path short-circuits to not_active_locally, preserving the original 409 behaviour the existing tests pin. Tests: 12 new (5 store primitive + 4 cancel-takeover unit + 3 HTTP including a regression guard verifying POST /stream?action=interrupt on a dead-owner run returns 202 instead of hanging on SSE). 244 directly-related tests pass; 36/36 blocking-IO gate pass. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com> * fix(runs): guard update_status and self-terminate on takeover Two defenses close a split-brain window where the original owner could overwrite a peer's takeover status: - update_status (SQL + memory store) now guards on status IN ('pending','running'). When takeover already set the row to 'error', the owner's final status write matches 0 rows and is dropped. - _persist_status: when update_status returns False, check whether the row exists before attempting recovery via put(). If the row exists (takeover by another worker), skip recovery instead of blindly upserting over the takeover. - Heartbeat _renew_leases: when update_lease returns False (row no longer pending/running or owner changed), cancel the local task so wasted CPU is bounded to the next heartbeat tick (~10s) instead of the full task lifetime. Also fix three reviewer feedback items: - Re-fetch the store row when cancel() returns lease_valid_elsewhere, so Retry-After uses the owner's freshly-renewed lease instead of a stale value from request start. - Fallback 'unknown' in takeover error message when owner_worker_id is NULL (pre-ownership data). - Remove dead else-10 branch from grace_seconds property (unreachable — all callers are downstream of the heartbeat_enabled guard). Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com> * test(runs): pin split-brain defences from update_status guard + heartbeat Three tests lock down the takeover authoritativeness so a late-running owner cannot overwrite a peer's claim: - update_status must reject writes when the store row is already terminal (taken over by another worker). - _persist_status must skip row-recovery via put() when the row exists but has been taken over. - Heartbeat _renew_leases must cancel the local task when update_lease returns False (row claimed by another worker). Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com> * fix(runs): precise outcome + log when local cancel loses to peer takeover Two reviewer precision nits on the split-brain defence: - _persist_status: branch the skip-reason log on existing["status"]. error → WARNING "peer takeover" (anomalous); interrupted/success → INFO "local cancel/completion race" (expected when user hits stop as the run finishes). Stops noisy false-positive takeover warnings in operator logs. - cancel() local path: when _persist_status returns False, re-check the store. If a peer's claim_for_takeover flipped the row to error between our in-memory cancel and the guarded update_status, surface taken_over instead of cancelled so the client sees a status consistent with the store. Test: test_cancel_returns_taken_over_when_peer_claims_during_local_cancel pins the race outcome. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com> * fix(runs): widen update_status guard, de-duplicate lease helpers, add coverage Round 3 of reviewer feedback: - Widen update_status guard to status IN ('pending','running','interrupted'). The original guard blocked interrupted→error (the rollback finalize path), losing the "Rolled back by user" message. interrupted is now permitted while error/success stay locked — takeover protection unchanged. - claim_for_takeover False now re-reads the store row to distinguish causes: owner renewed lease → lease_valid_elsewhere; row went terminal → not_cancellable; another worker already took it over → taken_over. - Extract _raise_lease_valid_elsewhere() helper to de-duplicate the 409+Retry-After block shared across cancel_run and stream_existing_run. - Extract _lease_expired_or_null() in persistence/run/sql.py to de-duplicate the lease-expiry SQL WHERE clause shared by claim_for_takeover and list_inflight_with_expired_lease. - 11 new tests: 5 SQL-layer claim_for_takeover (expired/valid/NULL/ terminal/nonexistent), 3 _compute_retry_after unit (NULL/unparseable/ normal), 2 claim re-read precision (terminal/takeover), 1 stream endpoint 409+Retry-After. Not addressed (non-blocking, reviewer agreed): - The 2–3 store.gets in the takeover cold path: optimizing the API to accept a pre-fetched record would couple the router to the manager more tightly than justified by the perf gain. - The lease-expiry inline loop in MemoryRunStore.list_inflight_with_- expired_lease pre-computes cutoff once for all rows; switching to the shared _is_lease_expired helper would recompute datetime.now() per row with no real benefit. 260 related tests pass; 36/36 blocking-IO gate pass; ruff clean. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com> * fix(runs): de-duplicate lease-expiry helper, restore defensive fallback Address final round of review feedback: - Extract is_lease_expired to deerflow.utils.time (no _ prefix, public utility). Manager and MemoryRunStore now import from the same place instead of the store reaching backward into the manager for a private function. - Restore defensive else-10 fallback in grace_seconds property (removed in an earlier round). The guard is unreachable for current callers but protects future ones from AttributeError. - Comment the transient in-memory interrupted vs store error state when a local cancel is superseded by a peer takeover. - Comment the max(1, ...) floor in _compute_retry_after — the floor is a lower bound, not a poll interval; clients should apply jitter. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.7 <noreply@anthropic.com> Co-authored-by: rayhpeng <rayhpeng@gmail.com>
96 lines
3.6 KiB
Python
96 lines
3.6 KiB
Python
"""ISO 8601 timestamp helpers for the Gateway and embedded runtime.
|
|
|
|
DeerFlow stores and serializes thread/run timestamps as ISO 8601 UTC
|
|
strings to match the LangGraph Platform schema (see
|
|
``langgraph_sdk.schema.Thread``, where ``created_at`` / ``updated_at``
|
|
are ``datetime`` and JSON-encode to ISO 8601). All timestamp generation
|
|
should funnel through :func:`now_iso` so the wire format stays
|
|
consistent across endpoints, the embedded ``RunManager``, and the
|
|
checkpoint metadata written by the Gateway.
|
|
|
|
:func:`coerce_iso` provides a forward-compatible read path for legacy
|
|
records that historically stored ``str(time.time())`` floats.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import re
|
|
from datetime import UTC, datetime, timedelta
|
|
|
|
__all__ = ["coerce_iso", "is_lease_expired", "now_iso"]
|
|
|
|
|
|
def is_lease_expired(lease_expires_at: str | None, *, grace_seconds: int) -> bool:
|
|
"""Return ``True`` when *lease_expires_at* has elapsed past grace.
|
|
|
|
A NULL lease (pre-ownership data) is always considered expired so
|
|
take-over (cancel from a non-owning worker) can reclaim it in the
|
|
same way reconciliation does. Unparseable timestamps are also
|
|
treated as expired (defence in depth).
|
|
"""
|
|
if lease_expires_at is None:
|
|
return True
|
|
try:
|
|
dt = datetime.fromisoformat(lease_expires_at)
|
|
if dt.tzinfo is None:
|
|
dt = dt.replace(tzinfo=UTC)
|
|
except (ValueError, TypeError):
|
|
return True
|
|
return dt < datetime.now(UTC) - timedelta(seconds=grace_seconds)
|
|
|
|
|
|
_UNIX_TIMESTAMP_PATTERN = re.compile(r"^\d{10}(?:\.\d+)?$")
|
|
"""Matches the unix-timestamp string shape historically written by
|
|
``str(time.time())`` (10-digit seconds with optional fractional part).
|
|
The 10-digit anchor avoids accidentally rewriting ISO years like
|
|
``"2026"`` and stays valid until the year 2286.
|
|
"""
|
|
|
|
|
|
def now_iso() -> str:
|
|
"""Return the current UTC time as an ISO 8601 string.
|
|
|
|
Example: ``"2026-04-27T03:19:46.511479+00:00"``.
|
|
"""
|
|
return datetime.now(UTC).isoformat()
|
|
|
|
|
|
def coerce_iso(value: object) -> str:
|
|
"""Best-effort coerce a stored timestamp to an ISO 8601 string.
|
|
|
|
Translates legacy unix-timestamp floats / strings written by older
|
|
DeerFlow versions into ISO without a one-shot migration. ISO strings
|
|
pass through unchanged; ``datetime`` instances are normalised to UTC
|
|
(tz-naive values are assumed to be UTC) and emitted via
|
|
``isoformat()`` so the wire format always uses the ``T`` separator;
|
|
empty values become ``""``; unrecognised values are stringified as a
|
|
last resort.
|
|
"""
|
|
if value is None or value == "":
|
|
return ""
|
|
if isinstance(value, bool):
|
|
# ``bool`` is a subclass of ``int`` — treat as garbage, not 0/1.
|
|
return str(value)
|
|
if isinstance(value, datetime):
|
|
# ``datetime`` must be handled before the ``int``/``float`` check;
|
|
# str(datetime) would produce ``"YYYY-MM-DD HH:MM:SS+00:00"``
|
|
# (space separator), which breaks strict ISO 8601 consumers.
|
|
if value.tzinfo is None:
|
|
value = value.replace(tzinfo=UTC)
|
|
else:
|
|
value = value.astimezone(UTC)
|
|
return value.isoformat()
|
|
if isinstance(value, (int, float)):
|
|
try:
|
|
return datetime.fromtimestamp(float(value), UTC).isoformat()
|
|
except (ValueError, OverflowError, OSError):
|
|
return str(value)
|
|
if isinstance(value, str):
|
|
if _UNIX_TIMESTAMP_PATTERN.match(value):
|
|
try:
|
|
return datetime.fromtimestamp(float(value), UTC).isoformat()
|
|
except (ValueError, OverflowError, OSError):
|
|
return value
|
|
return value
|
|
return str(value)
|