Skip to content

Commit 5ba8e44

Browse files
committed
Avoid loading job logs in dashboard APIs
Refactors dashboard job responses to separate summary payloads from full job details: most routes now return `JobSummaryResponse` without `stdout`/`stderr`, while only `GET /jobs/{job_id}` returns full logs. It also removes model-level deferred log columns in favor of query-time `defer(...)` options for list/locked queries, updates stream/job fetches accordingly, and adjusts tests to assert logs are persisted but not included in summary responses.
1 parent 54a9dbb commit 5ba8e44

4 files changed

Lines changed: 58 additions & 40 deletions

File tree

cea/interfaces/dashboard/lib/database/models.py

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,7 @@
55
from typing import Optional
66

77
from pydantic import AwareDatetime, computed_field
8-
from sqlalchemy import Index, Column, Text
9-
from sqlalchemy.orm import deferred
8+
from sqlalchemy import Index
109
from sqlmodel import Field, SQLModel, JSON, DateTime, BigInteger, select, inspect, text
1110

1211
import cea.scripts
@@ -97,8 +96,8 @@ class JobInfo(SQLModel, table=True):
9796
default_factory=get_current_time)
9897
start_time: Optional[AwareDatetime] = Field(sa_type=DateTime(timezone=True), nullable=True, default=None)
9998
end_time: Optional[AwareDatetime] = Field(sa_type=DateTime(timezone=True), nullable=True, default=None)
100-
stdout: Optional[str] = Field(default=None, sa_column=deferred(Column(Text), group='logs'))
101-
stderr: Optional[str] = Field(default=None, sa_column=deferred(Column(Text), group='logs'))
99+
stdout: Optional[str] = None
100+
stderr: Optional[str] = None
102101
project_id: str = Field(foreign_key="project.id", index=True)
103102
created_by: str = Field(foreign_key=f"{user_table_ref}.id", index=True)
104103
deleted_at: Optional[AwareDatetime] = Field(sa_type=DateTime(timezone=True), nullable=True, default=None, index=True)

cea/interfaces/dashboard/server/jobs.py

Lines changed: 50 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
import sqlalchemy.exc
1717
from fastapi import APIRouter, Depends, HTTPException, status, Request, Query
1818
from pydantic import BaseModel
19-
from sqlalchemy.orm import undefer_group
19+
from sqlalchemy.orm import defer
2020
from sqlmodel import select, desc
2121
from starlette.datastructures import UploadFile as _UploadFile
2222

@@ -42,7 +42,8 @@ class JobOutput(BaseModel):
4242
output: Any
4343

4444

45-
class JobInfoResponse(BaseModel):
45+
class JobSummaryResponse(BaseModel):
46+
"""Job shape without stdout/stderr -- used for list views where log text is never loaded."""
4647
id: str
4748
script: str
4849
parameters: Dict[str, Any]
@@ -51,27 +52,32 @@ class JobInfoResponse(BaseModel):
5152
created_time: Any
5253
start_time: Any = None
5354
end_time: Any = None
54-
stdout: str | None = None
55-
stderr: str | None = None
5655
project_id: str
5756
script_label: str | None = None
5857
scenario_name: str | None = None
5958
duration: float | None = None
6059

6160
@classmethod
62-
def from_job_info(
63-
cls,
64-
job: JobInfo,
65-
stdout: str | None = None,
66-
stderr: str | None = None,
67-
) -> "JobInfoResponse":
61+
def from_job_info(cls, job: JobInfo) -> "JobSummaryResponse":
6862
payload = job.model_dump(exclude={"stdout", "stderr"})
69-
return cls(**payload, stdout=stdout, stderr=stderr)
63+
return cls(**payload)
7064

7165
def to_event_payload(self) -> Dict[str, Any]:
7266
return self.model_dump(mode='json')
7367

7468

69+
class JobInfoResponse(JobSummaryResponse):
70+
"""Full job shape including stdout/stderr -- only GET /jobs/{job_id} returns this; every
71+
other route/event returns JobSummaryResponse, since nothing else reads log text from them."""
72+
stdout: str | None = None
73+
stderr: str | None = None
74+
75+
@classmethod
76+
def from_job_info(cls, job: JobInfo) -> "JobInfoResponse":
77+
payload = job.model_dump(exclude={"stdout", "stderr"})
78+
return cls(**payload, stdout=job.stdout, stderr=job.stderr)
79+
80+
7581
def get_cea_job_temp_prefix(job_id: str) -> str:
7682
"""Get the prefix used for temporary directories for a given job ID."""
7783
return f"cea_job_{job_id}_"
@@ -192,7 +198,10 @@ async def locked_owned_job(session: SessionDep, job_id: str, user_id: CEAUserID)
192198
normalized_id = _normalize_job_id(job_id)
193199
try:
194200
result = await session.execute(
195-
select(JobInfo).where(JobInfo.id == normalized_id).with_for_update()
201+
select(JobInfo)
202+
.options(defer(JobInfo.stdout), defer(JobInfo.stderr)) # type: ignore[arg-type]
203+
.where(JobInfo.id == normalized_id)
204+
.with_for_update()
196205
)
197206
except sqlalchemy.exc.OperationalError as e:
198207
logger.error(e)
@@ -214,14 +223,21 @@ async def get_jobs(
214223
offset: int = Query(0, ge=0, description="Number of jobs to skip"),
215224
state: int | None = Query(None, description="Filter by job state (0=PENDING, 1=STARTED, 2=SUCCESS, 3=ERROR, 4=CANCELED, 5=KILLED)"),
216225
exclude_deleted: bool = Query(True, description="Exclude deleted jobs from results")
217-
) -> List[JobInfo]:
226+
) -> List[JobSummaryResponse]:
218227
"""
219228
Get a paginated list of jobs for the current user and project with optional filtering.
220229
221230
Returns jobs ordered by creation time (most recent first), paginated by `limit` and `offset`.
222231
Jobs are filtered by deleted_at field rather than state to preserve completion states.
232+
233+
stdout/stderr are excluded from list results (can be large); fetch a single job via
234+
GET /jobs/{job_id} for full log output.
223235
"""
224-
query = select(JobInfo).where(JobInfo.project_id == project_id, JobInfo.created_by == user_id)
236+
query = (
237+
select(JobInfo)
238+
.options(defer(JobInfo.stdout), defer(JobInfo.stderr)) # type: ignore[arg-type]
239+
.where(JobInfo.project_id == project_id, JobInfo.created_by == user_id)
240+
)
225241

226242
# Filter by state if specified
227243
if state is not None:
@@ -233,7 +249,7 @@ async def get_jobs(
233249

234250
# Exclude deleted jobs by default based on deleted_at field
235251
if exclude_deleted:
236-
query = query.where(JobInfo.deleted_at.is_(None))
252+
query = query.where(JobInfo.deleted_at.is_(None)) # type: ignore[arg-type]
237253

238254
# Order by created_time descending (most recent first)
239255
query = query.order_by(desc(JobInfo.created_time))
@@ -242,22 +258,24 @@ async def get_jobs(
242258
query = query.limit(limit).offset(offset)
243259

244260
result = await session.execute(query)
245-
return list(result.scalars().all())
261+
jobs_list = result.scalars().all()
262+
263+
return [JobSummaryResponse.from_job_info(job) for job in jobs_list]
246264

247265

248266
@router.get("/{job_id}")
249267
async def get_job_info(session: SessionDep, job_id: str, user_id: CEAUserID) -> JobInfoResponse:
250268
"""Return a JobInfo by id"""
251269
normalized_id = _normalize_job_id(job_id)
252-
job = await session.get(JobInfo, normalized_id, options=[undefer_group('logs')])
270+
job = await session.get(JobInfo, normalized_id)
253271
job = _authorize_job_access(job, normalized_id, user_id)
254272

255-
return JobInfoResponse.from_job_info(job, stdout=job.stdout, stderr=job.stderr)
273+
return JobInfoResponse.from_job_info(job)
256274

257275

258276
@router.post("/new")
259277
async def create_new_job(request: Request, session: SessionDep, project_id: CEAProjectID, user_id: CEAUserID,
260-
settings: CEAServerSettings) -> JobInfoResponse:
278+
settings: CEAServerSettings) -> JobSummaryResponse:
261279
"""Post a new job to the list of jobs to complete"""
262280
content_type = request.headers.get("content-type", "")
263281

@@ -318,14 +336,14 @@ async def create_new_job(request: Request, session: SessionDep, project_id: CEAP
318336
await session.commit()
319337
await session.refresh(job)
320338

321-
job_payload = JobInfoResponse.from_job_info(job)
339+
job_payload = JobSummaryResponse.from_job_info(job)
322340
event_payload = job_payload.to_event_payload()
323341
await emit_with_retry("cea-job-created", event_payload, room=f"user-{job.created_by}")
324342
return job_payload
325343

326344

327345
@router.post("/started/{job_id}")
328-
async def set_job_started(session: SessionDep, job: LockedOwnedJob) -> JobInfoResponse:
346+
async def set_job_started(session: SessionDep, job: LockedOwnedJob) -> JobSummaryResponse:
329347
try:
330348
job.state = JobState.STARTED
331349
job.start_time = get_current_time()
@@ -337,15 +355,15 @@ async def set_job_started(session: SessionDep, job: LockedOwnedJob) -> JobInfoRe
337355
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(e))
338356

339357
# Emit event outside try-except so emit failures don't cause rollback
340-
job_payload = JobInfoResponse.from_job_info(job)
358+
job_payload = JobSummaryResponse.from_job_info(job)
341359
event_payload = job_payload.to_event_payload()
342360
await emit_with_retry("cea-worker-started", event_payload, room=f"user-{job.created_by}")
343361
return job_payload
344362

345363

346364
@router.post("/success/{job_id}")
347365
async def set_job_success(session: SessionDep, job: LockedOwnedJob,
348-
streams: CEAStreams, worker_processes: CEAWorkerProcesses, output: JobOutput) -> JobInfoResponse:
366+
streams: CEAStreams, worker_processes: CEAWorkerProcesses, output: JobOutput) -> JobSummaryResponse:
349367
try:
350368
job.state = JobState.SUCCESS
351369
job.error = None
@@ -368,7 +386,7 @@ async def set_job_success(session: SessionDep, job: LockedOwnedJob,
368386
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(e))
369387

370388
# Emit event outside try-except so emit failures don't cause rollback
371-
job_payload = JobInfoResponse.from_job_info(job, stdout=stdout_text)
389+
job_payload = JobSummaryResponse.from_job_info(job)
372390
event_payload = job_payload.to_event_payload()
373391
event_payload["output"] = output.output
374392
await emit_with_retry("cea-worker-success", event_payload, room=f"user-{job.created_by}")
@@ -377,7 +395,7 @@ async def set_job_success(session: SessionDep, job: LockedOwnedJob,
377395

378396
@router.post("/error/{job_id}")
379397
async def set_job_error(session: SessionDep, job: LockedOwnedJob,
380-
error: JobError, streams: CEAStreams, worker_processes: CEAWorkerProcesses) -> JobInfoResponse:
398+
error: JobError, streams: CEAStreams, worker_processes: CEAWorkerProcesses) -> JobSummaryResponse:
381399
message = error.message
382400
stacktrace = error.stacktrace
383401

@@ -404,7 +422,7 @@ async def set_job_error(session: SessionDep, job: LockedOwnedJob,
404422
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(e))
405423

406424
# Emit event outside try-except so emit failures don't cause rollback
407-
job_payload = JobInfoResponse.from_job_info(job, stdout=stdout_text, stderr=stacktrace)
425+
job_payload = JobSummaryResponse.from_job_info(job)
408426
event_payload = job_payload.to_event_payload()
409427
await emit_with_retry("cea-worker-error", event_payload, room=f"user-{job.created_by}")
410428

@@ -457,7 +475,7 @@ async def start_job(worker_processes: CEAWorkerProcesses, server_url: CEAServerU
457475

458476
@router.post("/cancel/{job_id}")
459477
async def cancel_job(session: SessionDep, job: LockedOwnedJob,
460-
worker_processes: CEAWorkerProcesses, streams: CEAStreams) -> JobInfoResponse:
478+
worker_processes: CEAWorkerProcesses, streams: CEAStreams) -> JobSummaryResponse:
461479
# Validate state: can only cancel PENDING or STARTED jobs (protected by row lock)
462480
if job.state not in (JobState.PENDING, JobState.STARTED):
463481
raise HTTPException(
@@ -496,13 +514,13 @@ async def cancel_job(session: SessionDep, job: LockedOwnedJob,
496514
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(e))
497515

498516
# Emit event outside try-except so emit failures don't cause rollback
499-
job_payload = JobInfoResponse.from_job_info(job, stdout=stdout_text)
517+
job_payload = JobSummaryResponse.from_job_info(job)
500518
event_payload = job_payload.to_event_payload()
501519
await emit_with_retry("cea-worker-canceled", event_payload, room=f"user-{job.created_by}")
502520
return job_payload
503521

504522

505-
async def kill_job(session, job_id: str, worker_processes, streams) -> JobInfoResponse:
523+
async def kill_job(session, job_id: str, worker_processes, streams) -> JobSummaryResponse:
506524
"""
507525
Kill a job (server-initiated termination, e.g., during shutdown).
508526
This is different from cancel_job which is user-initiated.
@@ -544,14 +562,14 @@ async def kill_job(session, job_id: str, worker_processes, streams) -> JobInfoRe
544562
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(e))
545563

546564
# Emit event outside try-except so emit failures don't cause rollback
547-
job_payload = JobInfoResponse.from_job_info(job, stdout=stdout_text)
565+
job_payload = JobSummaryResponse.from_job_info(job)
548566
event_payload = job_payload.to_event_payload()
549567
await emit_with_retry("cea-worker-killed", event_payload, room=f"user-{job.created_by}")
550568
return job_payload
551569

552570

553571
@router.delete("/{job_id}")
554-
async def delete_job(session: SessionDep, job: LockedOwnedJob, user_id: CEAUserID) -> JobInfoResponse:
572+
async def delete_job(session: SessionDep, job: LockedOwnedJob, user_id: CEAUserID) -> JobSummaryResponse:
555573
"""
556574
Mark a job as deleted (soft delete). The job row is not removed from the database,
557575
and the original completion state (SUCCESS/ERROR/CANCELED/KILLED) is preserved.
@@ -581,7 +599,7 @@ async def delete_job(session: SessionDep, job: LockedOwnedJob, user_id: CEAUserI
581599
await session.rollback()
582600
raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(e))
583601

584-
job_payload = JobInfoResponse.from_job_info(job)
602+
job_payload = JobSummaryResponse.from_job_info(job)
585603
event_payload = job_payload.to_event_payload()
586604
await emit_with_retry("cea-job-deleted", event_payload, room=f"user-{job.created_by}")
587605
return job_payload

cea/interfaces/dashboard/server/streams.py

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,6 @@
44
FIXME: when does this data get cleared?
55
"""
66
from fastapi import APIRouter, HTTPException, Request, status
7-
from sqlalchemy.orm import undefer_group
87

98
from cea.interfaces.dashboard.dependencies import CEAStreams, CEAUserID
109
from cea.interfaces.dashboard.lib.database.models import JobInfo
@@ -19,7 +18,7 @@
1918

2019
@router.get("/read/{job_id}")
2120
async def read_stream(session: SessionDep, streams: CEAStreams, job_id: str, user_id: CEAUserID):
22-
job = await session.get(JobInfo, job_id, options=[undefer_group('logs')])
21+
job = await session.get(JobInfo, job_id)
2322
if job is None:
2423
logger.info(f"read_stream: job {job_id} not found")
2524
return "" # Return empty string for non-existent jobs

cea/tests/test_dashboard_job_info.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -82,8 +82,10 @@ async def test_set_job_error_serialises_deferred_logs_without_lazy_loading(monke
8282

8383
assert response.state == JobState.ERROR
8484
assert response.error == "mesh failed"
85-
assert response.stdout == "stdout line"
86-
assert response.stderr == "traceback"
85+
# set_job_error returns JobSummaryResponse -- stdout/stderr are never in the
86+
# response/event payload, only persisted to the DB (checked below).
87+
assert not hasattr(response, "stdout")
88+
assert not hasattr(response, "stderr")
8789

8890
async with db_session() as session:
8991
stored_job = await session.get(JobInfo, job_id, options=[undefer_group("logs")])

0 commit comments

Comments
 (0)