Claim and process one item from claw.queue. Returns true if an item was processed.
()
| 86 | /// Claim and process one item from claw.queue. |
| 87 | /// Returns true if an item was processed. |
| 88 | fn 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 | ) |
no test coverage detected