Recipe — Unified scheduling, the priority queue, and the ScholarX RSS research feed¶
The gateway daemon runs one intelligent scheduler (CONCEPT:AU-OS.state.unified-scheduling-one-intelligent). Every
recurring job — the deploy/schedules.yml entries, the former fixed-interval
maintenance ticks, the self-evolution loop_cycle, and the ScholarX RSS research
feed — is a durable :Schedule node. The single scheduler tick evaluates them and
enqueues a scheduled_job :Task onto one hardened priority+scheduled queue
(CONCEPT:AU-KG.ingest.hardened-priority-scheduled-task) that the worker pool drains. Nothing recurring runs inline in the
scheduler thread anymore.
The queue (CONCEPT:AU-KG.ingest.hardened-priority-scheduled-task)¶
A :Task carries:
prio_bucket— discrete priority0(critical) …3(background). Workers claim the lowest non-empty bucket first (the L1 graph interpreter stripsORDER BY, so priority is N equality queries, not a sort).prioritize_taskand theprio/priorityarguments set it.scheduled+ eta — delayed execution. A task with a future eta waits asstatus='scheduled'until the per-minutepromotion_sweepmakes itpending. Application-level retry/backoff reuses this lane (a failure reschedules with exponential backoff); a task that exhaustsmax_attemptsbecomes adead_letter(distinct from the reaper's crash-requeue).blocked+depends_on— dependency gating. A blocked task is promoted once every dependency hascompleted; a terminally-failed dependency cancels it.
Controlling schedules (two surfaces)¶
# MCP
graph_schedules action=list
graph_schedules action=disable name=research_feed
graph_schedules action=prioritize name=loop_cycle priority=1
graph_schedules action=set_interval name=research_feed interval_s=900
graph_schedules action=run_now name=enrichment
# REST (the auto-mounted twin)
curl -s localhost:8080/graph/schedules -d '{"action":"list"}'
deploy/schedules.yml is the seed (desired state); the :Schedule node holds
live last-run / next-run / failure-backoff and survives restart and leader-failover.
The ScholarX RSS research feed (CONCEPT:AU-KG.research.scholarx-rss-research-feed)¶
A default-on research_feed schedule (KG_RESEARCH_FEED, cadence
KG_RESEARCH_FEED_INTERVAL, default 30 min) enqueues
LoopController.run_rss_feed_screen, which:
- reads the arXiv RSS feed (
get_recent_papers(days=1)— cheap title+abstract); - skips already-examined items via a
DeltaManifestseen-set keyed by arXiv id (every graded item, including rejects, is recorded so it is never re-graded); - grades each new item — keyword taxonomy (
score_paper) plus a ConceptMatcher novelty probe (_paper_novelty); on a GPU/embedder outage the novelty probe returnsNoneand grading degrades to keyword-only rather than failing; - enqueues a
research_paper_fetchtask for items at/above the relevance threshold, withprio_bucketderived from the grade — so the highest-graded papers are fetched and ingested first (priority = queue reordering). The fetch task downloads the full paper and ingests it viaResearchPipelineRunner.ingest_paper_full. Marginal items get a cheap abstract-only ingest inline.
Enable/disable or retune it like any schedule via graph_schedules.
Duplicate-tick safety (coalesce + collapse)¶
A scheduled job is an interval tick, not a backlog item — running a stale missed tick adds no value. Two mechanisms keep the queue from accumulating duplicates:
- Coalesce (per-schedule, at enqueue): a tick is not enqueued while a prior tick for the same schedule is still un-consumed (CONCEPT:AU-OS.state.unified-scheduling-one-intelligent).
- Collapse (self-healing, at each tick):
collapse_stale_tickscancels any schedule's active duplicate ticks down to ≤1, recovering from a backlog that pre-dates the coalescer or a window where its probe failed (CONCEPT:AU-OS.state.stale-tick-collapse).runningticks are never touched.
This pairs with the best-effort lane cap (CONCEPT:AU-ORCH.scheduling.low-value-high-volume): the maint lane
is capped at its floor coverage so a tick backlog can never crowd the throughput
lanes. See Ingestion Throughput.
Verifying end-to-end¶
# 1. trigger the feed now and watch the queue
graph_schedules action=run_now name=research_feed
graph_query "MATCH (t:Task {type:'research_paper_fetch'}) RETURN t.id, t.prio_bucket, t.status"
# 2. re-run: already-seen items are skipped (seen_skipped > 0)
graph_schedules action=run_now name=research_feed
# 3. the ingested papers land as Documents
graph_query "MATCH (a) WHERE a.id STARTS WITH 'article:scholarx:' RETURN count(a)"
See also: Delta-based ingestion, the gateway daemon map, the Loop engine.