Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
135 changes: 113 additions & 22 deletions erd_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -271,7 +271,14 @@ def guess_depth_from_spine(spine) -> int:
requested_priority INTEGER NOT NULL,
requested_at INTEGER NOT NULL,
started_at INTEGER,
state TEXT NOT NULL DEFAULT 'queued'
state TEXT NOT NULL DEFAULT 'queued',
-- Set when the request completes: the moment, and the word's branch
-- totals over every one of its requests as of that moment. A word whose
-- requests are all complete owns no live branch, so its totals are final.
completed_at INTEGER,
completed_branch_count INTEGER,
completed_direct_branch_count INTEGER,
completed_direct_done_branch_count INTEGER
);

CREATE INDEX IF NOT EXISTS idx_opener_work_priority_order
Expand Down Expand Up @@ -1709,6 +1716,30 @@ def _rebuildable(table, required):
"resolved_at": "INTEGER",
})
self._add_columns("opener_work", {"started_at": "INTEGER"})
opener_work_columns = {
row["name"] for row in
self._conn.execute("PRAGMA table_info(opener_work)")}
self._add_columns("opener_work", {
"completed_at": "INTEGER",
"completed_branch_count": "INTEGER",
"completed_direct_branch_count": "INTEGER",
"completed_direct_done_branch_count": "INTEGER",
Comment thread
ahernsean marked this conversation as resolved.
})
if "completed_branch_count" not in opener_work_columns:
# A queue from before completion totals were recorded: its
# completed requests still hold their memberships, so the totals
# and the moment of completion are taken from those.
self._conn.execute("""
UPDATE opener_work SET completed_at = (
SELECT MAX(membership.resolved_at)
FROM branch_opener_work AS membership
WHERE membership.opener_work_id = opener_work.opener_work_id)
WHERE state = 'complete'
""")
self._record_word_totals(
row["opener"] for row in self._conn.execute(
"SELECT DISTINCT opener FROM opener_work "
"WHERE state = 'complete'"))
if "claimed_at" in pending_columns:
self._conn.execute("""
UPDATE opener_work AS opener
Expand Down Expand Up @@ -2530,10 +2561,49 @@ def mark_openers_complete(self, openers):
if not openers:
return
placeholders = ",".join("?" for _ in openers)
self._conn.execute(
"UPDATE opener_work AS s SET state = 'complete' WHERE opener IN ("
+ placeholders + ") AND " + self._FINISHED_OPENER_WORK_PREDICATE,
openers)
self._conn.execute("BEGIN IMMEDIATE")
try:
self._conn.execute(
"UPDATE opener_work AS s SET state = 'complete', "
"completed_at = ? WHERE opener IN (" + placeholders + ") AND "
+ self._FINISHED_OPENER_WORK_PREDICATE,
(int(time.time()), *openers))
self._record_word_totals(openers)
self._conn.execute("COMMIT")
except Exception:
self._conn.execute("ROLLBACK")
raise

def _record_word_totals(self, openers):
"""Record each word's branch totals on its completed requests.

Every completed request of the word carries the totals as of its
latest completion, so a reader of a finished word takes them from its
request rows instead of rescanning every membership it ever held. The
caller holds the write transaction.
"""
openers = sorted(set(openers))
for start in range(0, len(openers), 500):
chunk = openers[start:start + 500]
self._conn.execute(f"""
UPDATE opener_work AS s
SET (completed_branch_count, completed_direct_branch_count,
completed_direct_done_branch_count) = (
SELECT COUNT(DISTINCT m.branch_id),
COUNT(DISTINCT CASE WHEN m.parent_branch_id IS NULL
THEN m.branch_id END),
COUNT(DISTINCT CASE
WHEN m.parent_branch_id IS NULL
AND p.status = 'done'
THEN m.branch_id END)
FROM opener_work AS word
JOIN branch_opener_work AS m
ON m.opener_work_id = word.opener_work_id
LEFT JOIN pending_branches AS p ON p.branch_id = m.branch_id
WHERE word.opener = s.opener)
WHERE s.state = 'complete'
AND s.opener IN ({",".join("?" for _ in chunk)})
""", chunk)

def _resolve_branch_memberships(self, branch_id: int = None,
withdraw: bool = False):
Expand Down Expand Up @@ -2578,11 +2648,16 @@ def _ready_opener_work_ids(self):
def _finish_opener_work_ids(self, opener_work_ids):
opener_work_ids = sorted(opener_work_ids)
if opener_work_ids:
placeholders = ",".join("?" for _ in opener_work_ids)
self._conn.execute(
"UPDATE opener_work SET state = 'complete' "
"WHERE opener_work_id IN ("
+ ",".join("?" for _ in opener_work_ids) + ")",
opener_work_ids)
"UPDATE opener_work SET state = 'complete', completed_at = ? "
"WHERE opener_work_id IN (" + placeholders + ")",
(int(time.time()), *opener_work_ids))
self._record_word_totals(
row["opener"] for row in self._conn.execute(
"SELECT opener FROM opener_work "
"WHERE opener_work_id IN (" + placeholders + ")",
opener_work_ids))

def _retire_exact_direct_response_groups(self, branch_id: int) -> list[str]:
"""Retire work below direct response groups whose exact result is done.
Expand Down Expand Up @@ -6326,19 +6401,33 @@ def opener_rows(self):
priority is governing the word's remaining queued work. It falls
back to the priority of the completed requests only once none are
outstanding.

A word whose every request is complete reads the totals recorded when
its last request completed; only words with an unfinished request are
counted from their memberships, so the cost follows the queue's
outstanding work rather than everything it has ever finished.
"""
opener_work_columns = {
row["name"] for row in self._conn.execute(
"PRAGMA table_info(opener_work)")
}
started_at = (
"MIN(CASE WHEN s.state != 'complete' THEN s.started_at END)"
if "started_at" in opener_work_columns else "NULL"
)
return self._conn.execute(f"""
return self._conn.execute("""
SELECT s.opener,
MIN(s.requested_at) AS requested_at,
NULL AS started_at,
MAX(s.requested_priority) AS requested_priority,
MAX(s.completed_at) AS completed_at,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve read-only compatibility with pre-upgrade queues

If the report server opens an existing queue before any writable swarm or CLI process has migrated it, this unconditional reference to completed_at fails with no such column and the Openers report returns a queue-source error. _open_report_queue() deliberately uses initialize_schema=False, so a stopped-swarm/report-only deployment cannot add the columns itself; retain a schema-aware fallback like the previous started_at handling or otherwise ensure the read works against the prior schema.

AGENTS.md reference: AGENTS.md:L518-L521

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not adopting this. The code is always updated before anything reads the queue: the deploy stops the swarm, pulls, and starts it, and erd_search.py start starts the supervisor (whose writable open migrates the queue) before the report server. A report server that ever did read an unmigrated queue would show a queue-source error on the Openers page until the swarm's next start, and the report opens the queue per request, so it recovers on its own with no restart.

COUNT(*) AS request_count,
0 AS has_active_request,
0 AS has_incomplete_request,
MAX(s.completed_branch_count) AS branch_count,
MAX(s.completed_direct_branch_count) AS direct_branch_count,
MAX(s.completed_direct_done_branch_count)
AS direct_done_branch_count
FROM opener_work s
GROUP BY s.opener
HAVING MAX(s.state != 'complete') = 0
UNION ALL
SELECT s.opener,
MIN(s.requested_at) AS requested_at,
{started_at} AS started_at,
MIN(CASE WHEN s.state != 'complete' THEN s.started_at END)
AS started_at,
COALESCE(
MAX(CASE WHEN s.state != 'complete'
THEN s.requested_priority END),
Expand All @@ -6350,17 +6439,19 @@ def opener_rows(self):
COUNT(DISTINCT m.branch_id) AS branch_count,
COUNT(DISTINCT CASE WHEN m.parent_branch_id IS NULL
THEN m.branch_id END)
AS direct_branch_count
,COUNT(DISTINCT CASE
AS direct_branch_count,
COUNT(DISTINCT CASE
WHEN m.parent_branch_id IS NULL
AND p.status = 'done'
THEN m.branch_id END) AS direct_done_branch_count
FROM opener_work s
LEFT JOIN branch_opener_work m
ON m.opener_work_id = s.opener_work_id
LEFT JOIN pending_branches p ON p.branch_id = m.branch_id
WHERE s.opener IN (
SELECT opener FROM opener_work WHERE state != 'complete')
GROUP BY s.opener
ORDER BY requested_priority DESC, s.opener
ORDER BY requested_priority DESC, opener
""").fetchall()

def opener_direct_branch_keys(self, openers):
Expand Down
Loading
Loading