mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-08-04 03:49:25 +00:00
`run_outcome_mapping.py` called itself "not a port implementation" and sat in a package of secondary adapters, while the half that actually invoked the use case lived as a closure in the composition root. It is one thing, and it is a primary adapter: the run runtime calls it the way HTTP calls the router and the clock calls the poller. `ScheduleRunCompletionListener` now holds the whole responsibility -- decide whether a finished run is ours, translate it, invoke the use case. Those are not two jobs: "ignore this run" is only meaningful as "do not call the service", so splitting them is what left the second half in a place where behaviour is not asserted. `build_run_completion_hook` drops to `return ScheduleRunCompletionListener(service)`. The composition root's own docstring says no adapter logic lives there; that is now true of it as well as of the routers it was written about. Placement --------- Kept in `app/adapters/schedule/` rather than moved beside the other two primary adapters. The context stays in one package; direction is stated by the class name and each module's first line, and the package `__init__` -- previously empty -- now lists which of its modules point which way, so a file added without that line is visibly a file whose direction nobody decided. A subdirectory for a single inbound module would have made the other four look like they had been sorted into something. Tests ----- This is the part that was not a rename. The conversion had 24 cases; the invocation had none, because the composition root is not where behaviour is asserted, so nothing covered "an ordinary chat run must not reach the service" as opposed to "produces no outcome object". The cases now drive `__call__` against a recording service, which asserts the same mappings plus what was done with them, and adds the two that were unreachable before: the service left entirely alone for a filtered run, and the completion stamped with a tz-aware current instant. `test_composition.py` gains `TestRunCompletionHook` for the assembly decision that remains -- including that the hook is bound to the service it was given, which a wrong wiring would type-check past. Confirmed by mutation that this case fails when the binding is broken. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
103 lines
4.1 KiB
Python
103 lines
4.1 KiB
Python
"""Primary adapter (inbound) -- the run runtime's completion callback.
|
|
|
|
The one file in this package that drives the domain rather than serving it.
|
|
Its siblings are secondary adapters the service calls out to; this one is
|
|
called *by* the run runtime, the same way the router is called by HTTP and the
|
|
poller by its clock. Kept here rather than beside those two so the context
|
|
stays in one place, with the direction stated by the name and by this line.
|
|
|
|
Its job is the filtering the legacy completion hook did inline. Every run in
|
|
the process reaches that callback, so most of them are none of this context's
|
|
business, and producing no call at all says exactly that -- not an error, just
|
|
nothing to write back. That is why `ScheduleService.handle_run_completion`
|
|
carries no guard clauses and never imports ``RunRecord``.
|
|
|
|
TODO(hexagonal): this depends on ``RunRecord``, a run-runtime type, rather
|
|
than on a contract published by the run context -- that context has not been
|
|
through a hexagonal slice yet. When it publishes one (a DTO, not its aggregate
|
|
and not its repository), replace ``_to_outcome``. Nothing else moves.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import UTC, datetime
|
|
from typing import TYPE_CHECKING
|
|
|
|
from deerflow.domain.schedule.model import RunStatus
|
|
from deerflow.domain.schedule.ports import RunOutcome
|
|
|
|
if TYPE_CHECKING:
|
|
from deerflow.domain.schedule.service import ScheduleService
|
|
from deerflow.runtime import RunRecord
|
|
|
|
# The runtime reports four terminal states; the domain has three, because
|
|
# `timeout` and `error` are the same fact to a scheduled task while
|
|
# `interrupted` is deliberately not -- a cancel or same-thread takeover ends
|
|
# the task CANCELLED, not FAILED.
|
|
_TERMINAL_STATUSES = {
|
|
"success": RunStatus.SUCCESS,
|
|
"error": RunStatus.FAILED,
|
|
"timeout": RunStatus.FAILED,
|
|
"interrupted": RunStatus.INTERRUPTED,
|
|
}
|
|
|
|
_INTERRUPTED_WITHOUT_ERROR = "run was interrupted before completion"
|
|
|
|
|
|
class ScheduleRunCompletionListener:
|
|
"""Turns a finished run into the write-back use case, or into nothing.
|
|
|
|
Deciding whether a run is ours and invoking the use case are one
|
|
responsibility, not two: "ignore this run" is only meaningful as "do not
|
|
call the service", so splitting them left the second half living in the
|
|
composition root, where behaviour is not asserted.
|
|
"""
|
|
|
|
def __init__(self, service: ScheduleService) -> None:
|
|
self._service = service
|
|
|
|
async def __call__(self, record: RunRecord) -> None:
|
|
outcome = self._to_outcome(record)
|
|
if outcome is None:
|
|
return
|
|
await self._service.handle_run_completion(outcome, now=datetime.now(UTC))
|
|
|
|
@staticmethod
|
|
def _to_outcome(record: RunRecord) -> RunOutcome | None:
|
|
"""Translate into domain vocabulary, or `None` to ignore the run.
|
|
|
|
`None` when the run is not a scheduled execution (no usable task
|
|
metadata, no owner) or has not reached a terminal state yet.
|
|
"""
|
|
metadata = record.metadata or {}
|
|
task_id = metadata.get("scheduled_task_id")
|
|
record_id = metadata.get("scheduled_task_run_id")
|
|
user_id = record.user_id
|
|
# `metadata` is a free-form dict a caller can influence, so the ids are
|
|
# type-checked rather than assumed; `user_id` is required because every
|
|
# task read is scoped by it.
|
|
if not isinstance(task_id, str) or not isinstance(record_id, str) or not user_id:
|
|
return None
|
|
|
|
status = _TERMINAL_STATUSES.get(str(record.status.value))
|
|
if status is None:
|
|
return None
|
|
|
|
if status is RunStatus.SUCCESS:
|
|
# A stale error left on a successful record must not be written back
|
|
# as the task's last_error.
|
|
error = None
|
|
elif status is RunStatus.INTERRUPTED:
|
|
error = record.error or _INTERRUPTED_WITHOUT_ERROR
|
|
else:
|
|
error = record.error
|
|
|
|
return RunOutcome(
|
|
task_id=task_id,
|
|
record_id=record_id,
|
|
run_id=record.run_id,
|
|
user_id=user_id,
|
|
status=status,
|
|
error=error,
|
|
)
|