Skip to main content

stateset_db/sqlite/
fulfillment.rs

1//! SQLite implementation for fulfillment (pick/pack/ship) management
2
3use crate::sqlite::{
4    map_db_error, parse_datetime_opt_row, parse_datetime_row, parse_decimal_opt_row,
5    parse_decimal_row, parse_enum_row, parse_uuid, parse_uuid_opt_row, parse_uuid_row,
6    with_immediate_transaction,
7};
8use chrono::Utc;
9use r2d2::Pool;
10use r2d2_sqlite::SqliteConnectionManager;
11use rusqlite::params;
12use rust_decimal::Decimal;
13use uuid::Uuid;
14
15use stateset_core::{
16    AddCarton, AddCartonItem, BatchResult, Carton, CartonItem, CommerceError, CompletePick,
17    CompleteShip, CreatePackTask, CreatePickTask, CreateShipTask, CreateWave, FulfillmentId,
18    FulfillmentRepository, OrderId, OrderItemId, PackStatus, PackTask, PackTaskFilter, PickStatus,
19    PickTask, PickTaskFilter, Result, ShipStatus, ShipTask, ShipTaskFilter, ShipmentId, Wave,
20    WaveFilter, WaveStatus, generate_carton_number, generate_wave_number,
21};
22
23/// SQLite fulfillment repository
24#[derive(Debug)]
25pub struct SqliteFulfillmentRepository {
26    pool: Pool<SqliteConnectionManager>,
27}
28
29impl SqliteFulfillmentRepository {
30    #[must_use]
31    pub const fn new(pool: Pool<SqliteConnectionManager>) -> Self {
32        Self { pool }
33    }
34
35    fn conn(&self) -> Result<r2d2::PooledConnection<SqliteConnectionManager>> {
36        self.pool.get().map_err(|e| CommerceError::DatabaseError(e.to_string()))
37    }
38
39    fn row_to_wave(row: &rusqlite::Row<'_>) -> rusqlite::Result<Wave> {
40        let id_str: String = row.get("id")?;
41        let status_str: String = row.get("status")?;
42        let started_str: Option<String> = row.get("started_at")?;
43        let completed_str: Option<String> = row.get("completed_at")?;
44
45        Ok(Wave {
46            id: FulfillmentId::from(parse_uuid_row(&id_str, "wave", "id")?),
47            wave_number: row.get("wave_number")?,
48            warehouse_id: row.get("warehouse_id")?,
49            status: parse_enum_row(&status_str, "wave", "status")?,
50            order_count: row.get("order_count")?,
51            pick_count: row.get("pick_count")?,
52            completed_pick_count: row.get("completed_pick_count")?,
53            priority: row.get("priority")?,
54            started_at: parse_datetime_opt_row(started_str, "wave", "started_at")?,
55            completed_at: parse_datetime_opt_row(completed_str, "wave", "completed_at")?,
56            notes: row.get("notes")?,
57            created_by: row.get("created_by")?,
58            created_at: parse_datetime_row(
59                &row.get::<_, String>("created_at")?,
60                "wave",
61                "created_at",
62            )?,
63            updated_at: parse_datetime_row(
64                &row.get::<_, String>("updated_at")?,
65                "wave",
66                "updated_at",
67            )?,
68        })
69    }
70
71    fn row_to_pick(row: &rusqlite::Row<'_>) -> rusqlite::Result<PickTask> {
72        let id_str: String = row.get("id")?;
73        let wave_id_str: Option<String> = row.get("wave_id")?;
74        let order_id_str: String = row.get("order_id")?;
75        let order_item_id_str: String = row.get("order_item_id")?;
76        let status_str: String = row.get("status")?;
77        let qty_req: String = row.get("quantity_requested")?;
78        let qty_pick: String = row.get("quantity_picked")?;
79        let qty_short: String = row.get("quantity_short")?;
80        let lot_id_str: Option<String> = row.get("lot_id")?;
81        let started_str: Option<String> = row.get("started_at")?;
82        let completed_str: Option<String> = row.get("completed_at")?;
83
84        Ok(PickTask {
85            id: parse_uuid_row(&id_str, "pick_task", "id")?,
86            wave_id: parse_uuid_opt_row(wave_id_str, "pick_task", "wave_id")?
87                .map(FulfillmentId::from),
88            order_id: OrderId::from(parse_uuid_row(&order_id_str, "pick_task", "order_id")?),
89            order_item_id: OrderItemId::from(parse_uuid_row(
90                &order_item_id_str,
91                "pick_task",
92                "order_item_id",
93            )?),
94            warehouse_id: row.get("warehouse_id")?,
95            status: parse_enum_row(&status_str, "pick_task", "status")?,
96            sku: row.get("sku")?,
97            product_name: row.get("product_name")?,
98            source_location_id: row.get("source_location_id")?,
99            source_location_code: row.get("source_location_code")?,
100            quantity_requested: parse_decimal_row(&qty_req, "pick_task", "quantity_requested")?,
101            quantity_picked: parse_decimal_row(&qty_pick, "pick_task", "quantity_picked")?,
102            quantity_short: parse_decimal_row(&qty_short, "pick_task", "quantity_short")?,
103            lot_id: parse_uuid_opt_row(lot_id_str, "pick_task", "lot_id")?,
104            serial_number: row.get("serial_number")?,
105            assigned_to: row.get("assigned_to")?,
106            priority: row.get("priority")?,
107            pick_sequence: row.get("pick_sequence")?,
108            started_at: parse_datetime_opt_row(started_str, "pick_task", "started_at")?,
109            completed_at: parse_datetime_opt_row(completed_str, "pick_task", "completed_at")?,
110            notes: row.get("notes")?,
111            created_at: parse_datetime_row(
112                &row.get::<_, String>("created_at")?,
113                "pick_task",
114                "created_at",
115            )?,
116            updated_at: parse_datetime_row(
117                &row.get::<_, String>("updated_at")?,
118                "pick_task",
119                "updated_at",
120            )?,
121        })
122    }
123
124    /// Read a single pick task by id from within a transaction (returns
125    /// `NotFound` smuggled through `rusqlite::Error` if it does not exist).
126    fn read_pick_by_id_tx(
127        tx: &rusqlite::Transaction<'_>,
128        id_str: &str,
129    ) -> std::result::Result<PickTask, rusqlite::Error> {
130        let mut stmt = tx.prepare("SELECT * FROM pick_tasks WHERE id = ?1")?;
131        let mut rows = stmt.query(params![id_str])?;
132        match rows.next()? {
133            Some(row) => Self::row_to_pick(row),
134            None => Err(rusqlite::Error::ToSqlConversionFailure(Box::new(CommerceError::NotFound))),
135        }
136    }
137
138    fn row_to_pack(row: &rusqlite::Row<'_>) -> rusqlite::Result<PackTask> {
139        let id_str: String = row.get("id")?;
140        let order_id_str: String = row.get("order_id")?;
141        let shipment_id_str: Option<String> = row.get("shipment_id")?;
142        let status_str: String = row.get("status")?;
143        let weight_str: Option<String> = row.get("total_weight_kg")?;
144        let started_str: Option<String> = row.get("started_at")?;
145        let completed_str: Option<String> = row.get("completed_at")?;
146
147        Ok(PackTask {
148            id: parse_uuid_row(&id_str, "pack_task", "id")?,
149            order_id: OrderId::from(parse_uuid_row(&order_id_str, "pack_task", "order_id")?),
150            shipment_id: parse_uuid_opt_row(shipment_id_str, "pack_task", "shipment_id")?
151                .map(ShipmentId::from),
152            status: parse_enum_row(&status_str, "pack_task", "status")?,
153            carton_count: row.get("carton_count")?,
154            total_weight_kg: parse_decimal_opt_row(weight_str, "pack_task", "total_weight_kg")?,
155            assigned_to: row.get("assigned_to")?,
156            packing_station: row.get("packing_station")?,
157            started_at: parse_datetime_opt_row(started_str, "pack_task", "started_at")?,
158            completed_at: parse_datetime_opt_row(completed_str, "pack_task", "completed_at")?,
159            notes: row.get("notes")?,
160            created_at: parse_datetime_row(
161                &row.get::<_, String>("created_at")?,
162                "pack_task",
163                "created_at",
164            )?,
165            updated_at: parse_datetime_row(
166                &row.get::<_, String>("updated_at")?,
167                "pack_task",
168                "updated_at",
169            )?,
170        })
171    }
172
173    fn row_to_carton(row: &rusqlite::Row<'_>) -> rusqlite::Result<Carton> {
174        let id_str: String = row.get("id")?;
175        let pack_task_id_str: String = row.get("pack_task_id")?;
176        let pkg_type_str: String = row.get("package_type")?;
177        let weight_str: Option<String> = row.get("weight_kg")?;
178        let length_str: Option<String> = row.get("length_cm")?;
179        let width_str: Option<String> = row.get("width_cm")?;
180        let height_str: Option<String> = row.get("height_cm")?;
181
182        Ok(Carton {
183            id: parse_uuid_row(&id_str, "carton", "id")?,
184            pack_task_id: parse_uuid_row(&pack_task_id_str, "carton", "pack_task_id")?,
185            carton_number: row.get("carton_number")?,
186            package_type: parse_enum_row(&pkg_type_str, "carton", "package_type")?,
187            weight_kg: parse_decimal_opt_row(weight_str, "carton", "weight_kg")?,
188            length_cm: parse_decimal_opt_row(length_str, "carton", "length_cm")?,
189            width_cm: parse_decimal_opt_row(width_str, "carton", "width_cm")?,
190            height_cm: parse_decimal_opt_row(height_str, "carton", "height_cm")?,
191            tracking_number: row.get("tracking_number")?,
192            label_printed: row.get::<_, i32>("label_printed")? != 0,
193            created_at: parse_datetime_row(
194                &row.get::<_, String>("created_at")?,
195                "carton",
196                "created_at",
197            )?,
198        })
199    }
200
201    fn row_to_carton_item(row: &rusqlite::Row<'_>) -> rusqlite::Result<CartonItem> {
202        let id_str: String = row.get("id")?;
203        let carton_id_str: String = row.get("carton_id")?;
204        let qty_str: String = row.get("quantity")?;
205        let lot_id_str: Option<String> = row.get("lot_id")?;
206
207        Ok(CartonItem {
208            id: parse_uuid_row(&id_str, "carton_item", "id")?,
209            carton_id: parse_uuid_row(&carton_id_str, "carton_item", "carton_id")?,
210            sku: row.get("sku")?,
211            quantity: parse_decimal_row(&qty_str, "carton_item", "quantity")?,
212            lot_id: parse_uuid_opt_row(lot_id_str, "carton_item", "lot_id")?,
213            serial_number: row.get("serial_number")?,
214        })
215    }
216
217    fn row_to_ship(row: &rusqlite::Row<'_>) -> rusqlite::Result<ShipTask> {
218        let id_str: String = row.get("id")?;
219        let order_id_str: String = row.get("order_id")?;
220        let shipment_id_str: String = row.get("shipment_id")?;
221        let pack_task_id_str: String = row.get("pack_task_id")?;
222        let status_str: String = row.get("status")?;
223        let cost_str: Option<String> = row.get("shipping_cost")?;
224        let shipped_str: Option<String> = row.get("shipped_at")?;
225
226        Ok(ShipTask {
227            id: parse_uuid_row(&id_str, "ship_task", "id")?,
228            order_id: OrderId::from(parse_uuid_row(&order_id_str, "ship_task", "order_id")?),
229            shipment_id: ShipmentId::from(parse_uuid_row(
230                &shipment_id_str,
231                "ship_task",
232                "shipment_id",
233            )?),
234            pack_task_id: parse_uuid_row(&pack_task_id_str, "ship_task", "pack_task_id")?,
235            status: parse_enum_row(&status_str, "ship_task", "status")?,
236            carrier: row.get("carrier")?,
237            service_level: row.get("service_level")?,
238            tracking_number: row.get("tracking_number")?,
239            label_url: row.get("label_url")?,
240            shipping_cost: parse_decimal_opt_row(cost_str, "ship_task", "shipping_cost")?,
241            assigned_to: row.get("assigned_to")?,
242            shipped_at: parse_datetime_opt_row(shipped_str, "ship_task", "shipped_at")?,
243            notes: row.get("notes")?,
244            created_at: parse_datetime_row(
245                &row.get::<_, String>("created_at")?,
246                "ship_task",
247                "created_at",
248            )?,
249            updated_at: parse_datetime_row(
250                &row.get::<_, String>("updated_at")?,
251                "ship_task",
252                "updated_at",
253            )?,
254        })
255    }
256}
257
258impl FulfillmentRepository for SqliteFulfillmentRepository {
259    // ========================================================================
260    // Wave Operations
261    // ========================================================================
262
263    fn create_wave(&self, input: CreateWave) -> Result<Wave> {
264        let conn = self.conn()?;
265        let now = Utc::now().to_rfc3339();
266        let id = FulfillmentId::new();
267        let wave_number = generate_wave_number();
268        let order_count = input.order_ids.len() as i32;
269
270        conn.execute(
271            "INSERT INTO waves (id, wave_number, warehouse_id, status, order_count, priority, notes, created_by, created_at, updated_at)
272             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?9)",
273            params![
274                id.to_string(),
275                wave_number,
276                input.warehouse_id,
277                WaveStatus::Draft.to_string(),
278                order_count,
279                input.priority.unwrap_or(0),
280                input.notes,
281                input.created_by,
282                now,
283            ],
284        ).map_err(map_db_error)?;
285
286        // Add order associations
287        for order_id in &input.order_ids {
288            conn.execute(
289                "INSERT INTO wave_orders (wave_id, order_id) VALUES (?1, ?2)",
290                params![id.to_string(), order_id.to_string()],
291            )
292            .map_err(map_db_error)?;
293        }
294
295        drop(conn);
296        self.get_wave(id)?
297            .ok_or_else(|| CommerceError::DatabaseError("Failed to create wave".into()))
298    }
299
300    fn get_wave(&self, id: FulfillmentId) -> Result<Option<Wave>> {
301        let conn = self.conn()?;
302        let mut stmt = conn.prepare("SELECT * FROM waves WHERE id = ?1").map_err(map_db_error)?;
303        let mut rows = stmt.query(params![id.to_string()]).map_err(map_db_error)?;
304
305        if let Some(row) = rows.next().map_err(map_db_error)? {
306            Ok(Some(Self::row_to_wave(row).map_err(map_db_error)?))
307        } else {
308            Ok(None)
309        }
310    }
311
312    fn get_wave_by_number(&self, number: &str) -> Result<Option<Wave>> {
313        let conn = self.conn()?;
314        let mut stmt =
315            conn.prepare("SELECT * FROM waves WHERE wave_number = ?1").map_err(map_db_error)?;
316        let mut rows = stmt.query(params![number]).map_err(map_db_error)?;
317
318        if let Some(row) = rows.next().map_err(map_db_error)? {
319            Ok(Some(Self::row_to_wave(row).map_err(map_db_error)?))
320        } else {
321            Ok(None)
322        }
323    }
324
325    fn list_waves(&self, filter: WaveFilter) -> Result<Vec<Wave>> {
326        let conn = self.conn()?;
327        let mut sql = "SELECT * FROM waves WHERE 1=1".to_string();
328        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
329
330        if let Some(warehouse_id) = filter.warehouse_id {
331            sql.push_str(" AND warehouse_id = ?");
332            params_vec.push(Box::new(warehouse_id));
333        }
334
335        if let Some(status) = filter.status {
336            sql.push_str(" AND status = ?");
337            params_vec.push(Box::new(status.to_string()));
338        }
339
340        sql.push_str(" ORDER BY priority DESC, created_at DESC");
341
342        if let Some(limit) = filter.limit {
343            sql.push_str(&format!(" LIMIT {limit}"));
344        }
345
346        let mut stmt = conn.prepare(&sql).map_err(map_db_error)?;
347        let params_refs: Vec<&dyn rusqlite::ToSql> =
348            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
349        let mut rows = stmt.query(params_refs.as_slice()).map_err(map_db_error)?;
350
351        let mut waves = Vec::new();
352        while let Some(row) = rows.next().map_err(map_db_error)? {
353            waves.push(Self::row_to_wave(row).map_err(map_db_error)?);
354        }
355        Ok(waves)
356    }
357
358    fn release_wave(&self, id: FulfillmentId) -> Result<Wave> {
359        let conn = self.conn()?;
360        let now = Utc::now().to_rfc3339();
361
362        conn.execute(
363            "UPDATE waves SET status = ?1, started_at = ?2 WHERE id = ?3 AND status = 'draft'",
364            params![WaveStatus::Released.to_string(), now, id.to_string()],
365        )
366        .map_err(map_db_error)?;
367
368        drop(conn);
369        self.get_wave(id)?
370            .ok_or_else(|| CommerceError::DatabaseError("Failed to release wave".into()))
371    }
372
373    fn complete_wave(&self, id: FulfillmentId) -> Result<Wave> {
374        let conn = self.conn()?;
375        let now = Utc::now().to_rfc3339();
376
377        conn.execute(
378            "UPDATE waves SET status = ?1, completed_at = ?2 WHERE id = ?3",
379            params![WaveStatus::Completed.to_string(), now, id.to_string()],
380        )
381        .map_err(map_db_error)?;
382
383        drop(conn);
384        self.get_wave(id)?
385            .ok_or_else(|| CommerceError::DatabaseError("Failed to complete wave".into()))
386    }
387
388    fn cancel_wave(&self, id: FulfillmentId) -> Result<Wave> {
389        let conn = self.conn()?;
390
391        conn.execute(
392            "UPDATE waves SET status = ?1 WHERE id = ?2",
393            params![WaveStatus::Cancelled.to_string(), id.to_string()],
394        )
395        .map_err(map_db_error)?;
396
397        drop(conn);
398        self.get_wave(id)?
399            .ok_or_else(|| CommerceError::DatabaseError("Failed to cancel wave".into()))
400    }
401
402    fn get_wave_orders(&self, wave_id: FulfillmentId) -> Result<Vec<OrderId>> {
403        let conn = self.conn()?;
404        let mut stmt = conn
405            .prepare("SELECT order_id FROM wave_orders WHERE wave_id = ?1")
406            .map_err(map_db_error)?;
407        let mut rows = stmt.query(params![wave_id.to_string()]).map_err(map_db_error)?;
408
409        let mut orders = Vec::new();
410        while let Some(row) = rows.next().map_err(map_db_error)? {
411            let id_str: String = row.get(0).map_err(map_db_error)?;
412            let id = OrderId::from(parse_uuid(&id_str, "wave_order", "order_id")?);
413            orders.push(id);
414        }
415        Ok(orders)
416    }
417
418    fn count_waves(&self, filter: WaveFilter) -> Result<u64> {
419        let conn = self.conn()?;
420        let mut sql = "SELECT COUNT(*) FROM waves WHERE 1=1".to_string();
421        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
422
423        if let Some(status) = filter.status {
424            sql.push_str(" AND status = ?");
425            params_vec.push(Box::new(status.to_string()));
426        }
427
428        let params_refs: Vec<&dyn rusqlite::ToSql> =
429            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
430        let count: i64 =
431            conn.query_row(&sql, params_refs.as_slice(), |row| row.get(0)).map_err(map_db_error)?;
432        Ok(count as u64)
433    }
434
435    // ========================================================================
436    // Pick Operations
437    // ========================================================================
438
439    fn create_pick(&self, input: CreatePickTask) -> Result<PickTask> {
440        let conn = self.conn()?;
441        let now = Utc::now().to_rfc3339();
442        let id = Uuid::new_v4();
443
444        conn.execute(
445            "INSERT INTO pick_tasks (id, wave_id, order_id, order_item_id, warehouse_id, status, sku, product_name,
446             source_location_id, quantity_requested, lot_id, serial_number, priority, notes, created_at, updated_at)
447             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?15)",
448            params![
449                id.to_string(),
450                input.wave_id.map(|id| id.to_string()),
451                input.order_id.to_string(),
452                input.order_item_id.to_string(),
453                input.warehouse_id,
454                PickStatus::Pending.to_string(),
455                input.sku,
456                input.product_name,
457                input.source_location_id,
458                input.quantity_requested.to_string(),
459                input.lot_id.map(|id| id.to_string()),
460                input.serial_number,
461                input.priority.unwrap_or(0),
462                input.notes,
463                now,
464            ],
465        ).map_err(map_db_error)?;
466
467        drop(conn);
468        self.get_pick(id)?
469            .ok_or_else(|| CommerceError::DatabaseError("Failed to create pick".into()))
470    }
471
472    fn get_pick(&self, id: Uuid) -> Result<Option<PickTask>> {
473        let conn = self.conn()?;
474        let mut stmt =
475            conn.prepare("SELECT * FROM pick_tasks WHERE id = ?1").map_err(map_db_error)?;
476        let mut rows = stmt.query(params![id.to_string()]).map_err(map_db_error)?;
477
478        if let Some(row) = rows.next().map_err(map_db_error)? {
479            Ok(Some(Self::row_to_pick(row).map_err(map_db_error)?))
480        } else {
481            Ok(None)
482        }
483    }
484
485    fn list_picks(&self, filter: PickTaskFilter) -> Result<Vec<PickTask>> {
486        let conn = self.conn()?;
487        let mut sql = "SELECT * FROM pick_tasks WHERE 1=1".to_string();
488        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
489
490        if let Some(warehouse_id) = filter.warehouse_id {
491            sql.push_str(" AND warehouse_id = ?");
492            params_vec.push(Box::new(warehouse_id));
493        }
494
495        if let Some(wave_id) = filter.wave_id {
496            sql.push_str(" AND wave_id = ?");
497            params_vec.push(Box::new(wave_id.to_string()));
498        }
499
500        if let Some(order_id) = filter.order_id {
501            sql.push_str(" AND order_id = ?");
502            params_vec.push(Box::new(order_id.to_string()));
503        }
504
505        if let Some(status) = filter.status {
506            sql.push_str(" AND status = ?");
507            params_vec.push(Box::new(status.to_string()));
508        }
509
510        if let Some(assigned_to) = filter.assigned_to {
511            sql.push_str(" AND assigned_to = ?");
512            params_vec.push(Box::new(assigned_to));
513        }
514
515        sql.push_str(" ORDER BY priority DESC, pick_sequence");
516
517        if let Some(limit) = filter.limit {
518            sql.push_str(&format!(" LIMIT {limit}"));
519        }
520
521        let mut stmt = conn.prepare(&sql).map_err(map_db_error)?;
522        let params_refs: Vec<&dyn rusqlite::ToSql> =
523            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
524        let mut rows = stmt.query(params_refs.as_slice()).map_err(map_db_error)?;
525
526        let mut picks = Vec::new();
527        while let Some(row) = rows.next().map_err(map_db_error)? {
528            picks.push(Self::row_to_pick(row).map_err(map_db_error)?);
529        }
530        Ok(picks)
531    }
532
533    fn assign_pick(&self, id: Uuid, assigned_to: &str) -> Result<PickTask> {
534        let conn = self.conn()?;
535
536        conn.execute(
537            "UPDATE pick_tasks SET assigned_to = ?1, status = ?2 WHERE id = ?3",
538            params![assigned_to, PickStatus::Assigned.to_string(), id.to_string()],
539        )
540        .map_err(map_db_error)?;
541
542        drop(conn);
543        self.get_pick(id)?
544            .ok_or_else(|| CommerceError::DatabaseError("Failed to assign pick".into()))
545    }
546
547    fn start_pick(&self, id: Uuid) -> Result<PickTask> {
548        let conn = self.conn()?;
549        let now = Utc::now().to_rfc3339();
550
551        conn.execute(
552            "UPDATE pick_tasks SET status = ?1, started_at = ?2 WHERE id = ?3",
553            params![PickStatus::InProgress.to_string(), now, id.to_string()],
554        )
555        .map_err(map_db_error)?;
556
557        drop(conn);
558        self.get_pick(id)?
559            .ok_or_else(|| CommerceError::DatabaseError("Failed to start pick".into()))
560    }
561
562    fn complete_pick(&self, input: CompletePick) -> Result<PickTask> {
563        let now = Utc::now().to_rfc3339();
564        let short_qty = input.quantity_short.unwrap_or(Decimal::ZERO);
565        let status =
566            if short_qty > Decimal::ZERO { PickStatus::Short } else { PickStatus::Completed };
567        let pick_id_str = input.pick_id.to_string();
568
569        // The status read, guards, pick UPDATE and wave-counter increment all run
570        // inside ONE `IMMEDIATE` transaction so concurrent completions serialize
571        // and only a real state transition folds into the wave counter.
572        with_immediate_transaction(&self.pool, |tx| {
573            let (status_str, requested_str): (String, String) = tx
574                .query_row(
575                    "SELECT status, quantity_requested FROM pick_tasks WHERE id = ?1",
576                    params![pick_id_str],
577                    |row| Ok((row.get(0)?, row.get(1)?)),
578                )
579                .map_err(|e| match e {
580                    rusqlite::Error::QueryReturnedNoRows => {
581                        rusqlite::Error::ToSqlConversionFailure(Box::new(CommerceError::NotFound))
582                    }
583                    other => other,
584                })?;
585            let current_status: PickStatus = parse_enum_row(&status_str, "pick_task", "status")?;
586            let requested = parse_decimal_row(&requested_str, "pick_task", "quantity_requested")?;
587
588            // A pick that is already finalized (Completed/Short) has already
589            // incremented the wave's completed_pick_count; re-completing it must
590            // be an idempotent no-op, never a double-count.
591            if matches!(current_status, PickStatus::Completed | PickStatus::Short) {
592                return Self::read_pick_by_id_tx(tx, &pick_id_str);
593            }
594            // A cancelled pick cannot be completed.
595            if current_status == PickStatus::Cancelled {
596                return Err(rusqlite::Error::ToSqlConversionFailure(Box::new(
597                    CommerceError::ValidationError("Cannot complete a cancelled pick task".into()),
598                )));
599            }
600            // Over-pick guard: cannot pick more than was requested.
601            if input.quantity_picked > requested {
602                return Err(rusqlite::Error::ToSqlConversionFailure(Box::new(
603                    CommerceError::ValidationError(format!(
604                        "Cannot pick {} of pick task {}: only {} were requested",
605                        input.quantity_picked, input.pick_id, requested
606                    )),
607                )));
608            }
609
610            tx.execute(
611                "UPDATE pick_tasks SET status = ?1, quantity_picked = ?2, quantity_short = ?3,
612                 lot_id = COALESCE(?4, lot_id), serial_number = COALESCE(?5, serial_number),
613                 completed_at = ?6 WHERE id = ?7",
614                params![
615                    status.to_string(),
616                    input.quantity_picked.to_string(),
617                    short_qty.to_string(),
618                    input.lot_id.map(|id| id.to_string()),
619                    input.serial_number,
620                    now,
621                    pick_id_str,
622                ],
623            )?;
624
625            let pick = Self::read_pick_by_id_tx(tx, &pick_id_str)?;
626            if let Some(wave_id) = pick.wave_id {
627                tx.execute(
628                    "UPDATE waves SET completed_pick_count = completed_pick_count + 1 WHERE id = ?1",
629                    params![wave_id.to_string()],
630                )?;
631            }
632            Ok(pick)
633        })
634    }
635
636    fn report_short(&self, id: Uuid, short_qty: Decimal, reason: &str) -> Result<PickTask> {
637        let conn = self.conn()?;
638
639        conn.execute(
640            "UPDATE pick_tasks SET status = ?1, quantity_short = ?2, notes = ?3 WHERE id = ?4",
641            params![PickStatus::Short.to_string(), short_qty.to_string(), reason, id.to_string()],
642        )
643        .map_err(map_db_error)?;
644
645        drop(conn);
646        self.get_pick(id)?
647            .ok_or_else(|| CommerceError::DatabaseError("Failed to report short".into()))
648    }
649
650    fn cancel_pick(&self, id: Uuid) -> Result<PickTask> {
651        let conn = self.conn()?;
652
653        conn.execute(
654            "UPDATE pick_tasks SET status = ?1 WHERE id = ?2",
655            params![PickStatus::Cancelled.to_string(), id.to_string()],
656        )
657        .map_err(map_db_error)?;
658
659        drop(conn);
660        self.get_pick(id)?
661            .ok_or_else(|| CommerceError::DatabaseError("Failed to cancel pick".into()))
662    }
663
664    fn get_picks_for_order(&self, order_id: OrderId) -> Result<Vec<PickTask>> {
665        self.list_picks(PickTaskFilter { order_id: Some(order_id), ..Default::default() })
666    }
667
668    fn get_picks_for_wave(&self, wave_id: FulfillmentId) -> Result<Vec<PickTask>> {
669        self.list_picks(PickTaskFilter { wave_id: Some(wave_id), ..Default::default() })
670    }
671
672    fn count_picks(&self, filter: PickTaskFilter) -> Result<u64> {
673        let conn = self.conn()?;
674        let mut sql = "SELECT COUNT(*) FROM pick_tasks WHERE 1=1".to_string();
675        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
676
677        if let Some(status) = filter.status {
678            sql.push_str(" AND status = ?");
679            params_vec.push(Box::new(status.to_string()));
680        }
681
682        let params_refs: Vec<&dyn rusqlite::ToSql> =
683            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
684        let count: i64 =
685            conn.query_row(&sql, params_refs.as_slice(), |row| row.get(0)).map_err(map_db_error)?;
686        Ok(count as u64)
687    }
688
689    // ========================================================================
690    // Pack Operations
691    // ========================================================================
692
693    fn create_pack(&self, input: CreatePackTask) -> Result<PackTask> {
694        let conn = self.conn()?;
695        let now = Utc::now().to_rfc3339();
696        let id = Uuid::new_v4();
697
698        conn.execute(
699            "INSERT INTO pack_tasks (id, order_id, status, notes, created_at, updated_at)
700             VALUES (?1, ?2, ?3, ?4, ?5, ?5)",
701            params![
702                id.to_string(),
703                input.order_id.to_string(),
704                PackStatus::Pending.to_string(),
705                input.notes,
706                now,
707            ],
708        )
709        .map_err(map_db_error)?;
710
711        drop(conn);
712        self.get_pack(id)?
713            .ok_or_else(|| CommerceError::DatabaseError("Failed to create pack".into()))
714    }
715
716    fn get_pack(&self, id: Uuid) -> Result<Option<PackTask>> {
717        let conn = self.conn()?;
718        let mut stmt =
719            conn.prepare("SELECT * FROM pack_tasks WHERE id = ?1").map_err(map_db_error)?;
720        let mut rows = stmt.query(params![id.to_string()]).map_err(map_db_error)?;
721
722        if let Some(row) = rows.next().map_err(map_db_error)? {
723            Ok(Some(Self::row_to_pack(row).map_err(map_db_error)?))
724        } else {
725            Ok(None)
726        }
727    }
728
729    fn list_packs(&self, filter: PackTaskFilter) -> Result<Vec<PackTask>> {
730        let conn = self.conn()?;
731        let mut sql = "SELECT * FROM pack_tasks WHERE 1=1".to_string();
732        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
733
734        if let Some(order_id) = filter.order_id {
735            sql.push_str(" AND order_id = ?");
736            params_vec.push(Box::new(order_id.to_string()));
737        }
738
739        if let Some(status) = filter.status {
740            sql.push_str(" AND status = ?");
741            params_vec.push(Box::new(status.to_string()));
742        }
743
744        sql.push_str(" ORDER BY created_at");
745
746        if let Some(limit) = filter.limit {
747            sql.push_str(&format!(" LIMIT {limit}"));
748        }
749
750        let mut stmt = conn.prepare(&sql).map_err(map_db_error)?;
751        let params_refs: Vec<&dyn rusqlite::ToSql> =
752            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
753        let mut rows = stmt.query(params_refs.as_slice()).map_err(map_db_error)?;
754
755        let mut packs = Vec::new();
756        while let Some(row) = rows.next().map_err(map_db_error)? {
757            packs.push(Self::row_to_pack(row).map_err(map_db_error)?);
758        }
759        Ok(packs)
760    }
761
762    fn assign_pack(&self, id: Uuid, assigned_to: &str) -> Result<PackTask> {
763        let conn = self.conn()?;
764
765        conn.execute(
766            "UPDATE pack_tasks SET assigned_to = ?1 WHERE id = ?2",
767            params![assigned_to, id.to_string()],
768        )
769        .map_err(map_db_error)?;
770
771        drop(conn);
772        self.get_pack(id)?
773            .ok_or_else(|| CommerceError::DatabaseError("Failed to assign pack".into()))
774    }
775
776    fn start_pack(&self, id: Uuid) -> Result<PackTask> {
777        let conn = self.conn()?;
778        let now = Utc::now().to_rfc3339();
779
780        conn.execute(
781            "UPDATE pack_tasks SET status = ?1, started_at = ?2 WHERE id = ?3",
782            params![PackStatus::InProgress.to_string(), now, id.to_string()],
783        )
784        .map_err(map_db_error)?;
785
786        drop(conn);
787        self.get_pack(id)?
788            .ok_or_else(|| CommerceError::DatabaseError("Failed to start pack".into()))
789    }
790
791    fn complete_pack(&self, id: Uuid) -> Result<PackTask> {
792        let conn = self.conn()?;
793        let now = Utc::now().to_rfc3339();
794
795        conn.execute(
796            "UPDATE pack_tasks SET status = ?1, completed_at = ?2 WHERE id = ?3",
797            params![PackStatus::Completed.to_string(), now, id.to_string()],
798        )
799        .map_err(map_db_error)?;
800
801        drop(conn);
802        self.get_pack(id)?
803            .ok_or_else(|| CommerceError::DatabaseError("Failed to complete pack".into()))
804    }
805
806    fn add_carton(&self, input: AddCarton) -> Result<Carton> {
807        let conn = self.conn()?;
808        let now = Utc::now().to_rfc3339();
809        let id = Uuid::new_v4();
810        let carton_number = generate_carton_number();
811
812        conn.execute(
813            "INSERT INTO cartons (id, pack_task_id, carton_number, package_type, weight_kg, length_cm, width_cm, height_cm, created_at)
814             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
815            params![
816                id.to_string(),
817                input.pack_task_id.to_string(),
818                carton_number,
819                input.package_type.to_string(),
820                input.weight_kg.map(|d| d.to_string()),
821                input.length_cm.map(|d| d.to_string()),
822                input.width_cm.map(|d| d.to_string()),
823                input.height_cm.map(|d| d.to_string()),
824                now,
825            ],
826        ).map_err(map_db_error)?;
827
828        // Update carton count
829        conn.execute(
830            "UPDATE pack_tasks SET carton_count = carton_count + 1 WHERE id = ?1",
831            params![input.pack_task_id.to_string()],
832        )
833        .map_err(map_db_error)?;
834
835        let mut stmt = conn.prepare("SELECT * FROM cartons WHERE id = ?1").map_err(map_db_error)?;
836        let mut rows = stmt.query(params![id.to_string()]).map_err(map_db_error)?;
837
838        if let Some(row) = rows.next().map_err(map_db_error)? {
839            Ok(Self::row_to_carton(row).map_err(map_db_error)?)
840        } else {
841            Err(CommerceError::DatabaseError("Failed to create carton".into()))
842        }
843    }
844
845    fn add_carton_item(&self, input: AddCartonItem) -> Result<CartonItem> {
846        let conn = self.conn()?;
847        let id = Uuid::new_v4();
848
849        conn.execute(
850            "INSERT INTO carton_items (id, carton_id, sku, quantity, lot_id, serial_number)
851             VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
852            params![
853                id.to_string(),
854                input.carton_id.to_string(),
855                input.sku,
856                input.quantity.to_string(),
857                input.lot_id.map(|id| id.to_string()),
858                input.serial_number,
859            ],
860        )
861        .map_err(map_db_error)?;
862
863        let mut stmt =
864            conn.prepare("SELECT * FROM carton_items WHERE id = ?1").map_err(map_db_error)?;
865        let mut rows = stmt.query(params![id.to_string()]).map_err(map_db_error)?;
866
867        if let Some(row) = rows.next().map_err(map_db_error)? {
868            Ok(Self::row_to_carton_item(row).map_err(map_db_error)?)
869        } else {
870            Err(CommerceError::DatabaseError("Failed to create carton item".into()))
871        }
872    }
873
874    fn get_cartons(&self, pack_task_id: Uuid) -> Result<Vec<Carton>> {
875        let conn = self.conn()?;
876        let mut stmt =
877            conn.prepare("SELECT * FROM cartons WHERE pack_task_id = ?1").map_err(map_db_error)?;
878        let mut rows = stmt.query(params![pack_task_id.to_string()]).map_err(map_db_error)?;
879
880        let mut cartons = Vec::new();
881        while let Some(row) = rows.next().map_err(map_db_error)? {
882            cartons.push(Self::row_to_carton(row).map_err(map_db_error)?);
883        }
884        Ok(cartons)
885    }
886
887    fn get_carton_items(&self, carton_id: Uuid) -> Result<Vec<CartonItem>> {
888        let conn = self.conn()?;
889        let mut stmt = conn
890            .prepare("SELECT * FROM carton_items WHERE carton_id = ?1")
891            .map_err(map_db_error)?;
892        let mut rows = stmt.query(params![carton_id.to_string()]).map_err(map_db_error)?;
893
894        let mut items = Vec::new();
895        while let Some(row) = rows.next().map_err(map_db_error)? {
896            items.push(Self::row_to_carton_item(row).map_err(map_db_error)?);
897        }
898        Ok(items)
899    }
900
901    fn mark_label_printed(&self, carton_id: Uuid) -> Result<Carton> {
902        let conn = self.conn()?;
903
904        conn.execute(
905            "UPDATE cartons SET label_printed = 1 WHERE id = ?1",
906            params![carton_id.to_string()],
907        )
908        .map_err(map_db_error)?;
909
910        let mut stmt = conn.prepare("SELECT * FROM cartons WHERE id = ?1").map_err(map_db_error)?;
911        let mut rows = stmt.query(params![carton_id.to_string()]).map_err(map_db_error)?;
912
913        if let Some(row) = rows.next().map_err(map_db_error)? {
914            Ok(Self::row_to_carton(row).map_err(map_db_error)?)
915        } else {
916            Err(CommerceError::NotFound)
917        }
918    }
919
920    fn cancel_pack(&self, id: Uuid) -> Result<PackTask> {
921        let conn = self.conn()?;
922
923        conn.execute(
924            "UPDATE pack_tasks SET status = ?1 WHERE id = ?2",
925            params![PackStatus::Cancelled.to_string(), id.to_string()],
926        )
927        .map_err(map_db_error)?;
928
929        drop(conn);
930        self.get_pack(id)?
931            .ok_or_else(|| CommerceError::DatabaseError("Failed to cancel pack".into()))
932    }
933
934    fn count_packs(&self, filter: PackTaskFilter) -> Result<u64> {
935        let conn = self.conn()?;
936        let mut sql = "SELECT COUNT(*) FROM pack_tasks WHERE 1=1".to_string();
937        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
938
939        if let Some(status) = filter.status {
940            sql.push_str(" AND status = ?");
941            params_vec.push(Box::new(status.to_string()));
942        }
943
944        let params_refs: Vec<&dyn rusqlite::ToSql> =
945            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
946        let count: i64 =
947            conn.query_row(&sql, params_refs.as_slice(), |row| row.get(0)).map_err(map_db_error)?;
948        Ok(count as u64)
949    }
950
951    // ========================================================================
952    // Ship Operations
953    // ========================================================================
954
955    fn create_ship(&self, input: CreateShipTask) -> Result<ShipTask> {
956        let conn = self.conn()?;
957        let now = Utc::now().to_rfc3339();
958        let id = Uuid::new_v4();
959
960        conn.execute(
961            "INSERT INTO ship_tasks (id, order_id, shipment_id, pack_task_id, status, carrier, service_level, notes, created_at, updated_at)
962             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?9)",
963            params![
964                id.to_string(),
965                input.order_id.to_string(),
966                input.shipment_id.to_string(),
967                input.pack_task_id.to_string(),
968                ShipStatus::Pending.to_string(),
969                input.carrier,
970                input.service_level,
971                input.notes,
972                now,
973            ],
974        ).map_err(map_db_error)?;
975
976        drop(conn);
977        self.get_ship(id)?
978            .ok_or_else(|| CommerceError::DatabaseError("Failed to create ship task".into()))
979    }
980
981    fn get_ship(&self, id: Uuid) -> Result<Option<ShipTask>> {
982        let conn = self.conn()?;
983        let mut stmt =
984            conn.prepare("SELECT * FROM ship_tasks WHERE id = ?1").map_err(map_db_error)?;
985        let mut rows = stmt.query(params![id.to_string()]).map_err(map_db_error)?;
986
987        if let Some(row) = rows.next().map_err(map_db_error)? {
988            Ok(Some(Self::row_to_ship(row).map_err(map_db_error)?))
989        } else {
990            Ok(None)
991        }
992    }
993
994    fn list_ships(&self, filter: ShipTaskFilter) -> Result<Vec<ShipTask>> {
995        let conn = self.conn()?;
996        let mut sql = "SELECT * FROM ship_tasks WHERE 1=1".to_string();
997        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
998
999        if let Some(order_id) = filter.order_id {
1000            sql.push_str(" AND order_id = ?");
1001            params_vec.push(Box::new(order_id.to_string()));
1002        }
1003
1004        if let Some(status) = filter.status {
1005            sql.push_str(" AND status = ?");
1006            params_vec.push(Box::new(status.to_string()));
1007        }
1008
1009        sql.push_str(" ORDER BY created_at");
1010
1011        if let Some(limit) = filter.limit {
1012            sql.push_str(&format!(" LIMIT {limit}"));
1013        }
1014
1015        let mut stmt = conn.prepare(&sql).map_err(map_db_error)?;
1016        let params_refs: Vec<&dyn rusqlite::ToSql> =
1017            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
1018        let mut rows = stmt.query(params_refs.as_slice()).map_err(map_db_error)?;
1019
1020        let mut ships = Vec::new();
1021        while let Some(row) = rows.next().map_err(map_db_error)? {
1022            ships.push(Self::row_to_ship(row).map_err(map_db_error)?);
1023        }
1024        Ok(ships)
1025    }
1026
1027    fn assign_ship(&self, id: Uuid, assigned_to: &str) -> Result<ShipTask> {
1028        let conn = self.conn()?;
1029
1030        conn.execute(
1031            "UPDATE ship_tasks SET assigned_to = ?1 WHERE id = ?2",
1032            params![assigned_to, id.to_string()],
1033        )
1034        .map_err(map_db_error)?;
1035
1036        drop(conn);
1037        self.get_ship(id)?
1038            .ok_or_else(|| CommerceError::DatabaseError("Failed to assign ship".into()))
1039    }
1040
1041    fn print_label(&self, id: Uuid, label_url: &str) -> Result<ShipTask> {
1042        let conn = self.conn()?;
1043
1044        conn.execute(
1045            "UPDATE ship_tasks SET status = ?1, label_url = ?2 WHERE id = ?3",
1046            params![ShipStatus::LabelPrinted.to_string(), label_url, id.to_string()],
1047        )
1048        .map_err(map_db_error)?;
1049
1050        drop(conn);
1051        self.get_ship(id)?
1052            .ok_or_else(|| CommerceError::DatabaseError("Failed to update ship".into()))
1053    }
1054
1055    fn complete_ship(&self, input: CompleteShip) -> Result<ShipTask> {
1056        let conn = self.conn()?;
1057        let now = Utc::now().to_rfc3339();
1058
1059        conn.execute(
1060            "UPDATE ship_tasks SET status = ?1, tracking_number = ?2, shipping_cost = ?3, shipped_at = ?4 WHERE id = ?5",
1061            params![
1062                ShipStatus::Shipped.to_string(),
1063                input.tracking_number,
1064                input.shipping_cost.map(|d| d.to_string()),
1065                now,
1066                input.ship_task_id.to_string(),
1067            ],
1068        ).map_err(map_db_error)?;
1069
1070        drop(conn);
1071        self.get_ship(input.ship_task_id)?
1072            .ok_or_else(|| CommerceError::DatabaseError("Failed to complete ship".into()))
1073    }
1074
1075    fn cancel_ship(&self, id: Uuid) -> Result<ShipTask> {
1076        let conn = self.conn()?;
1077
1078        conn.execute(
1079            "UPDATE ship_tasks SET status = ?1 WHERE id = ?2",
1080            params![ShipStatus::Cancelled.to_string(), id.to_string()],
1081        )
1082        .map_err(map_db_error)?;
1083
1084        drop(conn);
1085        self.get_ship(id)?
1086            .ok_or_else(|| CommerceError::DatabaseError("Failed to cancel ship".into()))
1087    }
1088
1089    fn count_ships(&self, filter: ShipTaskFilter) -> Result<u64> {
1090        let conn = self.conn()?;
1091        let mut sql = "SELECT COUNT(*) FROM ship_tasks WHERE 1=1".to_string();
1092        let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
1093
1094        if let Some(status) = filter.status {
1095            sql.push_str(" AND status = ?");
1096            params_vec.push(Box::new(status.to_string()));
1097        }
1098
1099        let params_refs: Vec<&dyn rusqlite::ToSql> =
1100            params_vec.iter().map(std::convert::AsRef::as_ref).collect();
1101        let count: i64 =
1102            conn.query_row(&sql, params_refs.as_slice(), |row| row.get(0)).map_err(map_db_error)?;
1103        Ok(count as u64)
1104    }
1105
1106    // ========================================================================
1107    // Workflow Helpers
1108    // ========================================================================
1109
1110    fn create_picks_for_order(
1111        &self,
1112        order_id: OrderId,
1113        warehouse_id: i32,
1114    ) -> Result<Vec<PickTask>> {
1115        let conn = self.conn()?;
1116
1117        let mut inputs = Vec::new();
1118        {
1119            // Get order items
1120            let mut stmt = conn
1121                .prepare("SELECT id, sku, name, quantity FROM order_items WHERE order_id = ?1")
1122                .map_err(map_db_error)?;
1123
1124            let mut rows = stmt.query(params![order_id.to_string()]).map_err(map_db_error)?;
1125
1126            while let Some(row) = rows.next().map_err(map_db_error)? {
1127                let item_id_str: String = row.get(0).map_err(map_db_error)?;
1128                let sku: String = row.get(1).map_err(map_db_error)?;
1129                let name: Option<String> = row.get(2).map_err(map_db_error)?;
1130                let qty: i32 = row.get(3).map_err(map_db_error)?;
1131
1132                // Find a location with inventory
1133                let location_id: i32 = conn
1134                    .query_row(
1135                        "SELECT l.id FROM locations l
1136                     JOIN location_inventory li ON l.id = li.location_id
1137                     WHERE l.warehouse_id = ?1 AND li.sku = ?2 AND l.is_pickable = 1
1138                     LIMIT 1",
1139                        params![warehouse_id, sku],
1140                        |row| row.get(0),
1141                    )
1142                    .unwrap_or(1); // Default to location 1 if not found
1143
1144                inputs.push(CreatePickTask {
1145                    wave_id: None,
1146                    order_id,
1147                    order_item_id: OrderItemId::from(parse_uuid(&item_id_str, "order_item", "id")?),
1148                    warehouse_id,
1149                    sku,
1150                    product_name: name,
1151                    source_location_id: location_id,
1152                    quantity_requested: Decimal::from(qty),
1153                    lot_id: None,
1154                    serial_number: None,
1155                    priority: None,
1156                    notes: None,
1157                });
1158            }
1159        }
1160
1161        drop(conn);
1162
1163        let mut picks = Vec::with_capacity(inputs.len());
1164        for input in inputs {
1165            picks.push(self.create_pick(input)?);
1166        }
1167
1168        Ok(picks)
1169    }
1170
1171    fn is_order_ready_to_pack(&self, order_id: OrderId) -> Result<bool> {
1172        let conn = self.conn()?;
1173
1174        let incomplete: i64 = conn.query_row(
1175            "SELECT COUNT(*) FROM pick_tasks WHERE order_id = ?1 AND status NOT IN ('completed', 'short', 'cancelled')",
1176            params![order_id.to_string()],
1177            |row| row.get(0),
1178        ).map_err(map_db_error)?;
1179
1180        Ok(incomplete == 0)
1181    }
1182
1183    fn is_order_ready_to_ship(&self, order_id: OrderId) -> Result<bool> {
1184        let conn = self.conn()?;
1185
1186        let completed: i64 = conn
1187            .query_row(
1188                "SELECT COUNT(*) FROM pack_tasks WHERE order_id = ?1 AND status = 'completed'",
1189                params![order_id.to_string()],
1190                |row| row.get(0),
1191            )
1192            .map_err(map_db_error)?;
1193
1194        Ok(completed > 0)
1195    }
1196
1197    // ========================================================================
1198    // Batch Operations
1199    // ========================================================================
1200
1201    fn create_waves_batch(&self, inputs: Vec<CreateWave>) -> Result<BatchResult<Wave>> {
1202        let mut result = BatchResult::new();
1203
1204        for (index, input) in inputs.into_iter().enumerate() {
1205            match self.create_wave(input) {
1206                Ok(wave) => result.record_success(wave),
1207                Err(e) => result.record_failure(index, None, &e),
1208            }
1209        }
1210
1211        Ok(result)
1212    }
1213
1214    fn get_picks_batch(&self, ids: Vec<Uuid>) -> Result<Vec<PickTask>> {
1215        let mut picks = Vec::new();
1216        for id in ids {
1217            if let Some(pick) = self.get_pick(id)? {
1218                picks.push(pick);
1219            }
1220        }
1221        Ok(picks)
1222    }
1223}
1224
1225#[cfg(test)]
1226mod tests {
1227    use super::*;
1228    use crate::SqliteDatabase;
1229    use rust_decimal_macros::dec;
1230    use stateset_core::{
1231        CompletePick, CreatePackTask, CreatePickTask, CreateWave, FulfillmentRepository, OrderId,
1232        OrderItemId, PickTaskFilter, WarehouseRepository, WaveFilter, WaveStatus,
1233    };
1234
1235    /// Build an in-memory DB and bootstrap a warehouse + location.
1236    /// Returns (fulfillment repo, warehouse id, location id).
1237    /// Pick tasks FK both warehouses(id) and locations(id) per migration 018.
1238    fn fresh_setup() -> (SqliteFulfillmentRepository, i32, i32) {
1239        let db = SqliteDatabase::in_memory().expect("in-memory");
1240        let wh = db
1241            .warehouse()
1242            .create_warehouse(stateset_core::CreateWarehouse {
1243                code: "WH-FULFIL".into(),
1244                name: "Fulfillment Test Warehouse".into(),
1245                warehouse_type: stateset_core::WarehouseType::Distribution,
1246                address: stateset_core::WarehouseAddress {
1247                    street1: "1 Test St".into(),
1248                    street2: None,
1249                    city: "Test City".into(),
1250                    state: "TC".into(),
1251                    postal_code: "00000".into(),
1252                    country: "US".into(),
1253                    phone: None,
1254                },
1255                timezone: None,
1256            })
1257            .expect("create warehouse");
1258        let loc = db
1259            .warehouse()
1260            .create_location(stateset_core::CreateLocation {
1261                warehouse_id: wh.id,
1262                code: Some("PICK-1".into()),
1263                location_type: stateset_core::LocationType::Pick,
1264                is_pickable: Some(true),
1265                is_receivable: Some(true),
1266                ..Default::default()
1267            })
1268            .expect("create location");
1269        (db.fulfillment(), wh.id, loc.id)
1270    }
1271
1272    fn fresh_repo() -> SqliteFulfillmentRepository {
1273        fresh_setup().0
1274    }
1275
1276    fn make_wave(
1277        repo: &SqliteFulfillmentRepository,
1278        warehouse_id: i32,
1279        orders: Vec<OrderId>,
1280    ) -> Wave {
1281        repo.create_wave(CreateWave {
1282            warehouse_id,
1283            order_ids: orders,
1284            priority: Some(5),
1285            notes: Some("test wave".into()),
1286            created_by: Some("alice".into()),
1287        })
1288        .expect("create wave")
1289    }
1290
1291    fn make_pick(
1292        repo: &SqliteFulfillmentRepository,
1293        warehouse_id: i32,
1294        location_id: i32,
1295        wave: Option<FulfillmentId>,
1296        order: OrderId,
1297        sku: &str,
1298    ) -> PickTask {
1299        repo.create_pick(CreatePickTask {
1300            wave_id: wave,
1301            order_id: order,
1302            order_item_id: OrderItemId::new(),
1303            warehouse_id,
1304            sku: sku.into(),
1305            product_name: Some(format!("Product {sku}")),
1306            source_location_id: location_id,
1307            quantity_requested: dec!(5),
1308            lot_id: None,
1309            serial_number: None,
1310            priority: Some(1),
1311            notes: None,
1312        })
1313        .expect("create pick")
1314    }
1315
1316    #[test]
1317    fn create_wave_starts_in_draft_with_orders() {
1318        let (repo, wh_id, _) = fresh_setup();
1319        let order_a = OrderId::new();
1320        let order_b = OrderId::new();
1321        let wave = make_wave(&repo, wh_id, vec![order_a, order_b]);
1322        assert_eq!(wave.warehouse_id, wh_id);
1323        assert_eq!(wave.status, WaveStatus::Draft);
1324        assert!(!wave.wave_number.is_empty());
1325
1326        let orders = repo.get_wave_orders(wave.id).expect("ok");
1327        assert_eq!(orders.len(), 2);
1328        assert!(orders.contains(&order_a) && orders.contains(&order_b));
1329    }
1330
1331    #[test]
1332    fn get_wave_and_get_wave_by_number_round_trip() {
1333        let (repo, wh_id, _) = fresh_setup();
1334        let wave = make_wave(&repo, wh_id, vec![OrderId::new()]);
1335        let by_id = repo.get_wave(wave.id).expect("ok").expect("found");
1336        assert_eq!(by_id.id, wave.id);
1337        let by_num = repo.get_wave_by_number(&wave.wave_number).expect("ok").expect("found");
1338        assert_eq!(by_num.id, wave.id);
1339        assert!(repo.get_wave_by_number("missing").expect("ok").is_none());
1340    }
1341
1342    #[test]
1343    fn complete_wave_transitions_status() {
1344        let (repo, wh_id, _) = fresh_setup();
1345        let wave = make_wave(&repo, wh_id, vec![OrderId::new()]);
1346        let done = repo.complete_wave(wave.id).expect("complete");
1347        assert_eq!(done.status, WaveStatus::Completed);
1348    }
1349
1350    #[test]
1351    fn cancel_wave_transitions_status() {
1352        let (repo, wh_id, _) = fresh_setup();
1353        let wave = make_wave(&repo, wh_id, vec![OrderId::new()]);
1354        let cancelled = repo.cancel_wave(wave.id).expect("cancel");
1355        assert_eq!(cancelled.status, WaveStatus::Cancelled);
1356    }
1357
1358    #[test]
1359    fn list_waves_filters_by_warehouse() {
1360        let (repo, wh_id, _) = fresh_setup();
1361        make_wave(&repo, wh_id, vec![OrderId::new()]);
1362        make_wave(&repo, wh_id, vec![OrderId::new()]);
1363        let waves = repo
1364            .list_waves(WaveFilter { warehouse_id: Some(wh_id), ..Default::default() })
1365            .expect("list");
1366        assert!(waves.len() >= 2);
1367        assert!(waves.iter().all(|w| w.warehouse_id == wh_id));
1368    }
1369
1370    #[test]
1371    fn create_pick_round_trips_and_lists_for_order() {
1372        let (repo, wh_id, loc_id) = fresh_setup();
1373        let order = OrderId::new();
1374        let pick = make_pick(&repo, wh_id, loc_id, None, order, "SKU-1");
1375        assert_eq!(pick.order_id, order);
1376        assert_eq!(pick.quantity_requested, dec!(5));
1377
1378        let by_id = repo.get_pick(pick.id).expect("ok").expect("found");
1379        assert_eq!(by_id.id, pick.id);
1380
1381        let for_order = repo.get_picks_for_order(order).expect("ok");
1382        assert_eq!(for_order.len(), 1);
1383    }
1384
1385    #[test]
1386    fn start_and_complete_pick_transitions() {
1387        let (repo, wh_id, loc_id) = fresh_setup();
1388        let order = OrderId::new();
1389        let pick = make_pick(&repo, wh_id, loc_id, None, order, "SKU-START");
1390        let started = repo.start_pick(pick.id).expect("start");
1391        assert_ne!(started.status, pick.status, "status should change after start");
1392
1393        let completed = repo
1394            .complete_pick(CompletePick {
1395                pick_id: pick.id,
1396                quantity_picked: dec!(5),
1397                quantity_short: None,
1398                short_reason: None,
1399                lot_id: None,
1400                serial_number: None,
1401                completed_by: Some("alice".into()),
1402            })
1403            .expect("complete");
1404        assert_eq!(completed.id, pick.id);
1405        assert_eq!(completed.quantity_picked, dec!(5));
1406    }
1407
1408    #[test]
1409    fn cancel_pick_changes_status() {
1410        let (repo, wh_id, loc_id) = fresh_setup();
1411        let pick = make_pick(&repo, wh_id, loc_id, None, OrderId::new(), "SKU-CN");
1412        let cancelled = repo.cancel_pick(pick.id).expect("cancel");
1413        assert_ne!(cancelled.status, pick.status);
1414    }
1415
1416    #[test]
1417    fn list_picks_filters_by_order() {
1418        let (repo, wh_id, loc_id) = fresh_setup();
1419        let order_a = OrderId::new();
1420        let order_b = OrderId::new();
1421        make_pick(&repo, wh_id, loc_id, None, order_a, "A1");
1422        make_pick(&repo, wh_id, loc_id, None, order_a, "A2");
1423        make_pick(&repo, wh_id, loc_id, None, order_b, "B1");
1424
1425        let picks_a = repo
1426            .list_picks(PickTaskFilter { order_id: Some(order_a), ..Default::default() })
1427            .expect("a");
1428        assert_eq!(picks_a.len(), 2);
1429    }
1430
1431    #[test]
1432    fn get_picks_for_wave_returns_picks() {
1433        let (repo, wh_id, loc_id) = fresh_setup();
1434        let wave = make_wave(&repo, wh_id, vec![OrderId::new()]);
1435        let order = OrderId::new();
1436        make_pick(&repo, wh_id, loc_id, Some(wave.id), order, "WV-1");
1437        make_pick(&repo, wh_id, loc_id, Some(wave.id), order, "WV-2");
1438        let picks = repo.get_picks_for_wave(wave.id).expect("ok");
1439        assert_eq!(picks.len(), 2);
1440    }
1441
1442    #[test]
1443    fn create_pack_task_round_trips() {
1444        let repo = fresh_repo();
1445        let order = OrderId::new();
1446        let pack = repo
1447            .create_pack(CreatePackTask { order_id: order, notes: Some("ship by tomorrow".into()) })
1448            .expect("create pack");
1449        assert_eq!(pack.order_id, order);
1450        let by_id = repo.get_pack(pack.id).expect("ok").expect("found");
1451        assert_eq!(by_id.id, pack.id);
1452    }
1453
1454    #[test]
1455    fn get_unknown_wave_returns_none() {
1456        let repo = fresh_repo();
1457        assert!(repo.get_wave(stateset_core::FulfillmentId::new()).expect("ok").is_none());
1458    }
1459
1460    #[test]
1461    fn get_unknown_pick_returns_none() {
1462        let repo = fresh_repo();
1463        assert!(repo.get_pick(Uuid::new_v4()).expect("ok").is_none());
1464    }
1465}