agentmux_srv\backend\storage/
content.rs1use rusqlite::params;
15use serde::{Deserialize, Serialize};
16
17use super::error::StoreError;
18use super::store::Store;
19
20#[derive(Debug, Clone, Serialize, Deserialize)]
22pub struct AgentContent {
23 pub agent_id: String,
24 pub content_type: String,
25 pub content: String,
26 pub updated_at: i64,
27}
28
29impl Store {
30 pub fn agent_content_get(
31 &self,
32 agent_id: &str,
33 content_type: &str,
34 ) -> Result<Option<AgentContent>, StoreError> {
35 let conn = self.conn.lock().unwrap();
36 let mut stmt = conn.prepare(
37 "SELECT agent_id, content_type, content, updated_at
38 FROM db_agent_content WHERE agent_id=?1 AND content_type=?2",
39 )?;
40 let local = {
41 let result = stmt.query_row(params![agent_id, content_type], |row| {
42 Ok(AgentContent {
43 agent_id: row.get(0)?,
44 content_type: row.get(1)?,
45 content: row.get(2)?,
46 updated_at: row.get(3)?,
47 })
48 });
49 match result {
50 Ok(content) => Some(content),
51 Err(rusqlite::Error::QueryReturnedNoRows) => None,
52 Err(e) => return Err(StoreError::Sqlite(e)),
53 }
54 };
55 drop(stmt);
56 drop(conn);
57 if local.is_some() {
58 return Ok(local);
59 }
60 if self.agent_def_exists_local(agent_id)? {
65 return Ok(local);
66 }
67 if let Some(reg) = self.shared_def_registry() {
68 if let Ok(Some(rec)) = reg.get(agent_id) {
69 if let Some(c) = rec
70 .data
71 .content
72 .iter()
73 .find(|c| c.content_type == content_type)
74 {
75 return Ok(Some(AgentContent {
76 agent_id: agent_id.to_string(),
77 content_type: c.content_type.clone(),
78 content: c.content.clone(),
79 updated_at: rec.data.updated_at,
80 }));
81 }
82 }
83 }
84 Ok(local)
85 }
86
87 pub fn agent_content_set(&self, content: &AgentContent) -> Result<(), StoreError> {
89 {
90 let conn = self.conn.lock().unwrap();
91 conn.execute(
92 "INSERT INTO db_agent_content (agent_id, content_type, content, updated_at)
93 VALUES (?1, ?2, ?3, ?4)
94 ON CONFLICT(agent_id, content_type) DO UPDATE SET content=?3, updated_at=?4",
95 params![
96 content.agent_id,
97 content.content_type,
98 content.content,
99 content.updated_at,
100 ],
101 )?;
102 }
103 self.registry_def_upsert(&content.agent_id);
107 Ok(())
108 }
109
110 pub(super) fn agent_content_get_all_local(
115 &self,
116 agent_id: &str,
117 ) -> Result<Vec<AgentContent>, StoreError> {
118 let conn = self.conn.lock().unwrap();
119 let mut stmt = conn.prepare(
120 "SELECT agent_id, content_type, content, updated_at
121 FROM db_agent_content WHERE agent_id=?1 ORDER BY content_type ASC",
122 )?;
123 let rows = stmt.query_map(params![agent_id], |row| {
124 Ok(AgentContent {
125 agent_id: row.get(0)?,
126 content_type: row.get(1)?,
127 content: row.get(2)?,
128 updated_at: row.get(3)?,
129 })
130 })?;
131 let mut contents = Vec::new();
132 for row in rows {
133 contents.push(row?);
134 }
135 Ok(contents)
136 }
137
138 pub fn agent_content_get_all(
139 &self,
140 agent_id: &str,
141 ) -> Result<Vec<AgentContent>, StoreError> {
142 let local = self.agent_content_get_all_local(agent_id)?;
143 if !local.is_empty() {
144 return Ok(local);
145 }
146 if self.agent_def_exists_local(agent_id)? {
152 return Ok(local);
153 }
154 if let Some(reg) = self.shared_def_registry() {
155 if let Ok(Some(rec)) = reg.get(agent_id) {
156 return Ok(rec
157 .data
158 .content
159 .iter()
160 .map(|c| AgentContent {
161 agent_id: agent_id.to_string(),
162 content_type: c.content_type.clone(),
163 content: c.content.clone(),
164 updated_at: rec.data.updated_at,
165 })
166 .collect());
167 }
168 }
169 Ok(local)
170 }
171
172 pub fn agent_content_delete(
174 &self,
175 agent_id: &str,
176 content_type: &str,
177 ) -> Result<bool, StoreError> {
178 let rows = {
179 let conn = self.conn.lock().unwrap();
180 conn.execute(
181 "DELETE FROM db_agent_content WHERE agent_id=?1 AND content_type=?2",
182 params![agent_id, content_type],
183 )?
184 };
185 self.registry_def_upsert(agent_id);
192 Ok(rows > 0)
193 }
194}