mirror of
https://github.com/langflow-ai/langflow.git
synced 2026-07-24 08:57:31 +08:00
feat(jobs): add append_event (per-job seq) and read_events(after_seq) to JobService
This commit is contained in:
@ -17,7 +17,7 @@ from langflow.services.database.models.jobs.crud import (
|
||||
get_latest_jobs_by_asset_ids,
|
||||
update_job_status,
|
||||
)
|
||||
from langflow.services.database.models.jobs.model import Job, JobStatus, JobType
|
||||
from langflow.services.database.models.jobs.model import Job, JobEvent, JobStatus, JobType
|
||||
from langflow.services.deps import session_scope
|
||||
from langflow.services.jobs.exceptions import DuplicateJobError
|
||||
|
||||
@ -229,6 +229,40 @@ class JobService(Service):
|
||||
await session.flush()
|
||||
return job
|
||||
|
||||
async def append_event(self, job_id: UUID, event_type: str, payload: dict) -> int:
|
||||
"""Append a durable event for a job and return its per-job seq.
|
||||
|
||||
seq is assigned as max(existing seq for job) + 1. UNIQUE(job_id, seq)
|
||||
on the table guards against concurrent double-assignment — a colliding
|
||||
writer raises IntegrityError, which the runner treats as a retryable
|
||||
append.
|
||||
"""
|
||||
async with session_scope() as session:
|
||||
stmt = select(func.max(JobEvent.seq)).where(JobEvent.job_id == job_id)
|
||||
result = await session.exec(stmt)
|
||||
current_max = result.one()
|
||||
next_seq = (current_max or 0) + 1
|
||||
event = JobEvent(job_id=job_id, seq=next_seq, event_type=event_type, payload=payload)
|
||||
session.add(event)
|
||||
await session.flush()
|
||||
return next_seq
|
||||
|
||||
async def read_events(self, job_id: UUID, after_seq: int = 0) -> list[JobEvent]:
|
||||
"""Return durable events for a job with seq > after_seq, ordered by seq.
|
||||
|
||||
``after_seq`` is the SSE Last-Event-ID cursor; pass 0 to read from the
|
||||
start.
|
||||
"""
|
||||
async with session_scope() as session:
|
||||
stmt = (
|
||||
select(JobEvent)
|
||||
.where(JobEvent.job_id == job_id)
|
||||
.where(JobEvent.seq > after_seq)
|
||||
.order_by(col(JobEvent.seq).asc())
|
||||
)
|
||||
result = await session.exec(stmt)
|
||||
return list(result.all())
|
||||
|
||||
async def get_latest_jobs_by_asset_ids(self, asset_ids: Sequence[UUID | str]) -> dict[UUID, Job]:
|
||||
"""Get the latest job for each asset ID in a single batch query.
|
||||
|
||||
|
||||
@ -65,3 +65,49 @@ async def test_set_error_persists_blob():
|
||||
async def test_set_result_returns_none_for_missing_job():
|
||||
service = JobService()
|
||||
assert await service.set_result(uuid4(), {"x": 1}) is None
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("client")
|
||||
async def test_append_event_assigns_monotonic_seq():
|
||||
service = JobService()
|
||||
job_id = uuid4()
|
||||
await service.create_job(job_id=job_id, flow_id=uuid4(), user_id=uuid4())
|
||||
|
||||
seq1 = await service.append_event(job_id, "run_started", {"a": 1})
|
||||
seq2 = await service.append_event(job_id, "vertex_started", {"b": 2})
|
||||
seq3 = await service.append_event(job_id, "run_finished", {"c": 3})
|
||||
|
||||
assert (seq1, seq2, seq3) == (1, 2, 3)
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("client")
|
||||
async def test_append_event_seq_is_per_job():
|
||||
service = JobService()
|
||||
job_a = uuid4()
|
||||
job_b = uuid4()
|
||||
await service.create_job(job_id=job_a, flow_id=uuid4(), user_id=uuid4())
|
||||
await service.create_job(job_id=job_b, flow_id=uuid4(), user_id=uuid4())
|
||||
|
||||
assert await service.append_event(job_a, "x", {}) == 1
|
||||
assert await service.append_event(job_b, "x", {}) == 1
|
||||
assert await service.append_event(job_a, "x", {}) == 2
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("client")
|
||||
async def test_read_events_after_seq_returns_ordered_tail():
|
||||
service = JobService()
|
||||
job_id = uuid4()
|
||||
await service.create_job(job_id=job_id, flow_id=uuid4(), user_id=uuid4())
|
||||
|
||||
await service.append_event(job_id, "e1", {"i": 1})
|
||||
await service.append_event(job_id, "e2", {"i": 2})
|
||||
await service.append_event(job_id, "e3", {"i": 3})
|
||||
|
||||
tail = await service.read_events(job_id, after_seq=1)
|
||||
assert [e.seq for e in tail] == [2, 3]
|
||||
assert [e.event_type for e in tail] == ["e2", "e3"]
|
||||
assert tail[0].payload == {"i": 2}
|
||||
|
||||
# after_seq=0 (or default) returns everything in order.
|
||||
all_events = await service.read_events(job_id, after_seq=0)
|
||||
assert [e.seq for e in all_events] == [1, 2, 3]
|
||||
|
||||
Reference in New Issue
Block a user