mirror of
https://github.com/tiennm99/DocsGPT.git
synced 2026-10-04 22:13:08 +00:00
The "Tool approval needed" toast (and several sibling surfaces) could linger after the state they represent was already gone. User-scoped SSE events (tool.approval.required, schedule.autopaused, attachment.queued, …) are durable and replayed on reconnect, but no terminal path emitted a matching clearing event and the reconciler only wrote operator-facing stack_logs — so a failed/expired message replayed its approval prompt with nothing to act on, and the toast trusted event presence over the actual message state. Backend — emit a user-facing event on every terminal path: - reconciler deletes pending_tool_state and publishes tool.approval.cleared when a stuck message is failed; cleanup_pending_tool_state does the same for TTL-reaped rows - reconciler now publishes source.ingest.failed (stalled ingest), schedule.run.failed (timeout/pending) and schedule.completed (once) - schedules PATCH-resume / DELETE publish schedule.resumed / .cancelled so a stale schedule.autopaused can't outvote them on replay - store_attachment gains the on_poison terminal-event hook the ingest tasks already have; mcp_oauth_task gains a soft/hard time limit so a hung flow self-reports mcp.oauth.failed - tighten the reconciler exemption to the (conversation_id, user_id) composite key Frontend — stop trusting event presence over truth: - notificationsSlice.resolveToolApproval evicts the matching tool.approval.required and persists its (stable) id dismissed so the backlog replay stays suppressed - ToolApprovalToast drops approval events older than the resumable TTL window as a backstop for a lost clearing event - schedulesSlice handles schedule.resumed / .cancelled / .completed - upload dismissals now outlive the SSE backlog retention window Tests cover the new clearing events at the reconciler, slice and dispatch layers.
500 lines
16 KiB
Python
500 lines
16 KiB
Python
from datetime import timedelta
|
|
|
|
from application.api.user.idempotency import with_idempotency
|
|
from application.celery_init import celery
|
|
from application.worker import (
|
|
agent_webhook_worker,
|
|
attachment_worker,
|
|
ingest_worker,
|
|
mcp_oauth,
|
|
remote_worker,
|
|
sync,
|
|
sync_worker,
|
|
)
|
|
|
|
|
|
# Shared decorator config for long-running, side-effecting tasks. ``acks_late``
|
|
# is also the celeryconfig default but stays explicit here so each task's
|
|
# durability story is grep-able next to the body. Combined with
|
|
# ``autoretry_for=(Exception,)`` and a bounded ``max_retries`` so a poison
|
|
# message can't loop forever.
|
|
DURABLE_TASK = dict(
|
|
bind=True,
|
|
acks_late=True,
|
|
autoretry_for=(Exception,),
|
|
retry_kwargs={"max_retries": 3, "countdown": 60},
|
|
retry_backoff=True,
|
|
)
|
|
|
|
|
|
# operation tag for the poison-path source.ingest.failed event, per task.
|
|
_INGEST_POISON_OPERATION = {
|
|
"ingest": "upload",
|
|
"ingest_remote": "upload",
|
|
"ingest_connector_task": "upload",
|
|
"reingest_source_task": "reingest",
|
|
}
|
|
|
|
|
|
def _emit_ingest_poison_event(task_name, bound):
|
|
"""Publish a terminal ``source.ingest.failed`` when the poison-guard trips.
|
|
|
|
The guard returns before the worker runs, so the worker's own failed
|
|
event never fires — without this the upload toast spins on "training".
|
|
"""
|
|
user = bound.get("user")
|
|
source_id = bound.get("source_id")
|
|
if not user or not source_id:
|
|
return
|
|
from application.events.publisher import publish_user_event
|
|
|
|
publish_user_event(
|
|
user,
|
|
"source.ingest.failed",
|
|
{
|
|
"source_id": str(source_id),
|
|
"filename": bound.get("filename") or "",
|
|
"operation": _INGEST_POISON_OPERATION.get(task_name, "upload"),
|
|
"error": "Ingestion stopped after repeated failures.",
|
|
},
|
|
scope={"kind": "source", "id": str(source_id)},
|
|
)
|
|
|
|
|
|
@celery.task(**DURABLE_TASK)
|
|
@with_idempotency(task_name="ingest", on_poison=_emit_ingest_poison_event)
|
|
def ingest(
|
|
self,
|
|
directory,
|
|
formats,
|
|
job_name,
|
|
user,
|
|
file_path,
|
|
filename,
|
|
file_name_map=None,
|
|
idempotency_key=None,
|
|
source_id=None,
|
|
):
|
|
resp = ingest_worker(
|
|
self,
|
|
directory,
|
|
formats,
|
|
job_name,
|
|
file_path,
|
|
filename,
|
|
user,
|
|
file_name_map=file_name_map,
|
|
idempotency_key=idempotency_key,
|
|
source_id=source_id,
|
|
)
|
|
return resp
|
|
|
|
|
|
@celery.task(**DURABLE_TASK)
|
|
@with_idempotency(task_name="ingest_remote", on_poison=_emit_ingest_poison_event)
|
|
def ingest_remote(
|
|
self, source_data, job_name, user, loader,
|
|
idempotency_key=None, source_id=None,
|
|
):
|
|
resp = remote_worker(
|
|
self, source_data, job_name, user, loader,
|
|
idempotency_key=idempotency_key,
|
|
source_id=source_id,
|
|
)
|
|
return resp
|
|
|
|
|
|
@celery.task(**DURABLE_TASK)
|
|
@with_idempotency(
|
|
task_name="reingest_source_task", on_poison=_emit_ingest_poison_event,
|
|
)
|
|
def reingest_source_task(self, source_id, user, idempotency_key=None):
|
|
from application.worker import reingest_source_worker
|
|
|
|
resp = reingest_source_worker(self, source_id, user)
|
|
return resp
|
|
|
|
|
|
# Beat-driven dispatch tasks default to ``acks_late=False``: a SIGKILL
|
|
# of a beat tick is harmless to redeliver only if the dispatch itself is
|
|
# idempotent. We keep these early-ACK so the broker doesn't replay a
|
|
# dispatch that already enqueued downstream work.
|
|
@celery.task(bind=True, acks_late=False)
|
|
def schedule_syncs(self, frequency):
|
|
resp = sync_worker(self, frequency)
|
|
return resp
|
|
|
|
|
|
@celery.task(bind=True)
|
|
def sync_source(
|
|
self,
|
|
source_data,
|
|
job_name,
|
|
user,
|
|
loader,
|
|
sync_frequency,
|
|
retriever,
|
|
doc_id,
|
|
):
|
|
resp = sync(
|
|
self,
|
|
source_data,
|
|
job_name,
|
|
user,
|
|
loader,
|
|
sync_frequency,
|
|
retriever,
|
|
doc_id,
|
|
)
|
|
return resp
|
|
|
|
|
|
def _emit_attachment_poison_event(task_name, bound):
|
|
"""Publish a terminal ``attachment.failed`` when the poison-guard trips.
|
|
|
|
Mirrors ``_emit_ingest_poison_event``: the guard returns before the
|
|
worker runs, so ``attachment_worker``'s own events never fire and the
|
|
upload toast would otherwise spin on "processing" forever.
|
|
"""
|
|
user = bound.get("user")
|
|
file_info = bound.get("file_info") or {}
|
|
attachment_id = file_info.get("attachment_id")
|
|
if not user or not attachment_id:
|
|
return
|
|
from application.events.publisher import publish_user_event
|
|
|
|
publish_user_event(
|
|
user,
|
|
"attachment.failed",
|
|
{
|
|
"attachment_id": str(attachment_id),
|
|
"filename": file_info.get("filename") or "",
|
|
"error": "Attachment processing stopped after repeated failures.",
|
|
},
|
|
scope={"kind": "attachment", "id": str(attachment_id)},
|
|
)
|
|
|
|
|
|
@celery.task(**DURABLE_TASK)
|
|
@with_idempotency(
|
|
task_name="store_attachment", on_poison=_emit_attachment_poison_event,
|
|
)
|
|
def store_attachment(self, file_info, user, idempotency_key=None):
|
|
resp = attachment_worker(self, file_info, user)
|
|
return resp
|
|
|
|
|
|
@celery.task(**DURABLE_TASK)
|
|
@with_idempotency(task_name="process_agent_webhook")
|
|
def process_agent_webhook(self, agent_id, payload, idempotency_key=None):
|
|
resp = agent_webhook_worker(self, agent_id, payload)
|
|
return resp
|
|
|
|
|
|
@celery.task(**DURABLE_TASK)
|
|
@with_idempotency(
|
|
task_name="ingest_connector_task", on_poison=_emit_ingest_poison_event,
|
|
)
|
|
def ingest_connector_task(
|
|
self,
|
|
job_name,
|
|
user,
|
|
source_type,
|
|
session_token=None,
|
|
file_ids=None,
|
|
folder_ids=None,
|
|
recursive=True,
|
|
retriever="classic",
|
|
operation_mode="upload",
|
|
doc_id=None,
|
|
sync_frequency="never",
|
|
idempotency_key=None,
|
|
source_id=None,
|
|
):
|
|
from application.worker import ingest_connector
|
|
|
|
resp = ingest_connector(
|
|
self,
|
|
job_name,
|
|
user,
|
|
source_type,
|
|
session_token=session_token,
|
|
file_ids=file_ids,
|
|
folder_ids=folder_ids,
|
|
recursive=recursive,
|
|
retriever=retriever,
|
|
operation_mode=operation_mode,
|
|
doc_id=doc_id,
|
|
sync_frequency=sync_frequency,
|
|
idempotency_key=idempotency_key,
|
|
source_id=source_id,
|
|
)
|
|
return resp
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def dispatch_scheduled_runs(self):
|
|
"""Beat-driven scheduler poller (body in scheduler_dispatcher)."""
|
|
from application.api.user.scheduler_dispatcher import dispatch_due_runs
|
|
|
|
return dispatch_due_runs()
|
|
|
|
|
|
@celery.task(
|
|
bind=True,
|
|
acks_late=True,
|
|
# Not DURABLE_TASK: agent runs have side effects; blind retry would double them.
|
|
autoretry_for=(),
|
|
max_retries=0,
|
|
)
|
|
def execute_scheduled_run(self, run_id):
|
|
"""Execute one scheduled run; soft-time-limit honors SCHEDULE_RUN_TIMEOUT."""
|
|
from application.api.user.scheduler_worker import execute_scheduled_run_body
|
|
|
|
return execute_scheduled_run_body(run_id, getattr(self.request, "id", None))
|
|
|
|
|
|
# Bind runtime soft-time-limit so the prefork worker can raise mid-agent.
|
|
try:
|
|
from application.core.settings import settings as _scheduler_settings
|
|
execute_scheduled_run.soft_time_limit = max(
|
|
30, int(_scheduler_settings.SCHEDULE_RUN_TIMEOUT),
|
|
)
|
|
execute_scheduled_run.time_limit = (
|
|
execute_scheduled_run.soft_time_limit + 60
|
|
)
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def cleanup_schedule_runs(self):
|
|
"""Trim ``schedule_runs`` per ``SCHEDULE_RUN_OUTPUT_RETENTION_DAYS``."""
|
|
from application.core.settings import settings
|
|
if not settings.POSTGRES_URI:
|
|
return {"deleted": 0, "skipped": "POSTGRES_URI not set"}
|
|
|
|
from application.storage.db.engine import get_engine
|
|
from application.storage.db.repositories.schedule_runs import (
|
|
ScheduleRunsRepository,
|
|
)
|
|
|
|
ttl_days = settings.SCHEDULE_RUN_OUTPUT_RETENTION_DAYS
|
|
engine = get_engine()
|
|
with engine.begin() as conn:
|
|
deleted = ScheduleRunsRepository(conn).cleanup_older_than(ttl_days)
|
|
return {"deleted": deleted, "ttl_days": ttl_days}
|
|
|
|
|
|
@celery.on_after_configure.connect
|
|
def setup_periodic_tasks(sender, **kwargs):
|
|
from application.core.settings import settings
|
|
|
|
sender.add_periodic_task(
|
|
timedelta(days=1),
|
|
schedule_syncs.s("daily"),
|
|
)
|
|
sender.add_periodic_task(
|
|
timedelta(weeks=1),
|
|
schedule_syncs.s("weekly"),
|
|
)
|
|
sender.add_periodic_task(
|
|
timedelta(days=30),
|
|
schedule_syncs.s("monthly"),
|
|
)
|
|
# Replaces Mongo's TTL index on pending_tool_state.expires_at.
|
|
sender.add_periodic_task(
|
|
timedelta(seconds=60),
|
|
cleanup_pending_tool_state.s(),
|
|
name="cleanup-pending-tool-state",
|
|
)
|
|
# Pure housekeeping for ``task_dedup`` / ``webhook_dedup`` — the
|
|
# upsert paths already handle stale rows, so cadence only bounds
|
|
# table size. Hourly is plenty for typical traffic.
|
|
sender.add_periodic_task(
|
|
timedelta(hours=1),
|
|
cleanup_idempotency_dedup.s(),
|
|
name="cleanup-idempotency-dedup",
|
|
)
|
|
sender.add_periodic_task(
|
|
timedelta(seconds=30),
|
|
reconciliation_task.s(),
|
|
name="reconciliation",
|
|
)
|
|
sender.add_periodic_task(
|
|
timedelta(hours=7),
|
|
version_check_task.s(),
|
|
name="version-check",
|
|
)
|
|
# Bound ``message_events`` growth — every streamed SSE chunk writes
|
|
# one row, so retained chats accumulate hundreds of rows per
|
|
# message. Reconnect-replay is only meaningful for streams the user
|
|
# could plausibly still be waiting on, so 14 days is generous.
|
|
sender.add_periodic_task(
|
|
timedelta(hours=24),
|
|
cleanup_message_events.s(),
|
|
name="cleanup-message-events",
|
|
)
|
|
sender.add_periodic_task(
|
|
timedelta(hours=24),
|
|
cleanup_orphan_memories.s(),
|
|
name="cleanup-orphan-memories",
|
|
)
|
|
# Scheduler dispatcher and run-log trim.
|
|
sender.add_periodic_task(
|
|
timedelta(seconds=max(15, settings.SCHEDULE_DISPATCHER_INTERVAL)),
|
|
dispatch_scheduled_runs.s(),
|
|
name="dispatch-scheduled-runs",
|
|
)
|
|
sender.add_periodic_task(
|
|
timedelta(hours=24),
|
|
cleanup_schedule_runs.s(),
|
|
name="cleanup-schedule-runs",
|
|
)
|
|
|
|
|
|
# Bound time limits so a hung OAuth discovery (user never finishes the
|
|
# consent flow, upstream never redirects) self-terminates instead of
|
|
# stranding the ``mcp.oauth.awaiting_redirect`` envelope forever. The
|
|
# soft limit raises inside ``mcp_oauth``'s ``try`` so it publishes a
|
|
# terminal ``mcp.oauth.failed``; the hard limit is the prefork backstop.
|
|
# Generous so a human actively clicking through OAuth isn't cut off.
|
|
@celery.task(bind=True, soft_time_limit=600, time_limit=660)
|
|
def mcp_oauth_task(self, config, user):
|
|
resp = mcp_oauth(self, config, user)
|
|
return resp
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def cleanup_pending_tool_state(self):
|
|
"""Revert stale ``resuming`` rows, then delete TTL-expired rows."""
|
|
from application.core.settings import settings
|
|
if not settings.POSTGRES_URI:
|
|
return {"deleted": 0, "reverted": 0, "skipped": "POSTGRES_URI not set"}
|
|
|
|
from application.storage.db.engine import get_engine
|
|
from application.storage.db.repositories.pending_tool_state import (
|
|
PendingToolStateRepository,
|
|
)
|
|
|
|
engine = get_engine()
|
|
with engine.begin() as conn:
|
|
repo = PendingToolStateRepository(conn)
|
|
reverted = repo.revert_stale_resuming(grace_seconds=600)
|
|
cleared = repo.cleanup_expired()
|
|
|
|
# Reaping the resumable state retires any awaiting-approval prompt
|
|
# tied to it. Without a clearing event the durable
|
|
# ``tool.approval.required`` envelope replays on reconnect and the UI
|
|
# toast lingers for a conversation that can no longer be resumed.
|
|
from application.events.publisher import publish_user_event
|
|
|
|
for row in cleared:
|
|
user_id = row.get("user_id")
|
|
conversation_id = row.get("conversation_id")
|
|
if not user_id or not conversation_id:
|
|
continue
|
|
publish_user_event(
|
|
str(user_id),
|
|
"tool.approval.cleared",
|
|
{"conversation_id": str(conversation_id), "reason": "expired"},
|
|
scope={"kind": "conversation", "id": str(conversation_id)},
|
|
)
|
|
return {"deleted": len(cleared), "reverted": reverted}
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def cleanup_idempotency_dedup(self):
|
|
"""Delete TTL-expired rows from ``task_dedup`` and ``webhook_dedup``.
|
|
|
|
Pure housekeeping — the upsert paths already ignore stale rows
|
|
(TTL-aware ``ON CONFLICT DO UPDATE``), so this only bounds table
|
|
growth and keeps SELECT planning tight on large deployments.
|
|
"""
|
|
from application.core.settings import settings
|
|
if not settings.POSTGRES_URI:
|
|
return {
|
|
"task_dedup_deleted": 0,
|
|
"webhook_dedup_deleted": 0,
|
|
"skipped": "POSTGRES_URI not set",
|
|
}
|
|
|
|
from application.storage.db.engine import get_engine
|
|
from application.storage.db.repositories.idempotency import (
|
|
IdempotencyRepository,
|
|
)
|
|
|
|
engine = get_engine()
|
|
with engine.begin() as conn:
|
|
return IdempotencyRepository(conn).cleanup_expired()
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def reconciliation_task(self):
|
|
"""Sweep stuck durability rows and escalate them to terminal status + alert."""
|
|
from application.api.user.reconciliation import run_reconciliation
|
|
|
|
return run_reconciliation()
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def cleanup_message_events(self):
|
|
"""Delete ``message_events`` rows older than the retention window.
|
|
|
|
Streamed answer responses write one journal row per SSE yield,
|
|
so unbounded growth would dominate Postgres for any retained-
|
|
conversations deployment. The reconnect-replay path only needs
|
|
rows for in-flight streams; 14 days covers paused/tool-action
|
|
flows comfortably.
|
|
"""
|
|
from application.core.settings import settings
|
|
if not settings.POSTGRES_URI:
|
|
return {"deleted": 0, "skipped": "POSTGRES_URI not set"}
|
|
|
|
from application.storage.db.engine import get_engine
|
|
from application.storage.db.repositories.message_events import (
|
|
MessageEventsRepository,
|
|
)
|
|
|
|
ttl_days = settings.MESSAGE_EVENTS_RETENTION_DAYS
|
|
engine = get_engine()
|
|
with engine.begin() as conn:
|
|
deleted = MessageEventsRepository(conn).cleanup_older_than(ttl_days)
|
|
return {"deleted": deleted, "ttl_days": ttl_days}
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def cleanup_orphan_memories(self):
|
|
"""Sweep orphan memories left by the 0009 FK-to-trigger orphan window.
|
|
|
|
A ``memories`` INSERT for a real ``tool_id`` racing a ``user_tools``
|
|
DELETE leaves a permanent orphan the dropped FK would have rejected.
|
|
Default-tool synthetic ids are preserved (legitimate built-in data).
|
|
"""
|
|
from application.core.settings import settings
|
|
if not settings.POSTGRES_URI:
|
|
return {"deleted": 0, "skipped": "POSTGRES_URI not set"}
|
|
|
|
from application.agents.default_tools import default_tool_ids
|
|
from application.storage.db.engine import get_engine
|
|
from application.storage.db.repositories.memories import MemoriesRepository
|
|
|
|
keep_tool_ids = list(default_tool_ids().values())
|
|
engine = get_engine()
|
|
with engine.begin() as conn:
|
|
deleted = MemoriesRepository(conn).delete_orphans(keep_tool_ids)
|
|
return {"deleted": deleted}
|
|
|
|
|
|
@celery.task(bind=True, acks_late=False)
|
|
def version_check_task(self):
|
|
"""Periodic anonymous version check.
|
|
|
|
Complements the ``worker_ready`` boot trigger so long-running
|
|
deployments (>6h cache TTL) still refresh advisories. ``run_check``
|
|
is fail-silent and coordinates across replicas via Redis lock +
|
|
cache (see ``application.updates.version_check``).
|
|
"""
|
|
from application.updates.version_check import run_check
|
|
run_check()
|