MCPcopy Create free account
hub / github.com/calebwin/pgclaw / process_queue_item

Function process_queue_item

src/worker.rs:88–166  ·  view source on GitHub ↗

Claim and process one item from claw.queue. Returns true if an item was processed.

()

Source from the content-addressed store, hash-verified

86/// Claim and process one item from claw.queue.
87/// Returns true if an item was processed.
88fn process_queue_item() -> bool {
89 let item = BackgroundWorker::transaction(|| {
90 Spi::connect(|client| {
91 let result = client.select(
92 "UPDATE claw.queue
93 SET claimed_at = now()
94 WHERE id = (
95 SELECT id FROM claw.queue
96 WHERE claimed_at IS NULL
97 ORDER BY id
98 FOR UPDATE SKIP LOCKED
99 LIMIT 1
100 )
101 RETURNING id, table_oid, pk_value, claw_col, row_data::text,
102 event, claw_value, agent_type",
103 None,
104 None,
105 ).expect("SPI failed");
106
107 let row = result.into_iter().next()?;
108
109 Some(QueueItem {
110 id: row.get::<i64>(1).ok()??,
111 table_oid: row.get::<pg_sys::Oid>(2).ok()??,
112 pk_value: row.get::<String>(3).ok()??,
113 claw_col: row.get::<String>(4).ok()??,
114 row_data: row.get::<String>(5).ok()??,
115 event: row.get::<String>(6).ok()??,
116 claw_value: row.get::<String>(7).ok()??,
117 agent_type: row.get::<String>(8).ok()??,
118 })
119 })
120 });
121
122 let item = match item {
123 Some(item) => item,
124 None => return false,
125 };
126
127 let result = if item.agent_type == "claudecode" {
128 process_claudecode_queue_item(&item)
129 } else {
130 process_model_queue_item(&item)
131 };
132
133 BackgroundWorker::transaction(|| {
134 Spi::connect(|mut client| {
135 match &result {
136 Ok(()) => {
137 client
138 .update(
139 &format!(
140 "UPDATE claw.queue SET done_at = now() WHERE id = {}",
141 item.id
142 ),
143 None,
144 None,
145 )

Callers 1

pgclaw_worker_mainFunction · 0.85

Calls 2

process_model_queue_itemFunction · 0.85

Tested by

no test coverage detected