Skip to main content

stateset_db/sqlite/
work_orders.rs

1//! SQLite Work Order repository implementation
2
3use super::{
4    build_in_clause, map_db_error, params_refs, parse_datetime, parse_datetime_opt,
5    parse_datetime_opt_row, parse_datetime_row, parse_decimal_opt_row, parse_decimal_row,
6    parse_decimal_strict, parse_enum, parse_enum_row, parse_uuid, parse_uuid_opt,
7    parse_uuid_opt_row, parse_uuid_row, uuid_params, with_immediate_transaction,
8};
9use chrono::Utc;
10use r2d2::Pool;
11use r2d2_sqlite::SqliteConnectionManager;
12use rust_decimal::Decimal;
13use stateset_core::{
14    AddWorkOrderMaterial, BatchResult, CommerceError, CreateWorkOrder, CreateWorkOrderTask,
15    ProductId, Result, TaskStatus, UpdateWorkOrder, UpdateWorkOrderTask, WorkOrder,
16    WorkOrderFilter, WorkOrderMaterial, WorkOrderPriority, WorkOrderRepository, WorkOrderStatus,
17    WorkOrderTask, validate_batch_size,
18};
19use uuid::Uuid;
20
21/// SQLite implementation of `WorkOrderRepository`
22#[derive(Debug)]
23pub struct SqliteWorkOrderRepository {
24    pool: Pool<SqliteConnectionManager>,
25}
26
27impl SqliteWorkOrderRepository {
28    #[must_use]
29    pub const fn new(pool: Pool<SqliteConnectionManager>) -> Self {
30        Self { pool }
31    }
32
33    fn load_tasks(&self, work_order_id: Uuid) -> Result<Vec<WorkOrderTask>> {
34        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
35
36        let mut stmt = conn
37            .prepare(
38                "SELECT id, work_order_id, sequence, task_name, status, estimated_hours,
39                        actual_hours, assigned_to, started_at, completed_at, notes, created_at, updated_at
40                 FROM manufacturing_work_order_tasks WHERE work_order_id = ? ORDER BY sequence",
41            )
42            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
43
44        let rows = stmt
45            .query_map([work_order_id.to_string()], |row| {
46                Ok(WorkOrderTask {
47                    id: parse_uuid_row(&row.get::<_, String>(0)?, "work_order_task", "id")?,
48                    work_order_id: parse_uuid_row(
49                        &row.get::<_, String>(1)?,
50                        "work_order_task",
51                        "work_order_id",
52                    )?,
53                    sequence: row.get(2)?,
54                    task_name: row.get(3)?,
55                    status: parse_enum_row(&row.get::<_, String>(4)?, "work_order_task", "status")?,
56                    estimated_hours: parse_decimal_opt_row(
57                        row.get(5)?,
58                        "work_order_task",
59                        "estimated_hours",
60                    )?,
61                    actual_hours: parse_decimal_opt_row(
62                        row.get(6)?,
63                        "work_order_task",
64                        "actual_hours",
65                    )?,
66                    assigned_to: parse_uuid_opt_row(
67                        row.get::<_, Option<String>>(7)?,
68                        "work_order_task",
69                        "assigned_to",
70                    )?,
71                    started_at: parse_datetime_opt_row(
72                        row.get(8)?,
73                        "work_order_task",
74                        "started_at",
75                    )?,
76                    completed_at: parse_datetime_opt_row(
77                        row.get(9)?,
78                        "work_order_task",
79                        "completed_at",
80                    )?,
81                    notes: row.get(10)?,
82                    created_at: parse_datetime_row(
83                        &row.get::<_, String>(11)?,
84                        "work_order_task",
85                        "created_at",
86                    )?,
87                    updated_at: parse_datetime_row(
88                        &row.get::<_, String>(12)?,
89                        "work_order_task",
90                        "updated_at",
91                    )?,
92                })
93            })
94            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
95
96        let mut tasks = Vec::new();
97        for row in rows {
98            tasks.push(row.map_err(|e| CommerceError::DatabaseError(e.to_string()))?);
99        }
100
101        Ok(tasks)
102    }
103
104    fn load_materials(&self, work_order_id: Uuid) -> Result<Vec<WorkOrderMaterial>> {
105        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
106
107        let mut stmt = conn
108            .prepare(
109                "SELECT id, work_order_id, component_id, component_sku, component_name,
110                        reserved_quantity, consumed_quantity, inventory_reservation_id, created_at, updated_at
111                 FROM manufacturing_work_order_materials WHERE work_order_id = ?",
112            )
113            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
114
115        let rows = stmt
116            .query_map([work_order_id.to_string()], |row| {
117                Ok(WorkOrderMaterial {
118                    id: parse_uuid_row(&row.get::<_, String>(0)?, "work_order_material", "id")?,
119                    work_order_id: parse_uuid_row(
120                        &row.get::<_, String>(1)?,
121                        "work_order_material",
122                        "work_order_id",
123                    )?,
124                    component_id: parse_uuid_opt_row(
125                        row.get::<_, Option<String>>(2)?,
126                        "work_order_material",
127                        "component_id",
128                    )?,
129                    component_sku: row.get(3)?,
130                    component_name: row.get(4)?,
131                    reserved_quantity: parse_decimal_row(
132                        &row.get::<_, String>(5)?,
133                        "work_order_material",
134                        "reserved_quantity",
135                    )?,
136                    consumed_quantity: parse_decimal_row(
137                        &row.get::<_, String>(6)?,
138                        "work_order_material",
139                        "consumed_quantity",
140                    )?,
141                    inventory_reservation_id: parse_uuid_opt_row(
142                        row.get::<_, Option<String>>(7)?,
143                        "work_order_material",
144                        "inventory_reservation_id",
145                    )?,
146                    created_at: parse_datetime_row(
147                        &row.get::<_, String>(8)?,
148                        "work_order_material",
149                        "created_at",
150                    )?,
151                    updated_at: parse_datetime_row(
152                        &row.get::<_, String>(9)?,
153                        "work_order_material",
154                        "updated_at",
155                    )?,
156                })
157            })
158            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
159
160        let mut materials = Vec::new();
161        for row in rows {
162            materials.push(row.map_err(|e| CommerceError::DatabaseError(e.to_string()))?);
163        }
164
165        Ok(materials)
166    }
167
168    fn get_task_internal(&self, task_id: Uuid) -> Result<WorkOrderTask> {
169        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
170
171        conn.query_row(
172            "SELECT id, work_order_id, sequence, task_name, status, estimated_hours,
173                    actual_hours, assigned_to, started_at, completed_at, notes, created_at, updated_at
174             FROM manufacturing_work_order_tasks WHERE id = ?",
175            [task_id.to_string()],
176            |row| {
177                Ok(WorkOrderTask {
178                    id: parse_uuid_row(&row.get::<_, String>(0)?, "work_order_task", "id")?,
179                    work_order_id: parse_uuid_row(&row.get::<_, String>(1)?, "work_order_task", "work_order_id")?,
180                    sequence: row.get(2)?,
181                    task_name: row.get(3)?,
182                    status: parse_enum_row(&row.get::<_, String>(4)?, "work_order_task", "status")?,
183                    estimated_hours: parse_decimal_opt_row(row.get(5)?, "work_order_task", "estimated_hours")?,
184                    actual_hours: parse_decimal_opt_row(row.get(6)?, "work_order_task", "actual_hours")?,
185                    assigned_to: parse_uuid_opt_row(
186                        row.get::<_, Option<String>>(7)?,
187                        "work_order_task",
188                        "assigned_to",
189                    )?,
190                    started_at: parse_datetime_opt_row(row.get(8)?, "work_order_task", "started_at")?,
191                    completed_at: parse_datetime_opt_row(row.get(9)?, "work_order_task", "completed_at")?,
192                    notes: row.get(10)?,
193                    created_at: parse_datetime_row(&row.get::<_, String>(11)?, "work_order_task", "created_at")?,
194                    updated_at: parse_datetime_row(&row.get::<_, String>(12)?, "work_order_task", "updated_at")?,
195                })
196            },
197        )
198        .map_err(|e| match e {
199            rusqlite::Error::QueryReturnedNoRows => CommerceError::NotFound,
200            _ => CommerceError::DatabaseError(e.to_string()),
201        })
202    }
203
204    fn get_material_internal(&self, material_id: Uuid) -> Result<WorkOrderMaterial> {
205        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
206
207        conn.query_row(
208            "SELECT id, work_order_id, component_id, component_sku, component_name,
209                    reserved_quantity, consumed_quantity, inventory_reservation_id, created_at, updated_at
210             FROM manufacturing_work_order_materials WHERE id = ?",
211            [material_id.to_string()],
212            |row| {
213                Ok(WorkOrderMaterial {
214                    id: parse_uuid_row(&row.get::<_, String>(0)?, "work_order_material", "id")?,
215                    work_order_id: parse_uuid_row(&row.get::<_, String>(1)?, "work_order_material", "work_order_id")?,
216                    component_id: parse_uuid_opt_row(
217                        row.get::<_, Option<String>>(2)?,
218                        "work_order_material",
219                        "component_id",
220                    )?,
221                    component_sku: row.get(3)?,
222                    component_name: row.get(4)?,
223                    reserved_quantity: parse_decimal_row(&row.get::<_, String>(5)?, "work_order_material", "reserved_quantity")?,
224                    consumed_quantity: parse_decimal_row(&row.get::<_, String>(6)?, "work_order_material", "consumed_quantity")?,
225                    inventory_reservation_id: parse_uuid_opt_row(
226                        row.get::<_, Option<String>>(7)?,
227                        "work_order_material",
228                        "inventory_reservation_id",
229                    )?,
230                    created_at: parse_datetime_row(&row.get::<_, String>(8)?, "work_order_material", "created_at")?,
231                    updated_at: parse_datetime_row(&row.get::<_, String>(9)?, "work_order_material", "updated_at")?,
232                })
233            },
234        )
235        .map_err(|e| match e {
236            rusqlite::Error::QueryReturnedNoRows => CommerceError::NotFound,
237            _ => CommerceError::DatabaseError(e.to_string()),
238        })
239    }
240}
241
242impl WorkOrderRepository for SqliteWorkOrderRepository {
243    fn create(&self, input: CreateWorkOrder) -> Result<WorkOrder> {
244        let id = Uuid::new_v4();
245        let work_order_number = WorkOrder::generate_work_order_number();
246        let now = Utc::now();
247        let priority = input.priority.unwrap_or(WorkOrderPriority::Normal);
248
249        // Insert work order in a scoped block to release connection before adding tasks
250        {
251            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
252
253            conn.execute(
254                "INSERT INTO manufacturing_work_orders (id, work_order_number, product_id, bom_id, work_center_id,
255                 assigned_to, status, priority, quantity_to_build, quantity_completed, scheduled_start, scheduled_end, notes, created_at, updated_at)
256                 VALUES (?, ?, ?, ?, ?, ?, 'planned', ?, ?, '0', ?, ?, ?, ?, ?)",
257                rusqlite::params![
258                    id.to_string(),
259                    work_order_number,
260                    input.product_id.to_string(),
261                    input.bom_id.map(|u| u.to_string()),
262                    input.work_center_id,
263                    input.assigned_to.map(|u| u.to_string()),
264                    priority.to_string(),
265                    input.quantity_to_build.to_string(),
266                    input.scheduled_start.map(|dt| dt.to_rfc3339()),
267                    input.scheduled_end.map(|dt| dt.to_rfc3339()),
268                    input.notes,
269                    now.to_rfc3339(),
270                    now.to_rfc3339(),
271                ],
272            )
273            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
274        } // Connection released here
275
276        // Create tasks if provided (after releasing the main connection)
277        let mut tasks = Vec::new();
278        if let Some(task_inputs) = input.tasks {
279            for task_input in task_inputs {
280                let task = self.add_task(id, task_input)?;
281                tasks.push(task);
282            }
283        }
284
285        Ok(WorkOrder {
286            id,
287            work_order_number,
288            product_id: input.product_id,
289            bom_id: input.bom_id,
290            work_center_id: input.work_center_id,
291            assigned_to: input.assigned_to,
292            status: WorkOrderStatus::Planned,
293            priority,
294            quantity_to_build: input.quantity_to_build,
295            quantity_completed: Decimal::ZERO,
296            scheduled_start: input.scheduled_start,
297            scheduled_end: input.scheduled_end,
298            actual_start: None,
299            actual_end: None,
300            notes: input.notes,
301            tasks,
302            materials: vec![],
303            version: 1,
304            created_at: now,
305            updated_at: now,
306        })
307    }
308
309    fn get(&self, id: Uuid) -> Result<Option<WorkOrder>> {
310        // Query work order in a scoped block to release connection before loading tasks/materials
311        let wo_data = {
312            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
313
314            let result = conn.query_row(
315                "SELECT id, work_order_number, product_id, bom_id, work_center_id, assigned_to,
316                        status, priority, quantity_to_build, quantity_completed, scheduled_start,
317                        scheduled_end, actual_start, actual_end, notes, created_at, updated_at
318                 FROM manufacturing_work_orders WHERE id = ?",
319                [id.to_string()],
320                |row| {
321                    Ok((
322                        row.get::<_, String>(0)?,
323                        row.get::<_, String>(1)?,
324                        row.get::<_, String>(2)?,
325                        row.get::<_, Option<String>>(3)?,
326                        row.get::<_, Option<String>>(4)?,
327                        row.get::<_, Option<String>>(5)?,
328                        row.get::<_, String>(6)?,
329                        row.get::<_, String>(7)?,
330                        row.get::<_, String>(8)?,
331                        row.get::<_, String>(9)?,
332                        row.get::<_, Option<String>>(10)?,
333                        row.get::<_, Option<String>>(11)?,
334                        row.get::<_, Option<String>>(12)?,
335                        row.get::<_, Option<String>>(13)?,
336                        row.get::<_, Option<String>>(14)?,
337                        row.get::<_, String>(15)?,
338                        row.get::<_, String>(16)?,
339                    ))
340                },
341            );
342
343            match result {
344                Ok(data) => Some(data),
345                Err(rusqlite::Error::QueryReturnedNoRows) => None,
346                Err(e) => return Err(CommerceError::DatabaseError(e.to_string())),
347            }
348        }; // Connection released here
349
350        match wo_data {
351            Some((
352                id_str,
353                work_order_number,
354                product_id,
355                bom_id,
356                work_center_id,
357                assigned_to,
358                status,
359                priority,
360                quantity_to_build,
361                quantity_completed,
362                scheduled_start,
363                scheduled_end,
364                actual_start,
365                actual_end,
366                notes,
367                created_at,
368                updated_at,
369            )) => {
370                let wo_id = parse_uuid(&id_str, "work_order", "id")?;
371                let tasks = self.load_tasks(wo_id)?;
372                let materials = self.load_materials(wo_id)?;
373
374                Ok(Some(WorkOrder {
375                    id: wo_id,
376                    work_order_number,
377                    product_id: ProductId::from(parse_uuid(
378                        &product_id,
379                        "work_order",
380                        "product_id",
381                    )?),
382                    bom_id: parse_uuid_opt(bom_id, "work_order", "bom_id")?,
383                    work_center_id,
384                    assigned_to: parse_uuid_opt(assigned_to, "work_order", "assigned_to")?,
385                    status: parse_enum(&status, "work_order", "status")?,
386                    priority: parse_enum(&priority, "work_order", "priority")?,
387                    quantity_to_build: parse_decimal_strict(
388                        &quantity_to_build,
389                        "work_order",
390                        "quantity_to_build",
391                    )?,
392                    quantity_completed: parse_decimal_strict(
393                        &quantity_completed,
394                        "work_order",
395                        "quantity_completed",
396                    )?,
397                    scheduled_start: parse_datetime_opt(
398                        scheduled_start,
399                        "work_order",
400                        "scheduled_start",
401                    )?,
402                    scheduled_end: parse_datetime_opt(
403                        scheduled_end,
404                        "work_order",
405                        "scheduled_end",
406                    )?,
407                    actual_start: parse_datetime_opt(actual_start, "work_order", "actual_start")?,
408                    actual_end: parse_datetime_opt(actual_end, "work_order", "actual_end")?,
409                    notes,
410                    tasks,
411                    materials,
412                    version: 1, // Default to 1 for backwards compatibility
413                    created_at: parse_datetime(&created_at, "work_order", "created_at")?,
414                    updated_at: parse_datetime(&updated_at, "work_order", "updated_at")?,
415                }))
416            }
417            None => Ok(None),
418        }
419    }
420
421    fn get_by_number(&self, work_order_number: &str) -> Result<Option<WorkOrder>> {
422        // Query ID in a scoped block to release connection before calling self.get()
423        let id_result = {
424            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
425
426            let result = conn.query_row(
427                "SELECT id FROM manufacturing_work_orders WHERE work_order_number = ?",
428                [work_order_number],
429                |row| row.get::<_, String>(0),
430            );
431
432            match result {
433                Ok(id_str) => Some(parse_uuid(&id_str, "work_order", "id")?),
434                Err(rusqlite::Error::QueryReturnedNoRows) => None,
435                Err(e) => return Err(CommerceError::DatabaseError(e.to_string())),
436            }
437        }; // Connection released here
438
439        match id_result {
440            Some(id) => self.get(id),
441            None => Ok(None),
442        }
443    }
444
445    fn update(&self, id: Uuid, input: UpdateWorkOrder) -> Result<WorkOrder> {
446        // Get existing work order first (releases connection after)
447        let existing = self.get(id)?.ok_or(CommerceError::NotFound)?;
448        let now = Utc::now();
449
450        let new_status = input.status.unwrap_or(existing.status);
451        let new_priority = input.priority.unwrap_or(existing.priority);
452        let new_assigned_to = input.assigned_to.or(existing.assigned_to);
453        let new_notes = input.notes.or(existing.notes);
454        let new_work_center_id = input.work_center_id.or(existing.work_center_id);
455        let new_scheduled_start = input.scheduled_start.or(existing.scheduled_start);
456        let new_scheduled_end = input.scheduled_end.or(existing.scheduled_end);
457
458        // Do the update in a scoped block
459        {
460            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
461
462            conn.execute(
463                "UPDATE manufacturing_work_orders SET status = ?, priority = ?, assigned_to = ?,
464                 work_center_id = ?, scheduled_start = ?, scheduled_end = ?, notes = ?, updated_at = ? WHERE id = ?",
465                rusqlite::params![
466                    new_status.to_string(),
467                    new_priority.to_string(),
468                    new_assigned_to.map(|u| u.to_string()),
469                    new_work_center_id,
470                    new_scheduled_start.map(|dt| dt.to_rfc3339()),
471                    new_scheduled_end.map(|dt| dt.to_rfc3339()),
472                    new_notes,
473                    now.to_rfc3339(),
474                    id.to_string(),
475                ],
476            )
477            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
478        } // Connection released here
479
480        // Fetch and return the updated work order
481        self.get(id)?.ok_or(CommerceError::NotFound)
482    }
483
484    fn list(&self, filter: WorkOrderFilter) -> Result<Vec<WorkOrder>> {
485        // Collect all IDs in a scoped block to release connection before calling self.get()
486        let ids = {
487            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
488
489            let limit = i64::from(super::effective_limit(filter.limit));
490            let offset = i64::from(filter.offset.unwrap_or(0));
491
492            let mut sql = "SELECT id FROM manufacturing_work_orders WHERE 1=1".to_string();
493            let mut params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
494
495            if let Some(product_id) = filter.product_id {
496                sql.push_str(" AND product_id = ?");
497                params.push(Box::new(product_id.to_string()));
498            }
499
500            if let Some(bom_id) = filter.bom_id {
501                sql.push_str(" AND bom_id = ?");
502                params.push(Box::new(bom_id.to_string()));
503            }
504
505            if let Some(status) = filter.status {
506                sql.push_str(" AND status = ?");
507                params.push(Box::new(status.to_string()));
508            }
509
510            if let Some(priority) = filter.priority {
511                sql.push_str(" AND priority = ?");
512                params.push(Box::new(priority.to_string()));
513            }
514
515            if let Some(assigned_to) = filter.assigned_to {
516                sql.push_str(" AND assigned_to = ?");
517                params.push(Box::new(assigned_to.to_string()));
518            }
519
520            if let Some(work_center_id) = filter.work_center_id {
521                sql.push_str(" AND work_center_id = ?");
522                params.push(Box::new(work_center_id));
523            }
524            if filter.overdue_only.unwrap_or(false) {
525                sql.push_str(" AND scheduled_end IS NOT NULL AND scheduled_end < ? AND status NOT IN ('completed', 'cancelled')");
526                params.push(Box::new(Utc::now().to_rfc3339()));
527            }
528
529            // Keyset cursor: (created_at, id) for stable DESC ordering
530            if let Some((cursor_created, cursor_id)) = &filter.after_cursor {
531                sql.push_str(" AND (created_at < ? OR (created_at = ? AND id < ?))");
532                params.push(Box::new(cursor_created.clone()));
533                params.push(Box::new(cursor_created.clone()));
534                params.push(Box::new(cursor_id.clone()));
535            }
536
537            sql.push_str(" ORDER BY created_at DESC, id DESC LIMIT ? OFFSET ?");
538            params.push(Box::new(limit));
539            // Offset pagination applies only in non-cursor mode.
540            params.push(Box::new(if filter.after_cursor.is_none() { offset } else { 0 }));
541
542            let mut stmt =
543                conn.prepare(&sql).map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
544
545            let param_refs: Vec<&dyn rusqlite::ToSql> =
546                params.iter().map(std::convert::AsRef::as_ref).collect();
547
548            let rows = stmt
549                .query_map(param_refs.as_slice(), |row| row.get::<_, String>(0))
550                .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
551
552            let mut id_list = Vec::new();
553            for row in rows {
554                let id_str = row.map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
555                id_list.push(parse_uuid(&id_str, "work_order", "id")?);
556            }
557            id_list
558        }; // Connection released here
559
560        // Now fetch each work order (each call gets its own connection)
561        let mut work_orders = Vec::new();
562        for id in ids {
563            if let Some(wo) = self.get(id)? {
564                work_orders.push(wo);
565            }
566        }
567
568        Ok(work_orders)
569    }
570
571    fn delete(&self, id: Uuid) -> Result<()> {
572        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
573
574        // Mark as cancelled instead of hard delete
575        conn.execute(
576            "UPDATE manufacturing_work_orders SET status = 'cancelled', updated_at = ? WHERE id = ?",
577            rusqlite::params![Utc::now().to_rfc3339(), id.to_string()],
578        )
579        .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
580
581        Ok(())
582    }
583
584    fn start(&self, id: Uuid) -> Result<WorkOrder> {
585        let now = Utc::now();
586
587        // Do the update in a scoped block
588        {
589            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
590
591            conn.execute(
592                "UPDATE manufacturing_work_orders SET status = 'in_progress', actual_start = ?, updated_at = ? WHERE id = ?",
593                rusqlite::params![now.to_rfc3339(), now.to_rfc3339(), id.to_string()],
594            )
595            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
596        } // Connection released here
597
598        // Fetch and return the updated work order
599        self.get(id)?.ok_or(CommerceError::NotFound)
600    }
601
602    fn complete(&self, id: Uuid, quantity_completed: Decimal) -> Result<WorkOrder> {
603        if quantity_completed <= Decimal::ZERO {
604            return Err(CommerceError::ValidationError(
605                "Completed quantity must be greater than zero".to_string(),
606            ));
607        }
608
609        let now = Utc::now();
610        let id_str = id.to_string();
611
612        // Read the current `quantity_completed`, accumulate, and write it back
613        // inside ONE `IMMEDIATE` transaction. IMMEDIATE takes the write lock up
614        // front, so two concurrent `complete` calls serialize instead of both
615        // reading the same starting quantity and one overwriting the other (a
616        // lost update that would under-count completed units).
617        with_immediate_transaction(&self.pool, |tx| {
618            let existing: (String, String, Option<String>) = tx
619                .query_row(
620                    "SELECT quantity_completed, quantity_to_build, actual_end FROM manufacturing_work_orders WHERE id = ?",
621                    [&id_str],
622                    |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
623                )
624                .map_err(|e| match e {
625                    rusqlite::Error::QueryReturnedNoRows => {
626                        rusqlite::Error::ToSqlConversionFailure(Box::new(CommerceError::NotFound))
627                    }
628                    other => other,
629                })?;
630
631            let existing_completed =
632                parse_decimal_strict(&existing.0, "work_order", "quantity_completed")
633                    .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
634            let quantity_to_build =
635                parse_decimal_strict(&existing.1, "work_order", "quantity_to_build")
636                    .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
637            let new_quantity_completed = existing_completed + quantity_completed;
638            let is_complete = new_quantity_completed >= quantity_to_build;
639            let new_status = if is_complete { "completed" } else { "partially_completed" };
640            let new_actual_end = if is_complete { Some(now.to_rfc3339()) } else { existing.2 };
641
642            tx.execute(
643                "UPDATE manufacturing_work_orders
644                 SET quantity_completed = ?, status = ?, actual_end = ?, updated_at = ?
645                 WHERE id = ?",
646                rusqlite::params![
647                    new_quantity_completed.to_string(),
648                    new_status,
649                    new_actual_end,
650                    now.to_rfc3339(),
651                    id_str,
652                ],
653            )?;
654            Ok(())
655        })?;
656
657        // Fetch and return the updated work order
658        self.get(id)?.ok_or(CommerceError::NotFound)
659    }
660
661    fn hold(&self, id: Uuid) -> Result<WorkOrder> {
662        self.update(
663            id,
664            UpdateWorkOrder { status: Some(WorkOrderStatus::OnHold), ..Default::default() },
665        )
666    }
667
668    fn resume(&self, id: Uuid) -> Result<WorkOrder> {
669        self.update(
670            id,
671            UpdateWorkOrder { status: Some(WorkOrderStatus::InProgress), ..Default::default() },
672        )
673    }
674
675    fn cancel(&self, id: Uuid) -> Result<WorkOrder> {
676        self.update(
677            id,
678            UpdateWorkOrder { status: Some(WorkOrderStatus::Cancelled), ..Default::default() },
679        )
680    }
681
682    fn add_task(&self, work_order_id: Uuid, task: CreateWorkOrderTask) -> Result<WorkOrderTask> {
683        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
684
685        let id = Uuid::new_v4();
686        let now = Utc::now();
687        let sequence = task.sequence.unwrap_or(1);
688
689        conn.execute(
690            "INSERT INTO manufacturing_work_order_tasks (id, work_order_id, sequence, task_name, status, estimated_hours, assigned_to, notes, created_at, updated_at)
691             VALUES (?, ?, ?, ?, 'pending', ?, ?, ?, ?, ?)",
692            rusqlite::params![
693                id.to_string(),
694                work_order_id.to_string(),
695                sequence,
696                task.task_name,
697                task.estimated_hours.map(|h| h.to_string()),
698                task.assigned_to.map(|u| u.to_string()),
699                task.notes,
700                now.to_rfc3339(),
701                now.to_rfc3339(),
702            ],
703        )
704        .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
705
706        Ok(WorkOrderTask {
707            id,
708            work_order_id,
709            sequence,
710            task_name: task.task_name,
711            status: TaskStatus::Pending,
712            estimated_hours: task.estimated_hours,
713            actual_hours: None,
714            assigned_to: task.assigned_to,
715            started_at: None,
716            completed_at: None,
717            notes: task.notes,
718            created_at: now,
719            updated_at: now,
720        })
721    }
722
723    fn update_task(&self, task_id: Uuid, task: UpdateWorkOrderTask) -> Result<WorkOrderTask> {
724        // Get existing task first (releases connection after)
725        let existing = self.get_task_internal(task_id)?;
726        let now = Utc::now();
727
728        let new_sequence = task.sequence.unwrap_or(existing.sequence);
729        let new_task_name = task.task_name.unwrap_or(existing.task_name);
730        let new_status = task.status.unwrap_or(existing.status);
731        let new_estimated = task.estimated_hours.or(existing.estimated_hours);
732        let new_actual = task.actual_hours.or(existing.actual_hours);
733        let new_assigned = task.assigned_to.or(existing.assigned_to);
734        let new_started = task.started_at.or(existing.started_at);
735        let new_completed = task.completed_at.or(existing.completed_at);
736        let new_notes = task.notes.or(existing.notes);
737
738        // Do the update in a scoped block
739        {
740            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
741
742            conn.execute(
743                "UPDATE manufacturing_work_order_tasks SET sequence = ?, task_name = ?, status = ?,
744                 estimated_hours = ?, actual_hours = ?, assigned_to = ?, started_at = ?,
745                 completed_at = ?, notes = ?, updated_at = ? WHERE id = ?",
746                rusqlite::params![
747                    new_sequence,
748                    new_task_name,
749                    new_status.to_string(),
750                    new_estimated.map(|h| h.to_string()),
751                    new_actual.map(|h| h.to_string()),
752                    new_assigned.map(|u| u.to_string()),
753                    new_started.map(|dt| dt.to_rfc3339()),
754                    new_completed.map(|dt| dt.to_rfc3339()),
755                    new_notes,
756                    now.to_rfc3339(),
757                    task_id.to_string(),
758                ],
759            )
760            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
761        } // Connection released here
762
763        self.get_task_internal(task_id)
764    }
765
766    fn remove_task(&self, task_id: Uuid) -> Result<()> {
767        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
768
769        conn.execute(
770            "DELETE FROM manufacturing_work_order_tasks WHERE id = ?",
771            [task_id.to_string()],
772        )
773        .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
774
775        Ok(())
776    }
777
778    fn get_tasks(&self, work_order_id: Uuid) -> Result<Vec<WorkOrderTask>> {
779        self.load_tasks(work_order_id)
780    }
781
782    fn start_task(&self, task_id: Uuid) -> Result<WorkOrderTask> {
783        let now = Utc::now();
784
785        // Do the update in a scoped block
786        {
787            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
788
789            conn.execute(
790                "UPDATE manufacturing_work_order_tasks SET status = 'in_progress', started_at = ?, updated_at = ? WHERE id = ?",
791                rusqlite::params![now.to_rfc3339(), now.to_rfc3339(), task_id.to_string()],
792            )
793            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
794        } // Connection released here
795
796        self.get_task_internal(task_id)
797    }
798
799    fn complete_task(&self, task_id: Uuid, actual_hours: Option<Decimal>) -> Result<WorkOrderTask> {
800        let now = Utc::now();
801
802        // Do the update in a scoped block
803        {
804            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
805
806            conn.execute(
807                "UPDATE manufacturing_work_order_tasks SET status = 'completed', actual_hours = ?, completed_at = ?, updated_at = ? WHERE id = ?",
808                rusqlite::params![
809                    actual_hours.map(|h| h.to_string()),
810                    now.to_rfc3339(),
811                    now.to_rfc3339(),
812                    task_id.to_string(),
813                ],
814            )
815            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
816        } // Connection released here
817
818        self.get_task_internal(task_id)
819    }
820
821    fn add_material(
822        &self,
823        work_order_id: Uuid,
824        material: AddWorkOrderMaterial,
825    ) -> Result<WorkOrderMaterial> {
826        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
827
828        let id = Uuid::new_v4();
829        let now = Utc::now();
830
831        conn.execute(
832            "INSERT INTO manufacturing_work_order_materials (id, work_order_id, component_id, component_sku, component_name, reserved_quantity, consumed_quantity, created_at, updated_at)
833             VALUES (?, ?, ?, ?, ?, ?, '0', ?, ?)",
834            rusqlite::params![
835                id.to_string(),
836                work_order_id.to_string(),
837                material.component_id.map(|u| u.to_string()),
838                material.component_sku,
839                material.component_name,
840                material.quantity.to_string(),
841                now.to_rfc3339(),
842                now.to_rfc3339(),
843            ],
844        )
845        .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
846
847        Ok(WorkOrderMaterial {
848            id,
849            work_order_id,
850            component_id: material.component_id,
851            component_sku: material.component_sku,
852            component_name: material.component_name,
853            reserved_quantity: material.quantity,
854            consumed_quantity: Decimal::ZERO,
855            inventory_reservation_id: None,
856            created_at: now,
857            updated_at: now,
858        })
859    }
860
861    fn consume_material(&self, material_id: Uuid, quantity: Decimal) -> Result<WorkOrderMaterial> {
862        // Get existing material first (releases connection after)
863        let existing = self.get_material_internal(material_id)?;
864        let now = Utc::now();
865
866        let new_consumed = existing.consumed_quantity + quantity;
867
868        // Do the update in a scoped block
869        {
870            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
871
872            conn.execute(
873                "UPDATE manufacturing_work_order_materials SET consumed_quantity = ?, updated_at = ? WHERE id = ?",
874                rusqlite::params![
875                    new_consumed.to_string(),
876                    now.to_rfc3339(),
877                    material_id.to_string(),
878                ],
879            )
880            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
881        } // Connection released here
882
883        self.get_material_internal(material_id)
884    }
885
886    fn get_materials(&self, work_order_id: Uuid) -> Result<Vec<WorkOrderMaterial>> {
887        self.load_materials(work_order_id)
888    }
889
890    fn count(&self, filter: WorkOrderFilter) -> Result<u64> {
891        let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
892
893        let mut sql = "SELECT COUNT(*) FROM manufacturing_work_orders WHERE 1=1".to_string();
894        let mut params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
895
896        if let Some(product_id) = filter.product_id {
897            sql.push_str(" AND product_id = ?");
898            params.push(Box::new(product_id.to_string()));
899        }
900
901        if let Some(bom_id) = filter.bom_id {
902            sql.push_str(" AND bom_id = ?");
903            params.push(Box::new(bom_id.to_string()));
904        }
905
906        if let Some(status) = filter.status {
907            sql.push_str(" AND status = ?");
908            params.push(Box::new(status.to_string()));
909        }
910
911        if let Some(priority) = filter.priority {
912            sql.push_str(" AND priority = ?");
913            params.push(Box::new(priority.to_string()));
914        }
915
916        if let Some(assigned_to) = filter.assigned_to {
917            sql.push_str(" AND assigned_to = ?");
918            params.push(Box::new(assigned_to.to_string()));
919        }
920
921        if let Some(work_center_id) = filter.work_center_id {
922            sql.push_str(" AND work_center_id = ?");
923            params.push(Box::new(work_center_id));
924        }
925        if filter.overdue_only.unwrap_or(false) {
926            sql.push_str(" AND scheduled_end IS NOT NULL AND scheduled_end < ? AND status NOT IN ('completed', 'cancelled')");
927            params.push(Box::new(Utc::now().to_rfc3339()));
928        }
929
930        let param_refs: Vec<&dyn rusqlite::ToSql> =
931            params.iter().map(std::convert::AsRef::as_ref).collect();
932
933        let count: i64 = conn
934            .query_row(&sql, param_refs.as_slice(), |row| row.get(0))
935            .map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
936
937        Ok(count as u64)
938    }
939
940    // === Batch Operations ===
941
942    fn create_batch(&self, inputs: Vec<CreateWorkOrder>) -> Result<BatchResult<WorkOrder>> {
943        validate_batch_size(&inputs)?;
944        let mut result = BatchResult::with_capacity(inputs.len());
945
946        for (index, input) in inputs.into_iter().enumerate() {
947            match self.create(input) {
948                Ok(work_order) => result.record_success(work_order),
949                Err(e) => result.record_failure(index, None, &e),
950            }
951        }
952
953        Ok(result)
954    }
955
956    fn create_batch_atomic(&self, inputs: Vec<CreateWorkOrder>) -> Result<Vec<WorkOrder>> {
957        validate_batch_size(&inputs)?;
958        if inputs.is_empty() {
959            return Ok(vec![]);
960        }
961
962        let mut conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
963        let tx = super::begin_immediate(&mut conn).map_err(map_db_error)?;
964        let mut results = Vec::with_capacity(inputs.len());
965
966        for input in inputs {
967            let id = Uuid::new_v4();
968            let work_order_number = WorkOrder::generate_work_order_number();
969            let now = Utc::now();
970            let priority = input.priority.unwrap_or(WorkOrderPriority::Normal);
971
972            tx.execute(
973                "INSERT INTO manufacturing_work_orders (id, work_order_number, product_id, bom_id, work_center_id,
974                 assigned_to, status, priority, quantity_to_build, quantity_completed, scheduled_start, scheduled_end, notes, created_at, updated_at)
975                 VALUES (?, ?, ?, ?, ?, ?, 'planned', ?, ?, '0', ?, ?, ?, ?, ?)",
976                rusqlite::params![
977                    id.to_string(),
978                    work_order_number,
979                    input.product_id.to_string(),
980                    input.bom_id.map(|u| u.to_string()),
981                    input.work_center_id,
982                    input.assigned_to.map(|u| u.to_string()),
983                    priority.to_string(),
984                    input.quantity_to_build.to_string(),
985                    input.scheduled_start.map(|dt| dt.to_rfc3339()),
986                    input.scheduled_end.map(|dt| dt.to_rfc3339()),
987                    input.notes,
988                    now.to_rfc3339(),
989                    now.to_rfc3339(),
990                ],
991            )
992            .map_err(map_db_error)?;
993
994            // Insert tasks if provided
995            let mut tasks = Vec::new();
996            if let Some(task_inputs) = input.tasks {
997                for task_input in task_inputs {
998                    let task_id = Uuid::new_v4();
999                    let sequence = task_input.sequence.unwrap_or(1);
1000
1001                    tx.execute(
1002                        "INSERT INTO manufacturing_work_order_tasks (id, work_order_id, sequence, task_name, status, estimated_hours, assigned_to, notes, created_at, updated_at)
1003                         VALUES (?, ?, ?, ?, 'pending', ?, ?, ?, ?, ?)",
1004                        rusqlite::params![
1005                            task_id.to_string(),
1006                            id.to_string(),
1007                            sequence,
1008                            task_input.task_name,
1009                            task_input.estimated_hours.map(|h| h.to_string()),
1010                            task_input.assigned_to.map(|u| u.to_string()),
1011                            task_input.notes,
1012                            now.to_rfc3339(),
1013                            now.to_rfc3339(),
1014                        ],
1015                    )
1016                    .map_err(map_db_error)?;
1017
1018                    tasks.push(WorkOrderTask {
1019                        id: task_id,
1020                        work_order_id: id,
1021                        sequence,
1022                        task_name: task_input.task_name,
1023                        status: TaskStatus::Pending,
1024                        estimated_hours: task_input.estimated_hours,
1025                        actual_hours: None,
1026                        assigned_to: task_input.assigned_to,
1027                        started_at: None,
1028                        completed_at: None,
1029                        notes: task_input.notes,
1030                        created_at: now,
1031                        updated_at: now,
1032                    });
1033                }
1034            }
1035
1036            results.push(WorkOrder {
1037                id,
1038                work_order_number,
1039                product_id: input.product_id,
1040                bom_id: input.bom_id,
1041                work_center_id: input.work_center_id,
1042                assigned_to: input.assigned_to,
1043                status: WorkOrderStatus::Planned,
1044                priority,
1045                quantity_to_build: input.quantity_to_build,
1046                quantity_completed: Decimal::ZERO,
1047                scheduled_start: input.scheduled_start,
1048                scheduled_end: input.scheduled_end,
1049                actual_start: None,
1050                actual_end: None,
1051                notes: input.notes,
1052                tasks,
1053                materials: vec![],
1054                version: 1,
1055                created_at: now,
1056                updated_at: now,
1057            });
1058        }
1059
1060        tx.commit().map_err(map_db_error)?;
1061        Ok(results)
1062    }
1063
1064    fn update_batch(
1065        &self,
1066        updates: Vec<(Uuid, UpdateWorkOrder)>,
1067    ) -> Result<BatchResult<WorkOrder>> {
1068        validate_batch_size(&updates)?;
1069        let mut result = BatchResult::with_capacity(updates.len());
1070
1071        for (index, (id, input)) in updates.into_iter().enumerate() {
1072            match self.update(id, input) {
1073                Ok(work_order) => result.record_success(work_order),
1074                Err(e) => result.record_failure(index, Some(id.to_string()), &e),
1075            }
1076        }
1077
1078        Ok(result)
1079    }
1080
1081    fn update_batch_atomic(&self, updates: Vec<(Uuid, UpdateWorkOrder)>) -> Result<Vec<WorkOrder>> {
1082        validate_batch_size(&updates)?;
1083        if updates.is_empty() {
1084            return Ok(vec![]);
1085        }
1086
1087        let mut conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
1088        let tx = super::begin_immediate(&mut conn).map_err(map_db_error)?;
1089        let mut updated_ids = Vec::with_capacity(updates.len());
1090
1091        for (id, input) in updates {
1092            let now = Utc::now();
1093
1094            // Get existing work order data
1095            let existing: (String, String, Option<String>, Option<String>, Option<String>, Option<String>) = tx
1096                .query_row(
1097                    "SELECT status, priority, assigned_to, work_center_id, scheduled_start, scheduled_end
1098                     FROM manufacturing_work_orders WHERE id = ?",
1099                    [id.to_string()],
1100                    |row| Ok((
1101                        row.get(0)?,
1102                        row.get(1)?,
1103                        row.get(2)?,
1104                        row.get(3)?,
1105                        row.get(4)?,
1106                        row.get(5)?,
1107                    )),
1108                )
1109                .map_err(|e| match e {
1110                    rusqlite::Error::QueryReturnedNoRows => CommerceError::NotFound,
1111                    _ => map_db_error(e),
1112                })?;
1113
1114            let new_status = input.status.map(|s| s.to_string()).unwrap_or(existing.0);
1115            let new_priority = input.priority.map(|p| p.to_string()).unwrap_or(existing.1);
1116            let new_assigned_to =
1117                input.assigned_to.map(|u| Some(u.to_string())).unwrap_or(existing.2);
1118            let new_work_center_id = input.work_center_id.or(existing.3);
1119            let new_scheduled_start =
1120                input.scheduled_start.map(|dt| Some(dt.to_rfc3339())).unwrap_or(existing.4);
1121            let new_scheduled_end =
1122                input.scheduled_end.map(|dt| Some(dt.to_rfc3339())).unwrap_or(existing.5);
1123            let new_notes = input.notes;
1124
1125            tx.execute(
1126                "UPDATE manufacturing_work_orders SET status = ?, priority = ?, assigned_to = ?,
1127                 work_center_id = ?, scheduled_start = ?, scheduled_end = ?, notes = COALESCE(?, notes), updated_at = ? WHERE id = ?",
1128                rusqlite::params![
1129                    new_status,
1130                    new_priority,
1131                    new_assigned_to,
1132                    new_work_center_id,
1133                    new_scheduled_start,
1134                    new_scheduled_end,
1135                    new_notes,
1136                    now.to_rfc3339(),
1137                    id.to_string(),
1138                ],
1139            )
1140            .map_err(map_db_error)?;
1141
1142            updated_ids.push(id);
1143        }
1144
1145        tx.commit().map_err(map_db_error)?;
1146
1147        // Fetch all updated work orders
1148        let mut results = Vec::with_capacity(updated_ids.len());
1149        for id in updated_ids {
1150            if let Some(wo) = self.get(id)? {
1151                results.push(wo);
1152            }
1153        }
1154
1155        Ok(results)
1156    }
1157
1158    fn delete_batch(&self, ids: Vec<Uuid>) -> Result<BatchResult<Uuid>> {
1159        validate_batch_size(&ids)?;
1160        let mut result = BatchResult::with_capacity(ids.len());
1161
1162        for (index, id) in ids.into_iter().enumerate() {
1163            match self.delete(id) {
1164                Ok(()) => result.record_success(id),
1165                Err(e) => result.record_failure(index, Some(id.to_string()), &e),
1166            }
1167        }
1168
1169        Ok(result)
1170    }
1171
1172    fn delete_batch_atomic(&self, ids: Vec<Uuid>) -> Result<()> {
1173        validate_batch_size(&ids)?;
1174        if ids.is_empty() {
1175            return Ok(());
1176        }
1177
1178        let mut conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
1179        let tx = super::begin_immediate(&mut conn).map_err(map_db_error)?;
1180
1181        let placeholders = build_in_clause(ids.len());
1182
1183        // Soft-delete semantics must match single-item delete.
1184        let now = Utc::now().to_rfc3339();
1185        let mut all_params: Vec<Box<dyn rusqlite::ToSql>> = vec![Box::new(now)];
1186        all_params.extend(uuid_params(&ids));
1187        let all_params_refs = params_refs(&all_params);
1188
1189        let sql = format!(
1190            "UPDATE manufacturing_work_orders SET status = 'cancelled', updated_at = ? WHERE id IN ({placeholders})"
1191        );
1192        tx.execute(&sql, all_params_refs.as_slice()).map_err(map_db_error)?;
1193
1194        tx.commit().map_err(map_db_error)?;
1195        Ok(())
1196    }
1197
1198    fn get_batch(&self, ids: Vec<Uuid>) -> Result<Vec<WorkOrder>> {
1199        validate_batch_size(&ids)?;
1200        if ids.is_empty() {
1201            return Ok(vec![]);
1202        }
1203
1204        // Query all work orders in a single query
1205        let work_order_data = {
1206            let conn = self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))?;
1207
1208            let placeholders = build_in_clause(ids.len());
1209            let sql = format!(
1210                "SELECT id, work_order_number, product_id, bom_id, work_center_id, assigned_to,
1211                        status, priority, quantity_to_build, quantity_completed, scheduled_start,
1212                        scheduled_end, actual_start, actual_end, notes, created_at, updated_at
1213                 FROM manufacturing_work_orders WHERE id IN ({placeholders})"
1214            );
1215
1216            let params = uuid_params(&ids);
1217            let param_refs = params_refs(&params);
1218
1219            let mut stmt = conn.prepare(&sql).map_err(map_db_error)?;
1220
1221            let rows = stmt
1222                .query_map(param_refs.as_slice(), |row| {
1223                    Ok((
1224                        row.get::<_, String>(0)?,
1225                        row.get::<_, String>(1)?,
1226                        row.get::<_, String>(2)?,
1227                        row.get::<_, Option<String>>(3)?,
1228                        row.get::<_, Option<String>>(4)?,
1229                        row.get::<_, Option<String>>(5)?,
1230                        row.get::<_, String>(6)?,
1231                        row.get::<_, String>(7)?,
1232                        row.get::<_, String>(8)?,
1233                        row.get::<_, String>(9)?,
1234                        row.get::<_, Option<String>>(10)?,
1235                        row.get::<_, Option<String>>(11)?,
1236                        row.get::<_, Option<String>>(12)?,
1237                        row.get::<_, Option<String>>(13)?,
1238                        row.get::<_, Option<String>>(14)?,
1239                        row.get::<_, String>(15)?,
1240                        row.get::<_, String>(16)?,
1241                    ))
1242                })
1243                .map_err(map_db_error)?;
1244
1245            let mut data = Vec::new();
1246            for row in rows {
1247                data.push(row.map_err(map_db_error)?);
1248            }
1249            data
1250        };
1251
1252        // Build work orders, loading tasks and materials for each
1253        let mut work_orders = Vec::with_capacity(work_order_data.len());
1254        for (
1255            id_str,
1256            work_order_number,
1257            product_id,
1258            bom_id,
1259            work_center_id,
1260            assigned_to,
1261            status,
1262            priority,
1263            quantity_to_build,
1264            quantity_completed,
1265            scheduled_start,
1266            scheduled_end,
1267            actual_start,
1268            actual_end,
1269            notes,
1270            created_at,
1271            updated_at,
1272        ) in work_order_data
1273        {
1274            let wo_id = parse_uuid(&id_str, "work_order", "id")?;
1275            let tasks = self.load_tasks(wo_id)?;
1276            let materials = self.load_materials(wo_id)?;
1277
1278            work_orders.push(WorkOrder {
1279                id: wo_id,
1280                work_order_number,
1281                product_id: ProductId::from(parse_uuid(&product_id, "work_order", "product_id")?),
1282                bom_id: parse_uuid_opt(bom_id, "work_order", "bom_id")?,
1283                work_center_id,
1284                assigned_to: parse_uuid_opt(assigned_to, "work_order", "assigned_to")?,
1285                status: parse_enum(&status, "work_order", "status")?,
1286                priority: parse_enum(&priority, "work_order", "priority")?,
1287                quantity_to_build: parse_decimal_strict(
1288                    &quantity_to_build,
1289                    "work_order",
1290                    "quantity_to_build",
1291                )?,
1292                quantity_completed: parse_decimal_strict(
1293                    &quantity_completed,
1294                    "work_order",
1295                    "quantity_completed",
1296                )?,
1297                scheduled_start: parse_datetime_opt(
1298                    scheduled_start,
1299                    "work_order",
1300                    "scheduled_start",
1301                )?,
1302                scheduled_end: parse_datetime_opt(scheduled_end, "work_order", "scheduled_end")?,
1303                actual_start: parse_datetime_opt(actual_start, "work_order", "actual_start")?,
1304                actual_end: parse_datetime_opt(actual_end, "work_order", "actual_end")?,
1305                notes,
1306                tasks,
1307                materials,
1308                version: 1,
1309                created_at: parse_datetime(&created_at, "work_order", "created_at")?,
1310                updated_at: parse_datetime(&updated_at, "work_order", "updated_at")?,
1311            });
1312        }
1313
1314        Ok(work_orders)
1315    }
1316}
1317
1318#[cfg(test)]
1319mod tests {
1320    use super::*;
1321    use crate::SqliteDatabase;
1322    use rust_decimal_macros::dec;
1323    use stateset_core::{
1324        AddWorkOrderMaterial, CreateWorkOrder, CreateWorkOrderTask, ProductId, TaskStatus,
1325        WorkOrderFilter, WorkOrderRepository, WorkOrderStatus,
1326    };
1327
1328    fn fresh_repo() -> SqliteWorkOrderRepository {
1329        SqliteDatabase::in_memory().expect("in-memory").work_orders()
1330    }
1331
1332    fn make_wo(repo: &SqliteWorkOrderRepository, qty: Decimal) -> WorkOrder {
1333        repo.create(CreateWorkOrder {
1334            product_id: ProductId::new(),
1335            quantity_to_build: qty,
1336            ..Default::default()
1337        })
1338        .expect("create wo")
1339    }
1340
1341    #[test]
1342    fn create_wo_starts_in_planned() {
1343        let repo = fresh_repo();
1344        let wo = make_wo(&repo, dec!(10));
1345        assert_eq!(wo.status, WorkOrderStatus::Planned);
1346        assert_eq!(wo.quantity_to_build, dec!(10));
1347        assert!(!wo.work_order_number.is_empty());
1348    }
1349
1350    #[test]
1351    fn create_wo_with_tasks_persists_them() {
1352        let repo = fresh_repo();
1353        let wo = repo
1354            .create(CreateWorkOrder {
1355                product_id: ProductId::new(),
1356                quantity_to_build: dec!(5),
1357                tasks: Some(vec![
1358                    CreateWorkOrderTask {
1359                        sequence: Some(1),
1360                        task_name: "Assembly".into(),
1361                        estimated_hours: Some(dec!(2)),
1362                        assigned_to: None,
1363                        notes: None,
1364                    },
1365                    CreateWorkOrderTask {
1366                        sequence: Some(2),
1367                        task_name: "QA".into(),
1368                        estimated_hours: Some(dec!(1)),
1369                        assigned_to: None,
1370                        notes: None,
1371                    },
1372                ]),
1373                ..Default::default()
1374            })
1375            .expect("create");
1376        let tasks = repo.get_tasks(wo.id).expect("tasks");
1377        assert_eq!(tasks.len(), 2);
1378    }
1379
1380    #[test]
1381    fn get_and_get_by_number_round_trip() {
1382        let repo = fresh_repo();
1383        let wo = make_wo(&repo, dec!(1));
1384        let by_id = repo.get(wo.id).expect("ok").expect("found");
1385        assert_eq!(by_id.id, wo.id);
1386        let by_num = repo.get_by_number(&wo.work_order_number).expect("ok").expect("found");
1387        assert_eq!(by_num.id, wo.id);
1388        assert!(repo.get_by_number("missing").expect("ok").is_none());
1389    }
1390
1391    #[test]
1392    fn start_transitions_to_in_progress() {
1393        let repo = fresh_repo();
1394        let wo = make_wo(&repo, dec!(1));
1395        let started = repo.start(wo.id).expect("start");
1396        assert_eq!(started.status, WorkOrderStatus::InProgress);
1397    }
1398
1399    #[test]
1400    fn complete_full_quantity_marks_completed() {
1401        let repo = fresh_repo();
1402        let wo = make_wo(&repo, dec!(5));
1403        repo.start(wo.id).expect("start");
1404        let done = repo.complete(wo.id, dec!(5)).expect("complete");
1405        assert_eq!(done.status, WorkOrderStatus::Completed);
1406        assert_eq!(done.quantity_completed, dec!(5));
1407    }
1408
1409    #[test]
1410    fn complete_partial_quantity_marks_partially_completed() {
1411        let repo = fresh_repo();
1412        let wo = make_wo(&repo, dec!(10));
1413        repo.start(wo.id).expect("start");
1414        let partial = repo.complete(wo.id, dec!(7)).expect("partial");
1415        assert_eq!(partial.status, WorkOrderStatus::PartiallyCompleted);
1416        assert_eq!(partial.quantity_completed, dec!(7));
1417    }
1418
1419    #[test]
1420    fn cancel_transitions_to_cancelled() {
1421        let repo = fresh_repo();
1422        let wo = make_wo(&repo, dec!(1));
1423        let cancelled = repo.cancel(wo.id).expect("cancel");
1424        assert_eq!(cancelled.status, WorkOrderStatus::Cancelled);
1425    }
1426
1427    #[test]
1428    fn list_filters_by_status() {
1429        let repo = fresh_repo();
1430        let planned = make_wo(&repo, dec!(1));
1431        let in_progress = make_wo(&repo, dec!(1));
1432        repo.start(in_progress.id).expect("start");
1433
1434        let pl = repo
1435            .list(WorkOrderFilter { status: Some(WorkOrderStatus::Planned), ..Default::default() })
1436            .expect("planned");
1437        let ip = repo
1438            .list(WorkOrderFilter {
1439                status: Some(WorkOrderStatus::InProgress),
1440                ..Default::default()
1441            })
1442            .expect("in_progress");
1443        assert!(pl.iter().any(|w| w.id == planned.id));
1444        assert!(ip.iter().any(|w| w.id == in_progress.id));
1445    }
1446
1447    #[test]
1448    fn add_task_appends_to_work_order() {
1449        let repo = fresh_repo();
1450        let wo = make_wo(&repo, dec!(1));
1451        let task = repo
1452            .add_task(
1453                wo.id,
1454                CreateWorkOrderTask {
1455                    sequence: Some(1),
1456                    task_name: "Solder".into(),
1457                    estimated_hours: Some(dec!(0.5)),
1458                    assigned_to: None,
1459                    notes: None,
1460                },
1461            )
1462            .expect("add task");
1463        assert_eq!(task.task_name, "Solder");
1464        let tasks = repo.get_tasks(wo.id).expect("tasks");
1465        assert_eq!(tasks.len(), 1);
1466    }
1467
1468    #[test]
1469    fn start_and_complete_task_transitions() {
1470        let repo = fresh_repo();
1471        let wo = make_wo(&repo, dec!(1));
1472        let task = repo
1473            .add_task(
1474                wo.id,
1475                CreateWorkOrderTask {
1476                    sequence: Some(1),
1477                    task_name: "Wash".into(),
1478                    estimated_hours: Some(dec!(1)),
1479                    assigned_to: None,
1480                    notes: None,
1481                },
1482            )
1483            .expect("add");
1484        let started = repo.start_task(task.id).expect("start");
1485        assert_eq!(started.status, TaskStatus::InProgress);
1486        let completed = repo.complete_task(task.id, Some(dec!(1.5))).expect("complete");
1487        assert_eq!(completed.status, TaskStatus::Completed);
1488    }
1489
1490    #[test]
1491    fn add_and_consume_material() {
1492        let repo = fresh_repo();
1493        let wo = make_wo(&repo, dec!(1));
1494        let mat = repo
1495            .add_material(
1496                wo.id,
1497                AddWorkOrderMaterial {
1498                    component_id: None,
1499                    component_sku: "PART-A".into(),
1500                    component_name: "Resistor".into(),
1501                    quantity: dec!(10),
1502                },
1503            )
1504            .expect("add mat");
1505        assert_eq!(mat.component_sku, "PART-A");
1506
1507        let consumed = repo.consume_material(mat.id, dec!(4)).expect("consume");
1508        assert_eq!(consumed.id, mat.id);
1509        let remaining_materials = repo.get_materials(wo.id).expect("materials");
1510        assert_eq!(remaining_materials.len(), 1);
1511    }
1512
1513    #[test]
1514    fn create_batch_returns_per_input_results() {
1515        let repo = fresh_repo();
1516        let result = repo
1517            .create_batch(vec![
1518                CreateWorkOrder {
1519                    product_id: ProductId::new(),
1520                    quantity_to_build: dec!(1),
1521                    ..Default::default()
1522                },
1523                CreateWorkOrder {
1524                    product_id: ProductId::new(),
1525                    quantity_to_build: dec!(2),
1526                    ..Default::default()
1527                },
1528            ])
1529            .expect("batch");
1530        assert_eq!(result.success_count, 2);
1531        assert_eq!(result.failure_count, 0);
1532    }
1533
1534    #[test]
1535    fn get_unknown_id_returns_none() {
1536        let repo = fresh_repo();
1537        assert!(repo.get(Uuid::new_v4()).expect("ok").is_none());
1538    }
1539
1540    #[test]
1541    fn get_batch_returns_only_existing() {
1542        let repo = fresh_repo();
1543        let w1 = make_wo(&repo, dec!(1));
1544        let w2 = make_wo(&repo, dec!(2));
1545        let stranger = Uuid::new_v4();
1546        let fetched = repo.get_batch(vec![w1.id, w2.id, stranger]).expect("ok");
1547        assert_eq!(fetched.len(), 2);
1548    }
1549
1550    #[test]
1551    fn concurrent_completions_are_not_lost() {
1552        // Ten completions of one unit each land simultaneously on a work order.
1553        // `complete` is a read-modify-write of `quantity_completed`; without
1554        // serialization the reads race and completions are silently overwritten
1555        // (lost updates). Every committed completion must be counted.
1556        use std::sync::{Arc, Barrier};
1557        use std::thread;
1558
1559        let db = Arc::new(SqliteDatabase::in_memory().expect("in-memory"));
1560        let wo = db
1561            .work_orders()
1562            .create(CreateWorkOrder {
1563                product_id: ProductId::new(),
1564                quantity_to_build: dec!(1000),
1565                ..Default::default()
1566            })
1567            .expect("create wo");
1568
1569        let thread_count = 10usize;
1570        let barrier = Arc::new(Barrier::new(thread_count));
1571        let mut handles = Vec::new();
1572        for _ in 0..thread_count {
1573            let db = Arc::clone(&db);
1574            let barrier = Arc::clone(&barrier);
1575            let id = wo.id;
1576            handles.push(thread::spawn(move || {
1577                let repo = db.work_orders();
1578                barrier.wait();
1579                repo.complete(id, dec!(1))
1580            }));
1581        }
1582        let results: Vec<_> = handles.into_iter().map(|h| h.join().expect("thread")).collect();
1583        // A failure is acceptable only if it is a transient lock (the caller
1584        // retries); a lost update is never acceptable.
1585        assert!(
1586            results
1587                .iter()
1588                .all(|r| r.is_ok() || format!("{:?}", r.as_ref().unwrap_err()).contains("locked")),
1589            "unexpected non-lock failure: {results:?}"
1590        );
1591
1592        let fetched = db.work_orders().get(wo.id).expect("get").expect("found");
1593        // Each completion adds exactly one unit; with the read+write serialized
1594        // every one lands, so the total equals the thread count (no lost updates).
1595        assert_eq!(
1596            fetched.quantity_completed,
1597            Decimal::from(thread_count as u64),
1598            "completions were lost to a race: {results:?}"
1599        );
1600    }
1601
1602    #[test]
1603    fn list_after_cursor_paginates_without_overlap() {
1604        let repo = fresh_repo();
1605        for _ in 0..3 {
1606            make_wo(&repo, dec!(1));
1607        }
1608        let all = repo.list(WorkOrderFilter::default()).expect("list all");
1609        assert_eq!(all.len(), 3);
1610
1611        let first_page =
1612            repo.list(WorkOrderFilter { limit: Some(2), ..Default::default() }).expect("page 1");
1613        assert_eq!(first_page.len(), 2);
1614        assert_eq!(first_page[0].id, all[0].id);
1615
1616        let last = &first_page[1];
1617        let second_page = repo
1618            .list(WorkOrderFilter {
1619                after_cursor: Some((last.created_at.to_rfc3339(), last.id.to_string())),
1620                ..Default::default()
1621            })
1622            .expect("page 2");
1623        assert_eq!(second_page.len(), 1);
1624        assert_eq!(second_page[0].id, all[2].id);
1625    }
1626}