From c573e4bfee0851880cf9eded2cc56870b7f25608 Mon Sep 17 00:00:00 2001 From: ogabrielluiz Date: Wed, 3 Jun 2026 23:15:46 -0300 Subject: [PATCH] feat(jobs): add append_event (per-job seq) and read_events(after_seq) to JobService --- .../base/langflow/services/jobs/service.py | 36 ++++++++++++++- .../unit/services/jobs/test_jobs_service.py | 46 +++++++++++++++++++ 2 files changed, 81 insertions(+), 1 deletion(-) diff --git a/src/backend/base/langflow/services/jobs/service.py b/src/backend/base/langflow/services/jobs/service.py index 4e2ddb196c..33bdabbd39 100644 --- a/src/backend/base/langflow/services/jobs/service.py +++ b/src/backend/base/langflow/services/jobs/service.py @@ -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. diff --git a/src/backend/tests/unit/services/jobs/test_jobs_service.py b/src/backend/tests/unit/services/jobs/test_jobs_service.py index c84508d20e..745f669349 100644 --- a/src/backend/tests/unit/services/jobs/test_jobs_service.py +++ b/src/backend/tests/unit/services/jobs/test_jobs_service.py @@ -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]