1use std::fs;
15use std::path::Path;
16
17use chio_kernel::{
18 ApprovalDecision, ApprovalFilter, ApprovalOutcome, ApprovalRequest, ApprovalStore,
19 ApprovalStoreError, ResolvedApproval, ThresholdApprovalCollectorProposal,
20 ThresholdApprovalCollectorState, ThresholdApprovalCollectorStore,
21 ThresholdApprovalCollectorStoreError,
22};
23use r2d2::Pool;
24use r2d2_sqlite::SqliteConnectionManager;
25use rusqlite::{params, OptionalExtension};
26
27pub struct SqliteApprovalStore {
32 pool: Pool<SqliteConnectionManager>,
33}
34
35const APPROVAL_STORE_SUPPORTED_SCHEMA_VERSION: i32 = 2;
37const APPROVAL_STORE_SCHEMA_KEY: &str = "approval";
41const APPROVAL_STORE_OWN_ANCHOR_TABLES: &[&str] = &["chio_hitl_pending"];
55
56const APPROVAL_STORE_COLOCATED_ANCHOR_TABLES: &[&str] = &[
67 "chio_hitl_pending",
68 "http_receipts",
69 "tool_receipts",
70 "chio_tool_receipts",
71];
72
73impl SqliteApprovalStore {
74 pub fn open(path: impl AsRef<Path>) -> Result<Self, ApprovalStoreError> {
77 Self::open_with_anchor_tables(path, APPROVAL_STORE_OWN_ANCHOR_TABLES)
78 }
79
80 pub fn open_colocated_with_receipt_store(
90 path: impl AsRef<Path>,
91 ) -> Result<Self, ApprovalStoreError> {
92 Self::open_with_anchor_tables(path, APPROVAL_STORE_COLOCATED_ANCHOR_TABLES)
93 }
94
95 fn open_with_anchor_tables(
96 path: impl AsRef<Path>,
97 anchor_tables: &[&str],
98 ) -> Result<Self, ApprovalStoreError> {
99 let path = path.as_ref();
100 if let Some(parent) = crate::sqlite_parent_dir_to_create(path) {
106 fs::create_dir_all(&parent)
107 .map_err(|e| ApprovalStoreError::Backend(format!("create dir: {e}")))?;
108 }
109 let manager = SqliteConnectionManager::file(path);
110 let pool = Pool::builder()
111 .max_size(8)
112 .build(manager)
113 .map_err(|e| ApprovalStoreError::Backend(format!("pool build: {e}")))?;
114 let store = Self { pool };
115 store.run_migrations(anchor_tables)?;
116 Ok(store)
117 }
118
119 pub fn open_in_memory() -> Result<Self, ApprovalStoreError> {
121 let manager = SqliteConnectionManager::memory();
122 let pool = Pool::builder()
123 .max_size(1)
124 .build(manager)
125 .map_err(|e| ApprovalStoreError::Backend(format!("pool build: {e}")))?;
126 let store = Self { pool };
127 store.run_migrations(APPROVAL_STORE_OWN_ANCHOR_TABLES)?;
128 Ok(store)
129 }
130
131 fn run_migrations(&self, anchor_tables: &[&str]) -> Result<(), ApprovalStoreError> {
132 let mut conn = self
133 .pool
134 .get()
135 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
136 let on_disk = crate::check_schema_version(
137 &conn,
138 APPROVAL_STORE_SCHEMA_KEY,
139 APPROVAL_STORE_SUPPORTED_SCHEMA_VERSION,
140 anchor_tables,
141 )
142 .map_err(|error| ApprovalStoreError::Backend(error.to_string()))?;
143 conn.execute_batch(
144 r#"
145 PRAGMA journal_mode = WAL;
146 PRAGMA synchronous = FULL;
147 PRAGMA busy_timeout = 5000;
148 PRAGMA foreign_keys = ON;
149 "#,
150 )
151 .map_err(|e| ApprovalStoreError::Backend(format!("migration setup: {e}")))?;
152 if on_disk < 2 {
153 let votes_table_exists: bool = conn
154 .query_row(
155 r#"
156 SELECT EXISTS(
157 SELECT 1 FROM sqlite_master
158 WHERE type = 'table'
159 AND name = 'chio_threshold_approval_collector_votes'
160 )
161 "#,
162 [],
163 |row| row.get(0),
164 )
165 .map_err(|e| ApprovalStoreError::Backend(format!("migration probe: {e}")))?;
166 if votes_table_exists {
167 let transaction = conn
168 .transaction()
169 .map_err(|e| ApprovalStoreError::Backend(format!("migration begin: {e}")))?;
170 transaction
171 .execute_batch(
172 r#"
173 ALTER TABLE chio_threshold_approval_collector_votes
174 RENAME TO chio_threshold_approval_collector_votes_v1;
175 CREATE TABLE chio_threshold_approval_collector_votes (
176 proposal_id TEXT NOT NULL,
177 token_id TEXT NOT NULL,
178 approver_fingerprint TEXT NOT NULL,
179 canonical_token_digest TEXT NOT NULL UNIQUE,
180 token_json BLOB NOT NULL,
181 received_at INTEGER NOT NULL,
182 PRIMARY KEY (proposal_id, token_id),
183 UNIQUE (proposal_id, approver_fingerprint),
184 UNIQUE (proposal_id, canonical_token_digest),
185 FOREIGN KEY (proposal_id)
186 REFERENCES chio_threshold_approval_collectors(proposal_id)
187 );
188 INSERT INTO chio_threshold_approval_collector_votes (
189 proposal_id, token_id, approver_fingerprint,
190 canonical_token_digest, token_json, received_at
191 )
192 SELECT proposal_id, token_id, approver_fingerprint,
193 canonical_token_digest, token_json, received_at
194 FROM chio_threshold_approval_collector_votes_v1;
195 DROP TABLE chio_threshold_approval_collector_votes_v1;
196 "#,
197 )
198 .map_err(|e| {
199 ApprovalStoreError::Backend(format!(
200 "threshold vote uniqueness migration: {e}"
201 ))
202 })?;
203 transaction
204 .commit()
205 .map_err(|e| ApprovalStoreError::Backend(format!("migration commit: {e}")))?;
206 }
207 }
208 conn.execute_batch(
209 r#"
210
211 CREATE TABLE IF NOT EXISTS chio_hitl_pending (
212 approval_id TEXT PRIMARY KEY,
213 policy_id TEXT NOT NULL,
214 subject_id TEXT NOT NULL,
215 tool_server TEXT NOT NULL,
216 tool_name TEXT NOT NULL,
217 parameter_hash TEXT NOT NULL,
218 expires_at INTEGER NOT NULL,
219 created_at INTEGER NOT NULL,
220 payload TEXT NOT NULL
221 );
222 CREATE INDEX IF NOT EXISTS idx_chio_hitl_pending_subject
223 ON chio_hitl_pending(subject_id);
224 CREATE INDEX IF NOT EXISTS idx_chio_hitl_pending_expires
225 ON chio_hitl_pending(expires_at);
226
227 CREATE TABLE IF NOT EXISTS chio_hitl_resolved (
228 approval_id TEXT PRIMARY KEY,
229 policy_id TEXT NOT NULL,
230 subject_id TEXT NOT NULL,
231 outcome TEXT NOT NULL,
232 resolved_at INTEGER NOT NULL,
233 approver_hex TEXT NOT NULL,
234 token_id TEXT NOT NULL
235 );
236 CREATE INDEX IF NOT EXISTS idx_chio_hitl_resolved_counts
237 ON chio_hitl_resolved(subject_id, policy_id, outcome);
238
239 CREATE TABLE IF NOT EXISTS chio_hitl_consumed_tokens (
240 token_id TEXT NOT NULL,
241 parameter_hash TEXT NOT NULL,
242 consumed_at INTEGER NOT NULL,
243 PRIMARY KEY (token_id, parameter_hash)
244 );
245
246 CREATE TABLE IF NOT EXISTS chio_threshold_approval_collectors (
247 proposal_id TEXT PRIMARY KEY,
248 request_id TEXT NOT NULL,
249 governed_intent_hash TEXT NOT NULL,
250 subject_fingerprint TEXT NOT NULL,
251 authorizing_capability_digest TEXT NOT NULL,
252 policy_hash TEXT NOT NULL,
253 threshold INTEGER NOT NULL CHECK (threshold > 0),
254 eligible_set_digest TEXT NOT NULL,
255 proposal_created_at INTEGER NOT NULL,
256 proposal_deadline INTEGER NOT NULL,
257 submitter_fingerprint TEXT,
258 require_submitter_separation INTEGER NOT NULL CHECK (
259 require_submitter_separation IN (0, 1)
260 ),
261 state TEXT NOT NULL CHECK (
262 state IN ('collecting', 'ready', 'delivered', 'cancelled')
263 ),
264 version INTEGER NOT NULL CHECK (version >= 0),
265 updated_at INTEGER NOT NULL,
266 proposal_json BLOB NOT NULL,
267 requirement_json BLOB NOT NULL,
268 record_json BLOB NOT NULL
269 );
270
271 CREATE TABLE IF NOT EXISTS chio_threshold_approval_collector_votes (
272 proposal_id TEXT NOT NULL,
273 token_id TEXT NOT NULL,
274 approver_fingerprint TEXT NOT NULL,
275 canonical_token_digest TEXT NOT NULL UNIQUE,
276 token_json BLOB NOT NULL,
277 received_at INTEGER NOT NULL,
278 PRIMARY KEY (proposal_id, token_id),
279 UNIQUE (proposal_id, approver_fingerprint),
280 UNIQUE (proposal_id, canonical_token_digest),
281 FOREIGN KEY (proposal_id)
282 REFERENCES chio_threshold_approval_collectors(proposal_id)
283 );
284 "#,
285 )
286 .map_err(|e| ApprovalStoreError::Backend(format!("migration: {e}")))?;
287 crate::stamp_schema_version(
288 &conn,
289 APPROVAL_STORE_SCHEMA_KEY,
290 APPROVAL_STORE_SUPPORTED_SCHEMA_VERSION,
291 )
292 .map_err(|error| ApprovalStoreError::Backend(error.to_string()))?;
293 Ok(())
294 }
295}
296
297fn collector_state_name(state: ThresholdApprovalCollectorState) -> &'static str {
298 match state {
299 ThresholdApprovalCollectorState::Collecting => "collecting",
300 ThresholdApprovalCollectorState::Ready => "ready",
301 ThresholdApprovalCollectorState::Delivered => "delivered",
302 ThresholdApprovalCollectorState::Cancelled => "cancelled",
303 }
304}
305
306fn collector_error(error: impl std::fmt::Display) -> ThresholdApprovalCollectorStoreError {
307 ThresholdApprovalCollectorStoreError::Backend(error.to_string())
308}
309
310fn encode_collector<T: serde::Serialize>(
311 value: &T,
312) -> Result<Vec<u8>, ThresholdApprovalCollectorStoreError> {
313 chio_core::canonical::canonical_json_bytes(value)
314 .map_err(|error| ThresholdApprovalCollectorStoreError::Serialization(error.to_string()))
315}
316
317fn decode_collector(
318 bytes: &[u8],
319) -> Result<ThresholdApprovalCollectorProposal, ThresholdApprovalCollectorStoreError> {
320 serde_json::from_slice(bytes)
321 .map_err(|error| ThresholdApprovalCollectorStoreError::Serialization(error.to_string()))
322}
323
324fn serialize_payload(request: &ApprovalRequest) -> Result<String, ApprovalStoreError> {
325 serde_json::to_string(request).map_err(|e| ApprovalStoreError::Serialization(e.to_string()))
326}
327
328fn deserialize_payload(raw: &str) -> Result<ApprovalRequest, ApprovalStoreError> {
329 serde_json::from_str(raw).map_err(|e| ApprovalStoreError::Serialization(e.to_string()))
330}
331
332impl ApprovalStore for SqliteApprovalStore {
333 fn store_pending(&self, request: &ApprovalRequest) -> Result<(), ApprovalStoreError> {
334 let payload = serialize_payload(request)?;
335 let conn = self
336 .pool
337 .get()
338 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
339 let returned_payload = conn
340 .query_row(
341 r#"
342 INSERT INTO chio_hitl_pending (approval_id, policy_id, subject_id, tool_server, tool_name, parameter_hash, expires_at, created_at, payload) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9) ON CONFLICT(approval_id) DO UPDATE SET payload = excluded.payload WHERE chio_hitl_pending.payload = excluded.payload RETURNING payload
343 "#,
344 params![
345 request.approval_id,
346 request.policy_id,
347 request.subject_id,
348 request.tool_server,
349 request.tool_name,
350 request.parameter_hash,
351 request.expires_at as i64,
352 request.created_at as i64,
353 payload,
354 ],
355 |row| row.get::<_, String>(0),
356 )
357 .optional()
358 .map_err(|e| ApprovalStoreError::Backend(format!("insert pending: {e}")))?;
359 if returned_payload.is_none() {
360 return Err(ApprovalStoreError::Backend(format!(
361 "approval_id {} already exists with different payload",
362 request.approval_id
363 )));
364 }
365 Ok(())
366 }
367
368 fn get_pending(&self, id: &str) -> Result<Option<ApprovalRequest>, ApprovalStoreError> {
369 let conn = self
370 .pool
371 .get()
372 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
373 let row: Option<String> = conn
374 .query_row(
375 "SELECT payload FROM chio_hitl_pending WHERE approval_id = ?1",
376 params![id],
377 |row| row.get::<_, String>(0),
378 )
379 .optional()
380 .map_err(|e| ApprovalStoreError::Backend(format!("select pending: {e}")))?;
381 match row {
382 Some(raw) => Ok(Some(deserialize_payload(&raw)?)),
383 None => Ok(None),
384 }
385 }
386
387 fn list_pending(
388 &self,
389 filter: &ApprovalFilter,
390 ) -> Result<Vec<ApprovalRequest>, ApprovalStoreError> {
391 let conn = self
392 .pool
393 .get()
394 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
395 let mut sql = String::from("SELECT payload FROM chio_hitl_pending WHERE 1=1");
396 if filter.subject_id.is_some() {
397 sql.push_str(" AND subject_id = :subject_id");
398 }
399 if filter.tool_server.is_some() {
400 sql.push_str(" AND tool_server = :tool_server");
401 }
402 if filter.tool_name.is_some() {
403 sql.push_str(" AND tool_name = :tool_name");
404 }
405 if filter.not_expired_at.is_some() {
406 sql.push_str(" AND expires_at > :not_expired_at");
407 }
408 sql.push_str(" ORDER BY created_at ASC");
409 if filter.limit.is_some() {
410 sql.push_str(" LIMIT :limit");
411 }
412
413 let mut stmt = conn
414 .prepare(&sql)
415 .map_err(|e| ApprovalStoreError::Backend(format!("prepare list: {e}")))?;
416
417 let mut params_vec: Vec<(&str, Box<dyn rusqlite::ToSql>)> = Vec::new();
418 if let Some(s) = &filter.subject_id {
419 params_vec.push((":subject_id", Box::new(s.clone())));
420 }
421 if let Some(s) = &filter.tool_server {
422 params_vec.push((":tool_server", Box::new(s.clone())));
423 }
424 if let Some(s) = &filter.tool_name {
425 params_vec.push((":tool_name", Box::new(s.clone())));
426 }
427 if let Some(t) = &filter.not_expired_at {
428 params_vec.push((":not_expired_at", Box::new(*t as i64)));
429 }
430 if let Some(limit) = &filter.limit {
431 params_vec.push((":limit", Box::new(*limit as i64)));
432 }
433
434 let refs: Vec<(&str, &dyn rusqlite::ToSql)> = params_vec
435 .iter()
436 .map(|(name, value)| (*name, value.as_ref()))
437 .collect();
438
439 let rows = stmt
440 .query_map(refs.as_slice(), |row| row.get::<_, String>(0))
441 .map_err(|e| ApprovalStoreError::Backend(format!("query list: {e}")))?;
442
443 let mut out = Vec::new();
444 for row in rows {
445 let raw = row.map_err(|e| ApprovalStoreError::Backend(format!("row: {e}")))?;
446 out.push(deserialize_payload(&raw)?);
447 }
448 Ok(out)
449 }
450
451 fn resolve(&self, id: &str, decision: &ApprovalDecision) -> Result<(), ApprovalStoreError> {
452 let mut conn = self
453 .pool
454 .get()
455 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
456 let tx = conn
457 .transaction()
458 .map_err(|e| ApprovalStoreError::Backend(format!("begin tx: {e}")))?;
459
460 let pending: Option<(String, String)> = tx
462 .query_row(
463 "SELECT policy_id, parameter_hash FROM chio_hitl_pending WHERE approval_id = ?1",
464 params![id],
465 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
466 )
467 .optional()
468 .map_err(|e| ApprovalStoreError::Backend(format!("select: {e}")))?;
469 let (policy_id, parameter_hash) = match pending {
470 Some(p) => p,
471 None => return Err(ApprovalStoreError::NotFound(id.to_string())),
472 };
473
474 let already: Option<i64> = tx
476 .query_row(
477 "SELECT 1 FROM chio_hitl_consumed_tokens WHERE token_id = ?1 AND parameter_hash = ?2",
478 params![decision.token.id, parameter_hash],
479 |row| row.get(0),
480 )
481 .optional()
482 .map_err(|e| ApprovalStoreError::Backend(format!("replay check: {e}")))?;
483 if already.is_some() {
484 return Err(ApprovalStoreError::Replay(id.to_string()));
485 }
486
487 let already_resolved: Option<i64> = tx
489 .query_row(
490 "SELECT 1 FROM chio_hitl_resolved WHERE approval_id = ?1",
491 params![id],
492 |row| row.get(0),
493 )
494 .optional()
495 .map_err(|e| ApprovalStoreError::Backend(format!("resolved check: {e}")))?;
496 if already_resolved.is_some() {
497 return Err(ApprovalStoreError::AlreadyResolved(id.to_string()));
498 }
499
500 let outcome = match decision.outcome {
501 ApprovalOutcome::Approved => "approved",
502 ApprovalOutcome::Denied => "denied",
503 };
504
505 tx.execute(
506 r#"INSERT INTO chio_hitl_resolved (
507 approval_id, policy_id, subject_id, outcome, resolved_at,
508 approver_hex, token_id
509 ) SELECT approval_id, policy_id, subject_id, ?2, ?3, ?4, ?5
510 FROM chio_hitl_pending WHERE approval_id = ?1"#,
511 params![
512 id,
513 outcome,
514 decision.received_at as i64,
515 decision.approver.to_hex(),
516 decision.token.id,
517 ],
518 )
519 .map_err(|e| ApprovalStoreError::Backend(format!("insert resolved: {e}")))?;
520
521 tx.execute(
522 "INSERT INTO chio_hitl_consumed_tokens (token_id, parameter_hash, consumed_at) VALUES (?1, ?2, ?3)",
523 params![decision.token.id, parameter_hash, decision.received_at as i64],
524 )
525 .map_err(|e| ApprovalStoreError::Backend(format!("insert consumed: {e}")))?;
526
527 tx.execute(
528 "DELETE FROM chio_hitl_pending WHERE approval_id = ?1",
529 params![id],
530 )
531 .map_err(|e| ApprovalStoreError::Backend(format!("delete pending: {e}")))?;
532
533 tx.commit()
534 .map_err(|e| ApprovalStoreError::Backend(format!("commit: {e}")))?;
535
536 let _ = policy_id;
538 Ok(())
539 }
540
541 fn count_approved(&self, subject_id: &str, policy_id: &str) -> Result<u64, ApprovalStoreError> {
542 let conn = self
543 .pool
544 .get()
545 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
546 let count: i64 = conn
547 .query_row(
548 "SELECT COUNT(*) FROM chio_hitl_resolved WHERE subject_id = ?1 AND policy_id = ?2 AND outcome = 'approved'",
549 params![subject_id, policy_id],
550 |row| row.get(0),
551 )
552 .map_err(|e| ApprovalStoreError::Backend(format!("count: {e}")))?;
553 Ok(count.max(0) as u64)
554 }
555
556 fn record_consumed(
557 &self,
558 token_id: &str,
559 parameter_hash: &str,
560 now: u64,
561 ) -> Result<(), ApprovalStoreError> {
562 let conn = self
563 .pool
564 .get()
565 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
566 let rows = conn.execute(
567 "INSERT OR IGNORE INTO chio_hitl_consumed_tokens (token_id, parameter_hash, consumed_at) VALUES (?1, ?2, ?3)",
568 params![token_id, parameter_hash, now as i64],
569 )
570 .map_err(|e| ApprovalStoreError::Backend(format!("insert consumed: {e}")))?;
571 if rows == 0 {
572 return Err(ApprovalStoreError::Replay(format!(
573 "token {token_id} already consumed"
574 )));
575 }
576 Ok(())
577 }
578
579 fn is_consumed(
580 &self,
581 token_id: &str,
582 parameter_hash: &str,
583 ) -> Result<bool, ApprovalStoreError> {
584 let conn = self
585 .pool
586 .get()
587 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
588 let row: Option<i64> = conn
589 .query_row(
590 "SELECT 1 FROM chio_hitl_consumed_tokens WHERE token_id = ?1 AND parameter_hash = ?2",
591 params![token_id, parameter_hash],
592 |row| row.get(0),
593 )
594 .optional()
595 .map_err(|e| ApprovalStoreError::Backend(format!("is_consumed: {e}")))?;
596 Ok(row.is_some())
597 }
598
599 fn get_resolution(&self, id: &str) -> Result<Option<ResolvedApproval>, ApprovalStoreError> {
600 let conn = self
601 .pool
602 .get()
603 .map_err(|e| ApprovalStoreError::Backend(format!("pool get: {e}")))?;
604 let row: Option<(String, String, i64, String, String)> = conn
605 .query_row(
606 r#"SELECT approval_id, outcome, resolved_at, approver_hex, token_id
607 FROM chio_hitl_resolved WHERE approval_id = ?1"#,
608 params![id],
609 |row| {
610 Ok((
611 row.get::<_, String>(0)?,
612 row.get::<_, String>(1)?,
613 row.get::<_, i64>(2)?,
614 row.get::<_, String>(3)?,
615 row.get::<_, String>(4)?,
616 ))
617 },
618 )
619 .optional()
620 .map_err(|e| ApprovalStoreError::Backend(format!("get_resolution: {e}")))?;
621 match row {
622 Some((approval_id, outcome_str, resolved_at, approver_hex, token_id)) => {
623 let outcome = match outcome_str.as_str() {
624 "approved" => ApprovalOutcome::Approved,
625 "denied" => ApprovalOutcome::Denied,
626 other => {
627 return Err(ApprovalStoreError::Serialization(format!(
628 "unknown outcome: {other}"
629 )))
630 }
631 };
632 Ok(Some(ResolvedApproval {
633 approval_id,
634 outcome,
635 resolved_at: resolved_at.max(0) as u64,
636 approver_hex,
637 token_id,
638 }))
639 }
640 None => Ok(None),
641 }
642 }
643}
644
645impl ThresholdApprovalCollectorStore for SqliteApprovalStore {
646 fn create(
647 &self,
648 proposal: &ThresholdApprovalCollectorProposal,
649 ) -> Result<(), ThresholdApprovalCollectorStoreError> {
650 let proposal_json = encode_collector(&proposal.proposal)?;
651 let requirement_json = encode_collector(&proposal.requirement)?;
652 let record_json = encode_collector(proposal)?;
653 let conn = self.pool.get().map_err(collector_error)?;
654 let existing: Option<Vec<u8>> = conn
655 .query_row(
656 "SELECT record_json FROM chio_threshold_approval_collectors WHERE proposal_id = ?1",
657 [&proposal.proposal.body.proposal_id],
658 |row| row.get(0),
659 )
660 .optional()
661 .map_err(collector_error)?;
662 if let Some(existing) = existing {
663 return if existing == record_json {
664 Ok(())
665 } else {
666 Err(ThresholdApprovalCollectorStoreError::Conflict(
667 "proposal id already exists with different content".to_string(),
668 ))
669 };
670 }
671 let body = &proposal.proposal.body;
672 conn.execute(
673 r#"
674 INSERT INTO chio_threshold_approval_collectors (
675 proposal_id, request_id, governed_intent_hash, subject_fingerprint,
676 authorizing_capability_digest, policy_hash, threshold,
677 eligible_set_digest, proposal_created_at, proposal_deadline,
678 submitter_fingerprint, require_submitter_separation, state,
679 version, updated_at, proposal_json, requirement_json, record_json
680 ) VALUES (
681 ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13,
682 ?14, ?15, ?16, ?17, ?18
683 )
684 "#,
685 params![
686 &body.proposal_id,
687 &body.request_id,
688 &body.governed_intent_hash,
689 body.subject.to_hex(),
690 &body.authorizing_capability_digest,
691 &body.policy_hash,
692 i64::from(body.threshold),
693 &body.eligible_set_digest,
694 i64::try_from(body.proposal_created_at).map_err(collector_error)?,
695 i64::try_from(body.proposal_deadline).map_err(collector_error)?,
696 proposal.submitter.as_ref().map(|key| key.to_hex()),
697 i64::from(proposal.require_submitter_separation),
698 collector_state_name(proposal.state),
699 i64::try_from(proposal.version).map_err(collector_error)?,
700 i64::try_from(proposal.updated_at).map_err(collector_error)?,
701 proposal_json,
702 requirement_json,
703 record_json,
704 ],
705 )
706 .map_err(collector_error)?;
707 Ok(())
708 }
709
710 fn get(
711 &self,
712 proposal_id: &str,
713 ) -> Result<Option<ThresholdApprovalCollectorProposal>, ThresholdApprovalCollectorStoreError>
714 {
715 let conn = self.pool.get().map_err(collector_error)?;
716 let record: Option<Vec<u8>> = conn
717 .query_row(
718 "SELECT record_json FROM chio_threshold_approval_collectors WHERE proposal_id = ?1",
719 [proposal_id],
720 |row| row.get(0),
721 )
722 .optional()
723 .map_err(collector_error)?;
724 record.map(|bytes| decode_collector(&bytes)).transpose()
725 }
726
727 fn append_token(
728 &self,
729 proposal_id: &str,
730 expected_version: u64,
731 token: &chio_core::capability::governance::GovernedApprovalToken,
732 replaced_token_id: Option<&str>,
733 next_state: ThresholdApprovalCollectorState,
734 updated_at: u64,
735 ) -> Result<ThresholdApprovalCollectorProposal, ThresholdApprovalCollectorStoreError> {
736 let mut conn = self.pool.get().map_err(collector_error)?;
737 let transaction = conn.transaction().map_err(collector_error)?;
738 let bytes: Option<Vec<u8>> = transaction
739 .query_row(
740 "SELECT record_json FROM chio_threshold_approval_collectors WHERE proposal_id = ?1",
741 [proposal_id],
742 |row| row.get(0),
743 )
744 .optional()
745 .map_err(collector_error)?;
746 let mut record = bytes
747 .as_deref()
748 .map(decode_collector)
749 .transpose()?
750 .ok_or_else(|| ThresholdApprovalCollectorStoreError::NotFound(proposal_id.into()))?;
751 if record.version != expected_version || record.state.is_terminal() {
752 return Err(ThresholdApprovalCollectorStoreError::Conflict(
753 "threshold approval proposal changed concurrently".to_string(),
754 ));
755 }
756 let previous_state = collector_state_name(record.state);
757 let token_digest = token.artifact_digest().map_err(|error| {
758 ThresholdApprovalCollectorStoreError::Serialization(error.to_string())
759 })?;
760 let token_json = encode_collector(token)?;
761 let write_result = if let Some(replaced_token_id) = replaced_token_id {
762 transaction.execute(
763 r#"
764 UPDATE chio_threshold_approval_collector_votes
765 SET token_id = ?1, approver_fingerprint = ?2,
766 canonical_token_digest = ?3, token_json = ?4, received_at = ?5
767 WHERE proposal_id = ?6 AND token_id = ?7
768 "#,
769 params![
770 &token.id,
771 token.approver.to_hex(),
772 token_digest,
773 token_json,
774 i64::try_from(updated_at).map_err(collector_error)?,
775 proposal_id,
776 replaced_token_id,
777 ],
778 )
779 } else {
780 transaction.execute(
781 r#"
782 INSERT INTO chio_threshold_approval_collector_votes (
783 proposal_id, token_id, approver_fingerprint,
784 canonical_token_digest, token_json, received_at
785 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
786 "#,
787 params![
788 proposal_id,
789 &token.id,
790 token.approver.to_hex(),
791 token_digest,
792 token_json,
793 i64::try_from(updated_at).map_err(collector_error)?,
794 ],
795 )
796 };
797 let changed_vote = write_result.map_err(|error| {
798 if matches!(
799 &error,
800 rusqlite::Error::SqliteFailure(
801 rusqlite::ffi::Error {
802 code: rusqlite::ErrorCode::ConstraintViolation,
803 ..
804 },
805 _
806 )
807 ) {
808 ThresholdApprovalCollectorStoreError::Conflict(
809 "threshold approval token id, digest, or signer is not unique".to_string(),
810 )
811 } else {
812 collector_error(error)
813 }
814 })?;
815 if changed_vote != 1 {
816 return Err(ThresholdApprovalCollectorStoreError::Conflict(
817 "threshold approval replacement token disappeared".to_string(),
818 ));
819 }
820 if let Some(replaced_token_id) = replaced_token_id {
821 let existing = record
822 .tokens
823 .iter_mut()
824 .find(|existing| existing.id == replaced_token_id)
825 .ok_or_else(|| {
826 ThresholdApprovalCollectorStoreError::Conflict(
827 "threshold approval replacement token disappeared".to_string(),
828 )
829 })?;
830 *existing = token.clone();
831 } else {
832 record.tokens.push(token.clone());
833 }
834 record.state = next_state;
835 record.version = record.version.checked_add(1).ok_or_else(|| {
836 ThresholdApprovalCollectorStoreError::Conflict("proposal version overflowed".into())
837 })?;
838 record.updated_at = updated_at;
839 let record_json = encode_collector(&record)?;
840 let changed = transaction
841 .execute(
842 r#"
843 UPDATE chio_threshold_approval_collectors
844 SET state = ?1, version = ?2, updated_at = ?3, record_json = ?4
845 WHERE proposal_id = ?5 AND version = ?6 AND state = ?7
846 "#,
847 params![
848 collector_state_name(next_state),
849 i64::try_from(record.version).map_err(collector_error)?,
850 i64::try_from(updated_at).map_err(collector_error)?,
851 record_json,
852 proposal_id,
853 i64::try_from(expected_version).map_err(collector_error)?,
854 previous_state,
855 ],
856 )
857 .map_err(collector_error)?;
858 if changed != 1 {
859 return Err(ThresholdApprovalCollectorStoreError::Conflict(
860 "threshold approval proposal changed concurrently".to_string(),
861 ));
862 }
863 transaction.commit().map_err(collector_error)?;
864 Ok(record)
865 }
866
867 fn transition(
868 &self,
869 proposal_id: &str,
870 expected_version: u64,
871 next_state: ThresholdApprovalCollectorState,
872 updated_at: u64,
873 ) -> Result<ThresholdApprovalCollectorProposal, ThresholdApprovalCollectorStoreError> {
874 let mut conn = self.pool.get().map_err(collector_error)?;
875 let transaction = conn.transaction().map_err(collector_error)?;
876 let bytes: Option<Vec<u8>> = transaction
877 .query_row(
878 "SELECT record_json FROM chio_threshold_approval_collectors WHERE proposal_id = ?1",
879 [proposal_id],
880 |row| row.get(0),
881 )
882 .optional()
883 .map_err(collector_error)?;
884 let mut record = bytes
885 .as_deref()
886 .map(decode_collector)
887 .transpose()?
888 .ok_or_else(|| ThresholdApprovalCollectorStoreError::NotFound(proposal_id.into()))?;
889 if record.version != expected_version || record.state.is_terminal() {
890 return Err(ThresholdApprovalCollectorStoreError::Conflict(
891 "threshold approval proposal changed concurrently".to_string(),
892 ));
893 }
894 let previous_state = collector_state_name(record.state);
895 record.state = next_state;
896 record.version = record.version.checked_add(1).ok_or_else(|| {
897 ThresholdApprovalCollectorStoreError::Conflict("proposal version overflowed".into())
898 })?;
899 record.updated_at = updated_at;
900 let record_json = encode_collector(&record)?;
901 let changed = transaction
902 .execute(
903 r#"
904 UPDATE chio_threshold_approval_collectors
905 SET state = ?1, version = ?2, updated_at = ?3, record_json = ?4
906 WHERE proposal_id = ?5 AND version = ?6 AND state = ?7
907 "#,
908 params![
909 collector_state_name(next_state),
910 i64::try_from(record.version).map_err(collector_error)?,
911 i64::try_from(updated_at).map_err(collector_error)?,
912 record_json,
913 proposal_id,
914 i64::try_from(expected_version).map_err(collector_error)?,
915 previous_state,
916 ],
917 )
918 .map_err(collector_error)?;
919 if changed != 1 {
920 return Err(ThresholdApprovalCollectorStoreError::Conflict(
921 "threshold approval proposal changed concurrently".to_string(),
922 ));
923 }
924 transaction.commit().map_err(collector_error)?;
925 Ok(record)
926 }
927}
928
929#[cfg(test)]
930#[allow(clippy::expect_used, clippy::unwrap_used)]
931mod tests {
932 use super::*;
933 use chio_core::capability::governance::{
934 GovernedApprovalDecision, GovernedApprovalToken, GovernedApprovalTokenBody,
935 ThresholdApprovalProposal, ThresholdApprovalProposalBody,
936 THRESHOLD_APPROVAL_PROPOSAL_SCHEMA,
937 };
938 use chio_core::capability::threshold_approval::{
939 ThresholdApprovalRequirement, ThresholdApproverIdentity,
940 };
941 use chio_core::crypto::{sha256_hex, Keypair};
942 use chio_kernel::ThresholdApprovalCollector;
943 use std::sync::Arc;
944 use std::time::{SystemTime, UNIX_EPOCH};
945
946 #[test]
947 fn open_colocated_creates_parent_dirs_for_a_file_uri_with_query() {
948 let nonce = SystemTime::now()
949 .duration_since(UNIX_EPOCH)
950 .expect("time before epoch")
951 .as_nanos();
952 let base = std::env::temp_dir().join(format!("chio-approval-uri-{nonce}"));
953 let db = base.join("nested").join("receipts.db");
954 let parent = db.parent().expect("db path has a parent");
955 assert!(
956 !parent.exists(),
957 "precondition: the parent dir must not exist yet"
958 );
959
960 let uri = format!("file:{}?mode=rwc", db.display());
965 let store = SqliteApprovalStore::open_colocated_with_receipt_store(uri.as_str())
966 .expect("open colocated approval store from a file: URI");
967 store
969 .store_pending(&sample_request("uri-1", "hash-uri"))
970 .expect("store a pending approval");
971
972 assert!(
973 parent.exists(),
974 "the real parent directory must be created before SQLite opens the URI"
975 );
976
977 let _ = fs::remove_dir_all(&base);
978 }
979
980 fn sample_request(id: &str, hash: &str) -> ApprovalRequest {
981 let subject = Keypair::generate();
982 let approver = Keypair::generate();
983 ApprovalRequest {
984 approval_id: id.into(),
985 policy_id: "policy-1".into(),
986 subject_id: "agent-1".into(),
987 capability_id: "cap-1".into(),
988 subject_public_key: Some(subject.public_key()),
989 tool_server: "srv".into(),
990 tool_name: "tool".into(),
991 action: "invoke".into(),
992 parameter_hash: hash.into(),
993 expires_at: 1_000_000,
994 callback_hint: None,
995 created_at: 42,
996 summary: "unit".into(),
997 governed_intent: None,
998 trusted_approvers: vec![approver.public_key()],
999 triggered_by: vec![],
1000 }
1001 }
1002
1003 #[test]
1004 fn store_and_list_round_trip() {
1005 let store = SqliteApprovalStore::open_in_memory().unwrap();
1006 let r1 = sample_request("a-1", "h-1");
1007 let r2 = sample_request("a-2", "h-2");
1008 store.store_pending(&r1).unwrap();
1009 store.store_pending(&r2).unwrap();
1010
1011 let all = store.list_pending(&ApprovalFilter::default()).unwrap();
1012 assert_eq!(all.len(), 2);
1013
1014 let fetched = store.get_pending("a-1").unwrap().unwrap();
1015 assert_eq!(fetched.approval_id, "a-1");
1016 assert_eq!(fetched.parameter_hash, "h-1");
1017 }
1018
1019 #[test]
1020 fn duplicate_pending_insert_is_idempotent_only_when_payload_matches() {
1021 let store = SqliteApprovalStore::open_in_memory().unwrap();
1022 let original = sample_request("dup-1", "hash-a");
1023 let identical = original.clone();
1024 let mut mismatched = original.clone();
1025 mismatched.parameter_hash = "hash-b".into();
1026
1027 store.store_pending(&original).unwrap();
1028 store.store_pending(&identical).unwrap();
1029
1030 let err = store.store_pending(&mismatched).unwrap_err();
1031 match err {
1032 ApprovalStoreError::Backend(message) => {
1033 assert!(message.contains("already exists with different payload"));
1034 }
1035 other => panic!("expected Backend mismatch error, got {other:?}"),
1036 }
1037
1038 let fetched = store.get_pending("dup-1").unwrap().unwrap();
1039 assert_eq!(fetched.parameter_hash, "hash-a");
1040 }
1041
1042 #[test]
1043 fn standalone_open_refuses_a_receipt_sidecar_that_colocated_open_adopts() {
1044 let dir = tempfile::tempdir().unwrap();
1053 let path = dir.path().join("sidecar.sqlite3");
1054 {
1055 let conn = rusqlite::Connection::open(&path).unwrap();
1056 conn.execute_batch(
1057 "CREATE TABLE http_receipts (id TEXT PRIMARY KEY, receipt_json TEXT NOT NULL);
1058 CREATE TABLE tool_receipts (id TEXT PRIMARY KEY, receipt_json TEXT NOT NULL);
1059 CREATE TABLE revoked_capabilities (capability_id TEXT PRIMARY KEY);",
1060 )
1061 .unwrap();
1062 let app_id: i32 = conn
1063 .query_row("PRAGMA application_id", [], |row| row.get(0))
1064 .unwrap();
1065 assert_eq!(
1066 app_id, 0,
1067 "fixture must be unstamped like a legacy database"
1068 );
1069 }
1070
1071 assert!(
1072 SqliteApprovalStore::open(&path).is_err(),
1073 "standalone approval open must refuse a receipt-only sidecar file"
1074 );
1075
1076 let store = SqliteApprovalStore::open_colocated_with_receipt_store(&path)
1077 .expect("co-located open must adopt the receipt sidecar file");
1078 store
1079 .store_pending(&sample_request("adopt-1", "hash-adopt"))
1080 .unwrap();
1081 assert!(store.get_pending("adopt-1").unwrap().is_some());
1082 }
1083
1084 #[test]
1085 fn standalone_open_reopens_a_genuine_approval_database() {
1086 let dir = tempfile::tempdir().unwrap();
1089 let path = dir.path().join("approval.sqlite3");
1090 {
1091 let store = SqliteApprovalStore::open(&path).unwrap();
1092 store
1093 .store_pending(&sample_request("reopen-1", "hash-reopen"))
1094 .unwrap();
1095 }
1096 let store = SqliteApprovalStore::open(&path)
1097 .expect("a genuine approval database must reopen standalone");
1098 assert!(store.get_pending("reopen-1").unwrap().is_some());
1099 }
1100
1101 #[test]
1102 fn open_migrates_v1_threshold_vote_ids_to_proposal_scope() {
1103 let dir = tempfile::tempdir().unwrap();
1104 let path = dir.path().join("approval-v1.sqlite3");
1105 drop(SqliteApprovalStore::open(&path).unwrap());
1106
1107 let connection = rusqlite::Connection::open(&path).unwrap();
1108 connection
1109 .execute_batch(
1110 r#"
1111 ALTER TABLE chio_threshold_approval_collector_votes
1112 RENAME TO chio_threshold_approval_collector_votes_v2;
1113 CREATE TABLE chio_threshold_approval_collector_votes (
1114 proposal_id TEXT NOT NULL,
1115 token_id TEXT NOT NULL UNIQUE,
1116 approver_fingerprint TEXT NOT NULL,
1117 canonical_token_digest TEXT NOT NULL UNIQUE,
1118 token_json BLOB NOT NULL,
1119 received_at INTEGER NOT NULL,
1120 UNIQUE (proposal_id, approver_fingerprint),
1121 UNIQUE (proposal_id, canonical_token_digest),
1122 FOREIGN KEY (proposal_id)
1123 REFERENCES chio_threshold_approval_collectors(proposal_id)
1124 );
1125 DROP TABLE chio_threshold_approval_collector_votes_v2;
1126 INSERT INTO chio_threshold_approval_collectors (
1127 proposal_id, request_id, governed_intent_hash,
1128 subject_fingerprint, authorizing_capability_digest,
1129 policy_hash, threshold, eligible_set_digest,
1130 proposal_created_at, proposal_deadline, submitter_fingerprint,
1131 require_submitter_separation, state, version, updated_at,
1132 proposal_json, requirement_json, record_json
1133 ) VALUES (
1134 'proposal-v1', 'request-v1', 'intent-v1', 'subject-v1',
1135 'capability-v1', 'policy-v1', 1, 'eligible-v1',
1136 100, 200, NULL, 0, 'collecting', 1, 100,
1137 X'01', X'01', X'01'
1138 );
1139 INSERT INTO chio_threshold_approval_collector_votes (
1140 proposal_id, token_id, approver_fingerprint,
1141 canonical_token_digest, token_json, received_at
1142 ) VALUES (
1143 'proposal-v1', 'shared-token', 'approver-v1',
1144 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa',
1145 X'01', 101
1146 );
1147 UPDATE chio_store_schema_versions
1148 SET version = 1
1149 WHERE store_key = 'approval';
1150 "#,
1151 )
1152 .unwrap();
1153 drop(connection);
1154
1155 let store = SqliteApprovalStore::open(&path).unwrap();
1156 let connection = store.pool.get().unwrap();
1157 let primary_key_columns: String = connection
1158 .query_row(
1159 r#"
1160 SELECT group_concat(name, ',')
1161 FROM (
1162 SELECT name
1163 FROM pragma_table_info('chio_threshold_approval_collector_votes')
1164 WHERE pk > 0
1165 ORDER BY pk
1166 )
1167 "#,
1168 [],
1169 |row| row.get(0),
1170 )
1171 .unwrap();
1172 let retained_votes: i64 = connection
1173 .query_row(
1174 "SELECT COUNT(*) FROM chio_threshold_approval_collector_votes",
1175 [],
1176 |row| row.get(0),
1177 )
1178 .unwrap();
1179 assert_eq!(primary_key_columns, "proposal_id,token_id");
1180 assert_eq!(retained_votes, 1);
1181 }
1182
1183 #[test]
1184 fn threshold_collector_recovers_votes_and_persists_delivery_before_return() {
1185 let dir = tempfile::tempdir().unwrap();
1186 let path = dir.path().join("threshold.sqlite3");
1187 let authority = Keypair::generate();
1188 let alice = Keypair::generate();
1189 let bob = Keypair::generate();
1190 let subject = Keypair::generate();
1191 let policy_hash = sha256_hex(b"threshold-policy");
1192 let requirement = ThresholdApprovalRequirement::new(
1193 policy_hash.clone(),
1194 2,
1195 vec![
1196 ThresholdApproverIdentity {
1197 identifier: "alice".to_string(),
1198 public_key: alice.public_key(),
1199 },
1200 ThresholdApproverIdentity {
1201 identifier: "bob".to_string(),
1202 public_key: bob.public_key(),
1203 },
1204 ],
1205 "directory-v1".to_string(),
1206 100,
1207 )
1208 .unwrap();
1209 let proposal = ThresholdApprovalProposal::sign(
1210 ThresholdApprovalProposalBody {
1211 schema: THRESHOLD_APPROVAL_PROPOSAL_SCHEMA.to_string(),
1212 proposal_id: "durable-proposal".to_string(),
1213 request_id: "durable-request".to_string(),
1214 governed_intent_hash: sha256_hex(b"durable-intent"),
1215 subject: subject.public_key(),
1216 authorizing_capability_digest: sha256_hex(b"durable-capability"),
1217 policy_hash: policy_hash.clone(),
1218 threshold: 2,
1219 eligible_set_digest: requirement.eligible_set_digest.clone(),
1220 proposal_created_at: 100,
1221 proposal_deadline: 200,
1222 policy_authority: authority.public_key(),
1223 },
1224 &authority,
1225 )
1226 .unwrap();
1227 let second_proposal = ThresholdApprovalProposal::sign(
1228 ThresholdApprovalProposalBody {
1229 schema: THRESHOLD_APPROVAL_PROPOSAL_SCHEMA.to_string(),
1230 proposal_id: "durable-proposal-b".to_string(),
1231 request_id: "durable-request-b".to_string(),
1232 governed_intent_hash: sha256_hex(b"durable-intent-b"),
1233 subject: subject.public_key(),
1234 authorizing_capability_digest: sha256_hex(b"durable-capability-b"),
1235 policy_hash: policy_hash.clone(),
1236 threshold: 2,
1237 eligible_set_digest: requirement.eligible_set_digest.clone(),
1238 proposal_created_at: 100,
1239 proposal_deadline: 200,
1240 policy_authority: authority.public_key(),
1241 },
1242 &authority,
1243 )
1244 .unwrap();
1245 let make_token = |proposal: &ThresholdApprovalProposal,
1246 approver: &Keypair,
1247 id: &str,
1248 expires_at: u64| {
1249 GovernedApprovalToken::sign(
1250 GovernedApprovalTokenBody {
1251 id: id.to_string(),
1252 approver: approver.public_key(),
1253 subject: proposal.body.subject.clone(),
1254 governed_intent_hash: proposal.body.governed_intent_hash.clone(),
1255 request_id: proposal.body.request_id.clone(),
1256 threshold_proposal_hash: Some(proposal.artifact_digest().unwrap()),
1257 issued_at: 101,
1258 expires_at,
1259 decision: GovernedApprovalDecision::Approved,
1260 },
1261 approver,
1262 )
1263 .unwrap()
1264 };
1265
1266 {
1267 let store = Arc::new(SqliteApprovalStore::open(&path).unwrap());
1268 let collector = ThresholdApprovalCollector::new(
1269 store,
1270 policy_hash.clone(),
1271 vec![authority.public_key()],
1272 );
1273 collector
1274 .create_proposal(proposal.clone(), requirement.clone(), None, false, 100)
1275 .unwrap();
1276 }
1277
1278 {
1279 let store = Arc::new(SqliteApprovalStore::open(&path).unwrap());
1280 let collector = ThresholdApprovalCollector::new(
1281 store,
1282 policy_hash.clone(),
1283 vec![authority.public_key()],
1284 );
1285 assert!(collector
1286 .get_proposal("durable-proposal")
1287 .unwrap()
1288 .is_some());
1289 collector
1290 .submit_token(
1291 "durable-proposal",
1292 make_token(&proposal, &alice, "token-alice", 120),
1293 110,
1294 )
1295 .unwrap();
1296 }
1297
1298 {
1299 let store = Arc::new(SqliteApprovalStore::open(&path).unwrap());
1300 let collector = ThresholdApprovalCollector::new(
1301 store,
1302 policy_hash.clone(),
1303 vec![authority.public_key()],
1304 );
1305 let recovered = collector.get_proposal("durable-proposal").unwrap().unwrap();
1306 assert_eq!(recovered.tokens.len(), 1);
1307 let ready = collector
1308 .submit_token(
1309 "durable-proposal",
1310 make_token(&proposal, &bob, "token-bob", 199),
1311 111,
1312 )
1313 .unwrap();
1314 assert_eq!(ready.state, ThresholdApprovalCollectorState::Ready);
1315 }
1316
1317 {
1318 let store = Arc::new(SqliteApprovalStore::open(&path).unwrap());
1319 let collector = ThresholdApprovalCollector::new(
1320 store,
1321 policy_hash.clone(),
1322 vec![authority.public_key()],
1323 );
1324 let ready = collector.get_proposal("durable-proposal").unwrap().unwrap();
1325 assert_eq!(ready.state, ThresholdApprovalCollectorState::Ready);
1326 assert!(collector.deliver("durable-proposal", 150).is_err());
1327 let refreshed = collector
1328 .submit_token(
1329 "durable-proposal",
1330 make_token(&proposal, &alice, "token-alice-fresh", 199),
1331 150,
1332 )
1333 .unwrap();
1334 assert_eq!(refreshed.state, ThresholdApprovalCollectorState::Ready);
1335 assert_eq!(refreshed.tokens.len(), 2);
1336 let delivered = collector.deliver("durable-proposal", 151).unwrap();
1337 assert_eq!(delivered.tokens.len(), 2);
1338 assert!(delivered
1339 .tokens
1340 .iter()
1341 .any(|token| token.id == "token-alice-fresh"));
1342 }
1343
1344 let store = Arc::new(SqliteApprovalStore::open(&path).unwrap());
1345 let collector =
1346 ThresholdApprovalCollector::new(store, policy_hash, vec![authority.public_key()]);
1347 assert_eq!(
1348 collector
1349 .get_proposal("durable-proposal")
1350 .unwrap()
1351 .unwrap()
1352 .state,
1353 ThresholdApprovalCollectorState::Delivered
1354 );
1355 collector
1356 .create_proposal(second_proposal.clone(), requirement, None, false, 100)
1357 .unwrap();
1358 collector
1359 .submit_token(
1360 "durable-proposal-b",
1361 make_token(&second_proposal, &alice, "token-alice-fresh", 199),
1362 110,
1363 )
1364 .unwrap();
1365 let second_ready = collector
1366 .submit_token(
1367 "durable-proposal-b",
1368 make_token(&second_proposal, &bob, "token-bob", 199),
1369 111,
1370 )
1371 .unwrap();
1372 assert_eq!(second_ready.state, ThresholdApprovalCollectorState::Ready);
1373 }
1374}