1#![allow(clippy::too_many_arguments)]
7
8use sqlx::{Row, SqlitePool};
9
10pub use nyx_agent_types::project::{ProjectRecord, ProjectRuntimeProfile};
11
12use crate::store::StoreError;
13
14pub const DEFAULT_PROJECT_ID: &str = "default-project";
19pub const DEFAULT_PROJECT_NAME: &str = "default";
20
21#[derive(Debug, Clone, Default)]
23pub enum ProjectPatchOption<T> {
24 #[default]
25 Unset,
26 Set(T),
27}
28
29#[derive(Debug, Default)]
30pub struct ProjectPatch {
31 pub description: ProjectPatchOption<Option<String>>,
32 pub target_base_url: ProjectPatchOption<Option<String>>,
33 pub env_config_json: ProjectPatchOption<Option<String>>,
34 pub runtime_profile_json: ProjectPatchOption<Option<String>>,
35 pub updated_at: i64,
36}
37
38pub struct ProjectStore<'a> {
39 pool: &'a SqlitePool,
40}
41
42impl<'a> ProjectStore<'a> {
43 pub fn new(pool: &'a SqlitePool) -> Self {
44 Self { pool }
45 }
46
47 pub async fn create(
51 &self,
52 id: &str,
53 name: &str,
54 description: Option<&str>,
55 target_base_url: Option<&str>,
56 env_config_json: Option<&str>,
57 now_ms: i64,
58 ) -> Result<ProjectRecord, StoreError> {
59 self.create_with_runtime_profile(
60 id,
61 name,
62 description,
63 target_base_url,
64 env_config_json,
65 None,
66 now_ms,
67 )
68 .await
69 }
70
71 pub async fn create_with_runtime_profile(
73 &self,
74 id: &str,
75 name: &str,
76 description: Option<&str>,
77 target_base_url: Option<&str>,
78 env_config_json: Option<&str>,
79 runtime_profile_json: Option<&str>,
80 now_ms: i64,
81 ) -> Result<ProjectRecord, StoreError> {
82 sqlx::query(
83 r#"
84 INSERT INTO projects (
85 id, name, description, target_base_url, env_config_json, runtime_profile_json,
86 created_at, updated_at
87 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
88 "#,
89 )
90 .bind(id)
91 .bind(name)
92 .bind(description)
93 .bind(target_base_url)
94 .bind(env_config_json)
95 .bind(runtime_profile_json)
96 .bind(now_ms)
97 .bind(now_ms)
98 .execute(self.pool)
99 .await?;
100 Ok(ProjectRecord {
101 id: id.to_string(),
102 name: name.to_string(),
103 description: description.map(str::to_string),
104 target_base_url: target_base_url.map(str::to_string),
105 env_config_json: env_config_json.map(str::to_string),
106 runtime_profile: parse_runtime_profile_json(runtime_profile_json.map(str::to_string))?,
107 default_launch_profile: None,
108 created_at: now_ms,
109 updated_at: now_ms,
110 })
111 }
112
113 pub async fn ensure_default(&self, now_ms: i64) -> Result<ProjectRecord, StoreError> {
118 if let Some(existing) = self.get(DEFAULT_PROJECT_ID).await? {
119 return Ok(existing);
120 }
121 match self.create(DEFAULT_PROJECT_ID, DEFAULT_PROJECT_NAME, None, None, None, now_ms).await
122 {
123 Ok(rec) => Ok(rec),
124 Err(StoreError::Sqlx(sqlx::Error::Database(db_err)))
126 if db_err.code().as_deref() == Some("2067")
127 || db_err.code().as_deref() == Some("1555") =>
128 {
129 self.get(DEFAULT_PROJECT_ID)
130 .await?
131 .ok_or(StoreError::Sqlx(sqlx::Error::RowNotFound))
132 }
133 Err(other) => Err(other),
134 }
135 }
136
137 pub async fn get(&self, id: &str) -> Result<Option<ProjectRecord>, StoreError> {
138 let row = sqlx::query(
139 r#"
140 SELECT projects.id, projects.name, projects.description, projects.target_base_url,
141 projects.env_config_json, projects.runtime_profile_json,
142 projects.created_at, projects.updated_at,
143 lp.id AS lp_id, lp.project_id AS lp_project_id, lp.name AS lp_name,
144 lp.mode AS lp_mode, lp.build_steps_json AS lp_build_steps_json,
145 lp.start_steps_json AS lp_start_steps_json,
146 lp.seed_steps_json AS lp_seed_steps_json,
147 lp.reset_steps_json AS lp_reset_steps_json,
148 lp.login_steps_json AS lp_login_steps_json,
149 lp.stop_steps_json AS lp_stop_steps_json,
150 lp.health_checks_json AS lp_health_checks_json,
151 lp.target_urls_json AS lp_target_urls_json,
152 lp.env_refs_json AS lp_env_refs_json,
153 lp.working_dirs_json AS lp_working_dirs_json,
154 lp.readiness AS lp_readiness, lp.created_at AS lp_created_at,
155 lp.updated_at AS lp_updated_at, lp.is_default AS lp_is_default
156 FROM projects
157 LEFT JOIN project_launch_profiles lp
158 ON lp.project_id = projects.id AND lp.is_default = 1
159 WHERE projects.id = ?
160 "#,
161 )
162 .bind(id)
163 .fetch_optional(self.pool)
164 .await?;
165 row.map(row_to_project_record).transpose()
166 }
167
168 pub async fn get_by_name(&self, name: &str) -> Result<Option<ProjectRecord>, StoreError> {
169 let row = sqlx::query(
170 r#"
171 SELECT projects.id, projects.name, projects.description, projects.target_base_url,
172 projects.env_config_json, projects.runtime_profile_json,
173 projects.created_at, projects.updated_at,
174 lp.id AS lp_id, lp.project_id AS lp_project_id, lp.name AS lp_name,
175 lp.mode AS lp_mode, lp.build_steps_json AS lp_build_steps_json,
176 lp.start_steps_json AS lp_start_steps_json,
177 lp.seed_steps_json AS lp_seed_steps_json,
178 lp.reset_steps_json AS lp_reset_steps_json,
179 lp.login_steps_json AS lp_login_steps_json,
180 lp.stop_steps_json AS lp_stop_steps_json,
181 lp.health_checks_json AS lp_health_checks_json,
182 lp.target_urls_json AS lp_target_urls_json,
183 lp.env_refs_json AS lp_env_refs_json,
184 lp.working_dirs_json AS lp_working_dirs_json,
185 lp.readiness AS lp_readiness, lp.created_at AS lp_created_at,
186 lp.updated_at AS lp_updated_at, lp.is_default AS lp_is_default
187 FROM projects
188 LEFT JOIN project_launch_profiles lp
189 ON lp.project_id = projects.id AND lp.is_default = 1
190 WHERE projects.name = ?
191 "#,
192 )
193 .bind(name)
194 .fetch_optional(self.pool)
195 .await?;
196 row.map(row_to_project_record).transpose()
197 }
198
199 pub async fn list(&self) -> Result<Vec<ProjectRecord>, StoreError> {
200 let rows = sqlx::query(
201 r#"
202 SELECT projects.id, projects.name, projects.description, projects.target_base_url,
203 projects.env_config_json, projects.runtime_profile_json,
204 projects.created_at, projects.updated_at,
205 lp.id AS lp_id, lp.project_id AS lp_project_id, lp.name AS lp_name,
206 lp.mode AS lp_mode, lp.build_steps_json AS lp_build_steps_json,
207 lp.start_steps_json AS lp_start_steps_json,
208 lp.seed_steps_json AS lp_seed_steps_json,
209 lp.reset_steps_json AS lp_reset_steps_json,
210 lp.login_steps_json AS lp_login_steps_json,
211 lp.stop_steps_json AS lp_stop_steps_json,
212 lp.health_checks_json AS lp_health_checks_json,
213 lp.target_urls_json AS lp_target_urls_json,
214 lp.env_refs_json AS lp_env_refs_json,
215 lp.working_dirs_json AS lp_working_dirs_json,
216 lp.readiness AS lp_readiness, lp.created_at AS lp_created_at,
217 lp.updated_at AS lp_updated_at, lp.is_default AS lp_is_default
218 FROM projects
219 LEFT JOIN project_launch_profiles lp
220 ON lp.project_id = projects.id AND lp.is_default = 1
221 ORDER BY projects.name
222 "#,
223 )
224 .fetch_all(self.pool)
225 .await?;
226 rows.into_iter().map(row_to_project_record).collect()
227 }
228
229 pub async fn update(&self, id: &str, patch: &ProjectPatch) -> Result<bool, StoreError> {
231 let Some(existing) = self.get(id).await? else {
232 return Ok(false);
233 };
234 let description = match &patch.description {
235 ProjectPatchOption::Unset => existing.description,
236 ProjectPatchOption::Set(v) => v.clone(),
237 };
238 let target_base_url = match &patch.target_base_url {
239 ProjectPatchOption::Unset => existing.target_base_url,
240 ProjectPatchOption::Set(v) => v.clone(),
241 };
242 let env_config_json = match &patch.env_config_json {
243 ProjectPatchOption::Unset => existing.env_config_json,
244 ProjectPatchOption::Set(v) => v.clone(),
245 };
246 let runtime_profile_json = match &patch.runtime_profile_json {
247 ProjectPatchOption::Unset => {
248 existing.runtime_profile.as_ref().map(serde_json::to_string).transpose()?
249 }
250 ProjectPatchOption::Set(v) => v.clone(),
251 };
252 sqlx::query(
253 r#"
254 UPDATE projects SET
255 description = ?,
256 target_base_url = ?,
257 env_config_json = ?,
258 runtime_profile_json = ?,
259 updated_at = ?
260 WHERE id = ?
261 "#,
262 )
263 .bind(description)
264 .bind(target_base_url)
265 .bind(env_config_json)
266 .bind(runtime_profile_json)
267 .bind(patch.updated_at)
268 .bind(id)
269 .execute(self.pool)
270 .await?;
271 Ok(true)
272 }
273
274 pub async fn delete(&self, id: &str) -> Result<u64, StoreError> {
275 let res = sqlx::query!("DELETE FROM projects WHERE id = ?", id).execute(self.pool).await?;
276 Ok(res.rows_affected())
277 }
278}
279
280fn row_to_project_record(row: sqlx::sqlite::SqliteRow) -> Result<ProjectRecord, StoreError> {
281 Ok(ProjectRecord {
282 id: row.try_get("id")?,
283 name: row.try_get("name")?,
284 description: row.try_get("description")?,
285 target_base_url: row.try_get("target_base_url")?,
286 env_config_json: row.try_get("env_config_json")?,
287 runtime_profile: parse_runtime_profile_json(row.try_get("runtime_profile_json")?)?,
288 default_launch_profile: row_to_default_launch_profile(&row)?,
289 created_at: row.try_get::<i64, _>("created_at")?,
290 updated_at: row.try_get::<i64, _>("updated_at")?,
291 })
292}
293
294fn row_to_default_launch_profile(
295 row: &sqlx::sqlite::SqliteRow,
296) -> Result<Option<nyx_agent_types::product::ProjectLaunchProfile>, StoreError> {
297 let id: Option<String> = row.try_get("lp_id")?;
298 let Some(id) = id else {
299 return Ok(None);
300 };
301 Ok(Some(nyx_agent_types::product::ProjectLaunchProfile {
302 id,
303 project_id: row.try_get("lp_project_id")?,
304 name: row.try_get("lp_name")?,
305 mode: row.try_get("lp_mode")?,
306 build_steps: serde_json::from_str(&row.try_get::<String, _>("lp_build_steps_json")?)?,
307 start_steps: serde_json::from_str(&row.try_get::<String, _>("lp_start_steps_json")?)?,
308 seed_steps: serde_json::from_str(&row.try_get::<String, _>("lp_seed_steps_json")?)?,
309 reset_steps: serde_json::from_str(&row.try_get::<String, _>("lp_reset_steps_json")?)?,
310 login_steps: serde_json::from_str(&row.try_get::<String, _>("lp_login_steps_json")?)?,
311 stop_steps: serde_json::from_str(&row.try_get::<String, _>("lp_stop_steps_json")?)?,
312 health_checks: serde_json::from_str(&row.try_get::<String, _>("lp_health_checks_json")?)?,
313 target_urls: serde_json::from_str(&row.try_get::<String, _>("lp_target_urls_json")?)?,
314 env_refs: serde_json::from_str(&row.try_get::<String, _>("lp_env_refs_json")?)?,
315 working_dirs: serde_json::from_str(&row.try_get::<String, _>("lp_working_dirs_json")?)?,
316 readiness: row.try_get("lp_readiness")?,
317 created_at: row.try_get::<i64, _>("lp_created_at")?,
318 updated_at: row.try_get::<i64, _>("lp_updated_at")?,
319 is_default: row.try_get::<i64, _>("lp_is_default")? != 0,
320 }))
321}
322
323fn parse_runtime_profile_json(
324 runtime_profile_json: Option<String>,
325) -> Result<Option<ProjectRuntimeProfile>, StoreError> {
326 runtime_profile_json
327 .map(|json| serde_json::from_str(&json))
328 .transpose()
329 .map_err(StoreError::ProjectRuntimeProfileJson)
330}
331
332#[cfg(test)]
333mod tests {
334 use super::*;
335 use crate::store::testutil::fresh_store;
336
337 #[tokio::test]
338 async fn create_then_get_roundtrips() {
339 let (_tmp, s) = fresh_store().await;
340 let rec = s
341 .projects()
342 .create("p-1", "acme", Some("desc"), Some("http://x"), None, 1_000)
343 .await
344 .expect("create");
345 let got = s.projects().get("p-1").await.expect("get").expect("row");
346 assert_eq!(got, rec);
347 }
348
349 #[tokio::test]
350 async fn get_by_name_resolves_id() {
351 let (_tmp, s) = fresh_store().await;
352 s.projects().create("p-1", "acme", None, None, None, 1_000).await.expect("create");
353 let got = s.projects().get_by_name("acme").await.expect("name").expect("row");
354 assert_eq!(got.id, "p-1");
355 }
356
357 #[tokio::test]
358 async fn list_returns_alphabetical_by_name() {
359 let (_tmp, s) = fresh_store().await;
360 s.projects().create("z", "zeta", None, None, None, 1_000).await.expect("z");
362 s.projects().create("a", "alpha", None, None, None, 1_000).await.expect("a");
363 let names: Vec<_> =
364 s.projects().list().await.expect("list").into_iter().map(|r| r.name).collect();
365 assert_eq!(names, vec!["alpha", "default", "zeta"]);
366 }
367
368 #[tokio::test]
369 async fn ensure_default_is_idempotent() {
370 let (_tmp, s) = fresh_store().await;
371 let first = s.projects().get(DEFAULT_PROJECT_ID).await.expect("get").expect("row");
372 let again = s.projects().ensure_default(9_999).await.expect("ensure");
373 assert_eq!(first.id, again.id);
374 assert_eq!(first.created_at, again.created_at, "must not reset created_at");
375 }
376
377 #[tokio::test]
378 async fn update_patches_subset() {
379 let (_tmp, s) = fresh_store().await;
380 s.projects().create("p-1", "acme", None, None, None, 1_000).await.expect("create");
381 let patch = ProjectPatch {
382 description: ProjectPatchOption::Set(Some("now described".to_string())),
383 target_base_url: ProjectPatchOption::Set(Some("http://acme".to_string())),
384 env_config_json: ProjectPatchOption::Unset,
385 runtime_profile_json: ProjectPatchOption::Unset,
386 updated_at: 5_000,
387 };
388 assert!(s.projects().update("p-1", &patch).await.expect("update"));
389 let got = s.projects().get("p-1").await.expect("get").expect("row");
390 assert_eq!(got.description.as_deref(), Some("now described"));
391 assert_eq!(got.target_base_url.as_deref(), Some("http://acme"));
392 assert_eq!(got.env_config_json, None);
393 assert_eq!(got.updated_at, 5_000);
394 }
395
396 #[tokio::test]
397 async fn update_returns_false_when_missing() {
398 let (_tmp, s) = fresh_store().await;
399 let patch = ProjectPatch { updated_at: 1, ..Default::default() };
400 assert!(!s.projects().update("ghost", &patch).await.expect("update"));
401 }
402
403 #[tokio::test]
404 async fn runtime_profile_json_roundtrips() {
405 let (_tmp, s) = fresh_store().await;
406 let profile_json = r#"{
407 "build_commands":[{"command":"npm ci","repo_name":"web","timeout_seconds":120}],
408 "start_commands":[{"command":"npm run dev","working_directory":"frontend"}],
409 "health_check_url":"http://localhost:3000/health",
410 "target_base_url":"http://localhost:3000",
411 "allowed_hosts":["localhost","127.0.0.1"],
412 "env_vars":[{"name":"NODE_ENV","value":"test","secret":false}],
413 "env_file":".env.test",
414 "timeout_seconds":300
415 }"#;
416 s.projects()
417 .create_with_runtime_profile(
418 "p-1",
419 "acme",
420 None,
421 Some("http://localhost:3000"),
422 None,
423 Some(profile_json),
424 1_000,
425 )
426 .await
427 .expect("create");
428
429 let got = s.projects().get("p-1").await.expect("get").expect("row");
430 let profile = got.runtime_profile.expect("profile");
431 assert_eq!(profile.build_commands[0].command, "npm ci");
432 assert_eq!(profile.build_commands[0].repo_name.as_deref(), Some("web"));
433 assert_eq!(profile.start_commands[0].working_directory.as_deref(), Some("frontend"));
434 assert_eq!(profile.allowed_hosts, vec!["localhost", "127.0.0.1"]);
435 assert_eq!(profile.env_vars[0].name, "NODE_ENV");
436 assert_eq!(profile.env_file.as_deref(), Some(".env.test"));
437 assert_eq!(profile.timeout_seconds, Some(300));
438 }
439
440 #[tokio::test]
441 async fn update_can_set_and_clear_runtime_profile() {
442 let (_tmp, s) = fresh_store().await;
443 s.projects().create("p-1", "acme", None, None, None, 1_000).await.expect("create");
444
445 let profile_json = r#"{"start_commands":[{"command":"cargo run"}]}"#;
446 let patch = ProjectPatch {
447 runtime_profile_json: ProjectPatchOption::Set(Some(profile_json.to_string())),
448 updated_at: 2_000,
449 ..Default::default()
450 };
451 assert!(s.projects().update("p-1", &patch).await.expect("set"));
452 let got = s.projects().get("p-1").await.expect("get").expect("row");
453 assert_eq!(got.runtime_profile.expect("profile").start_commands[0].command, "cargo run");
454
455 let clear = ProjectPatch {
456 runtime_profile_json: ProjectPatchOption::Set(None),
457 updated_at: 3_000,
458 ..Default::default()
459 };
460 assert!(s.projects().update("p-1", &clear).await.expect("clear"));
461 let got = s.projects().get("p-1").await.expect("get").expect("row");
462 assert!(got.runtime_profile.is_none());
463 }
464
465 #[tokio::test]
466 async fn delete_cascades_to_repos() {
467 use crate::store::testutil::sample_repo_for_project;
468 let (_tmp, s) = fresh_store().await;
469 let p = s
470 .projects()
471 .create("p-doomed", "doomed", None, None, None, 1_000)
472 .await
473 .expect("create");
474 let r = sample_repo_for_project("attached", &p.id);
475 s.repos().upsert(&r).await.expect("upsert");
476 let affected = s.projects().delete("p-doomed").await.expect("delete");
477 assert_eq!(affected, 1);
478 assert!(
479 s.repos().get("attached").await.expect("get").is_none(),
480 "FK cascade must drop repo"
481 );
482 }
483}