Rthur2003 commited on
Commit
45fa2f8
·
1 Parent(s): db51cdb

feat: background analysis jobs servisi ile queueing progress tracking ve caching eklendi

Browse files
app/services/analysis_jobs.py CHANGED
@@ -63,6 +63,8 @@ class Progress:
63
  self._started: Dict[str, float] = {}
64
  self._took: Dict[str, float] = {}
65
  self.cancelled = False
 
 
66
  # Returns True when nobody is waiting for the result any more.
67
  self.gone: Callable[[], Awaitable[bool]] = _never
68
 
@@ -143,6 +145,8 @@ class CpuGate:
143
 
144
  @asynccontextmanager
145
  async def slot(self, ticket: str, progress: Optional[Progress] = None) -> AsyncIterator[None]:
 
 
146
  event = asyncio.Event()
147
  self._events[ticket] = event
148
  self._waiting.append(ticket)
@@ -159,6 +163,7 @@ class CpuGate:
159
  started = time.monotonic()
160
  try:
161
  if progress is not None:
 
162
  progress.check()
163
  if await progress.gone():
164
  raise JobCancelled("job_abandoned")
@@ -216,22 +221,26 @@ class Job:
216
  created: float = field(default_factory=time.time)
217
  updated: float = field(default_factory=time.time)
218
  last_seen: float = field(default_factory=time.monotonic)
219
- status: str = "queued" # queued | running | done | error
220
  cached: bool = False
221
  response: Optional[Dict[str, Any]] = None
222
  task: Optional[asyncio.Task] = None
223
 
 
 
 
 
224
  def to_dict(self) -> Dict[str, Any]:
225
  data: Dict[str, Any] = {
226
  "jobId": self.id,
227
  "sourceType": self.source_type,
228
- "status": self.status,
229
  "steps": self.progress.snapshot(),
230
- "elapsedSec": round((self.updated if self.status in ("done", "error") else time.time()) - self.created, 2),
231
  "response": self.response,
232
  }
233
  position = cpu_gate.position(self.id)
234
- if self.status == "queued" and position is not None:
235
  data["queue"] = {"position": position, "estimatedWaitSec": cpu_gate.estimated_wait(position)}
236
  if self.cached:
237
  data["cached"] = True
@@ -252,25 +261,25 @@ class JobStore:
252
 
253
  def _prune(self) -> None:
254
  cutoff = time.time() - JOB_TTL_SEC
255
- for job_id in [j.id for j in self._jobs.values() if j.status in ("done", "error") and j.updated < cutoff]:
256
  del self._jobs[job_id]
257
  if len(self._jobs) > MAX_JOBS:
258
- finished = sorted((j for j in self._jobs.values() if j.status in ("done", "error")), key=lambda j: j.updated)
259
  for job in finished[: len(self._jobs) - MAX_JOBS]:
260
  del self._jobs[job.id]
261
 
262
  def pending_count(self) -> int:
263
- return sum(1 for j in self._jobs.values() if j.status in ("queued", "running"))
264
 
265
  def running_count(self) -> int:
266
- return sum(1 for j in self._jobs.values() if j.status == "running")
267
 
268
  def start(self, source_type: str, steps: List[str], work: Work, key: Optional[str] = None) -> Job:
269
  """Queue a job; the same file or link already in flight is shared, a recent one is served from cache."""
270
  self._prune()
271
  if key and key in self._inflight:
272
  shared = self._jobs.get(self._inflight[key])
273
- if shared and shared.status in ("queued", "running"):
274
  shared.last_seen = time.monotonic()
275
  return shared
276
 
@@ -305,6 +314,11 @@ class JobStore:
305
  job.progress.skip_pending()
306
  job.response = {"result": None, "warnings": [], "errors": [stop.code]}
307
  job.status = "error"
 
 
 
 
 
308
  except Exception: # noqa: BLE001 — the job must always settle
309
  job.progress.skip_pending()
310
  job.response = {"result": None, "warnings": [], "errors": ["internal_error"]}
@@ -330,10 +344,10 @@ class JobStore:
330
  job = self._jobs.get(job_id)
331
  if job is None:
332
  return None
333
- if job.status in ("queued", "running"):
334
  job.progress.cancelled = True
335
  # A job still waiting for the CPU never gets there.
336
- if job.status == "queued" and job.task is not None:
337
  job.task.cancel()
338
  return job
339
 
 
63
  self._started: Dict[str, float] = {}
64
  self._took: Dict[str, float] = {}
65
  self.cancelled = False
66
+ # queued: waiting for the CPU; running: downloading or analysing.
67
+ self.phase = "queued"
68
  # Returns True when nobody is waiting for the result any more.
69
  self.gone: Callable[[], Awaitable[bool]] = _never
70
 
 
145
 
146
  @asynccontextmanager
147
  async def slot(self, ticket: str, progress: Optional[Progress] = None) -> AsyncIterator[None]:
148
+ if progress is not None:
149
+ progress.phase = "queued"
150
  event = asyncio.Event()
151
  self._events[ticket] = event
152
  self._waiting.append(ticket)
 
163
  started = time.monotonic()
164
  try:
165
  if progress is not None:
166
+ progress.phase = "running"
167
  progress.check()
168
  if await progress.gone():
169
  raise JobCancelled("job_abandoned")
 
221
  created: float = field(default_factory=time.time)
222
  updated: float = field(default_factory=time.time)
223
  last_seen: float = field(default_factory=time.monotonic)
224
+ status: str = "active" # active | done | error
225
  cached: bool = False
226
  response: Optional[Dict[str, Any]] = None
227
  task: Optional[asyncio.Task] = None
228
 
229
+ @property
230
+ def finished(self) -> bool:
231
+ return self.status in ("done", "error")
232
+
233
  def to_dict(self) -> Dict[str, Any]:
234
  data: Dict[str, Any] = {
235
  "jobId": self.id,
236
  "sourceType": self.source_type,
237
+ "status": self.status if self.finished else self.progress.phase,
238
  "steps": self.progress.snapshot(),
239
+ "elapsedSec": round((self.updated if self.finished else time.time()) - self.created, 2),
240
  "response": self.response,
241
  }
242
  position = cpu_gate.position(self.id)
243
+ if not self.finished and position is not None:
244
  data["queue"] = {"position": position, "estimatedWaitSec": cpu_gate.estimated_wait(position)}
245
  if self.cached:
246
  data["cached"] = True
 
261
 
262
  def _prune(self) -> None:
263
  cutoff = time.time() - JOB_TTL_SEC
264
+ for job_id in [j.id for j in self._jobs.values() if j.finished and j.updated < cutoff]:
265
  del self._jobs[job_id]
266
  if len(self._jobs) > MAX_JOBS:
267
+ finished = sorted((j for j in self._jobs.values() if j.finished), key=lambda j: j.updated)
268
  for job in finished[: len(self._jobs) - MAX_JOBS]:
269
  del self._jobs[job.id]
270
 
271
  def pending_count(self) -> int:
272
+ return sum(1 for j in self._jobs.values() if not j.finished)
273
 
274
  def running_count(self) -> int:
275
+ return sum(1 for j in self._jobs.values() if not j.finished and j.progress.phase == "running")
276
 
277
  def start(self, source_type: str, steps: List[str], work: Work, key: Optional[str] = None) -> Job:
278
  """Queue a job; the same file or link already in flight is shared, a recent one is served from cache."""
279
  self._prune()
280
  if key and key in self._inflight:
281
  shared = self._jobs.get(self._inflight[key])
282
+ if shared and not shared.finished:
283
  shared.last_seen = time.monotonic()
284
  return shared
285
 
 
314
  job.progress.skip_pending()
315
  job.response = {"result": None, "warnings": [], "errors": [stop.code]}
316
  job.status = "error"
317
+ except asyncio.CancelledError:
318
+ # Cancelled while still waiting for the CPU.
319
+ job.progress.skip_pending()
320
+ job.response = {"result": None, "warnings": [], "errors": ["cancelled"]}
321
+ job.status = "error"
322
  except Exception: # noqa: BLE001 — the job must always settle
323
  job.progress.skip_pending()
324
  job.response = {"result": None, "warnings": [], "errors": ["internal_error"]}
 
344
  job = self._jobs.get(job_id)
345
  if job is None:
346
  return None
347
+ if not job.finished:
348
  job.progress.cancelled = True
349
  # A job still waiting for the CPU never gets there.
350
+ if cpu_gate.position(job.id) is not None and job.task is not None:
351
  job.task.cancel()
352
  return job
353
 
app/services/feature_extractor.py CHANGED
@@ -84,6 +84,22 @@ class AudioMeta:
84
  channels: int
85
 
86
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
87
  def extract_features(
88
  source: Union[Path, bytes, io.BytesIO],
89
  *,
 
84
  channels: int
85
 
86
 
87
+ def decode_clip(source: Union[Path, bytes, io.BytesIO]) -> bytes:
88
+ """The clip every layer analyses (first 60 s, mono, 22.05 kHz) as a float WAV.
89
+
90
+ Decoding an MP3 is the slowest part of reading it, and each layer used
91
+ to decode the upload on its own. They now read this instead: loading it
92
+ back at 22.05 kHz returns exactly the samples extract_features would
93
+ have produced from the original. Raises ValueError for short or silent audio.
94
+ """
95
+ import soundfile as sf
96
+
97
+ y, sr = _load_audio(source, _TARGET_SR)
98
+ buf = io.BytesIO()
99
+ sf.write(buf, y, sr, format="WAV", subtype="FLOAT")
100
+ return buf.getvalue()
101
+
102
+
103
  def extract_features(
104
  source: Union[Path, bytes, io.BytesIO],
105
  *,