const crypto = require("crypto"); const { Op } = require("sequelize"); const db = require("../../models"); const fault = (message, code, status = 400) => Object.assign(new Error(message), { code, status }); const ttlMinutes = () => Number(process.env.INVENTORY_RESERVATION_TTL_MINUTES || 15); const assertQuantity = (value) => { if (!Number.isInteger(value) || value <= 0) throw fault("Quantity must be a positive integer", "INVALID_QUANTITY"); }; const lockBalance = (warehouseId, variantId, transaction) => db.InventoryBalance.findOne({ where: { warehouse_id: warehouseId, variant_id: variantId }, transaction, lock: transaction.LOCK.UPDATE }); const ledger = (data, transaction) => db.InventoryTransaction.create({ id: crypto.randomUUID(), occurred_at: new Date(), ...data }, { transaction }); const availabilityStatus = ({ on_hand, reserved, low_stock_threshold }) => { const available = Number(on_hand) - Number(reserved); return available <= 0 ? "OUT_OF_STOCK" : available <= Number(low_stock_threshold) ? "LOW_STOCK" : "IN_STOCK"; }; async function adjustStock({ eventId, warehouseId, variantId, quantityDelta, reason, actorUserId, requestId, transaction: externalTransaction }) { if (!eventId || !Number.isInteger(quantityDelta) || quantityDelta === 0) throw fault("A non-zero integer adjustment and eventId are required", "INVALID_ADJUSTMENT"); const execute = async transaction => { const prior = await db.InventoryTransaction.findOne({ where: { event_id: eventId }, transaction, lock: transaction.LOCK.UPDATE }); if (prior) return { idempotent: true, transaction: prior }; let balance = await lockBalance(warehouseId, variantId, transaction); if (!balance) balance = await db.InventoryBalance.create({ id: crypto.randomUUID(), warehouse_id: warehouseId, variant_id: variantId, on_hand: 0, reserved: 0 }, { transaction }); const next = Number(balance.on_hand) + quantityDelta; if (next < 0 || next < Number(balance.reserved)) throw fault("Adjustment would violate available stock", "INSUFFICIENT_STOCK", 409); await balance.update({ on_hand: next }, { transaction }); const entry = await ledger({ event_id: eventId, warehouse_id: warehouseId, variant_id: variantId, type: "ADJUSTMENT", quantity_delta: quantityDelta, reason, actor_user_id: actorUserId, request_id: requestId }, transaction); return { balance, transaction: entry, idempotent: false }; }; return externalTransaction ? execute(externalTransaction) : db.sequelize.transaction(execute); } async function selectWarehouse(variantId, quantity, transaction) { const warehouses = await db.Warehouse.findAll({ where: { status: "ACTIVE" }, order: [["is_default", "DESC"], ["code", "ASC"]], transaction }); for (const warehouse of warehouses) { const balance = await lockBalance(warehouse.id, variantId, transaction); if (balance && Number(balance.on_hand) - Number(balance.reserved) >= quantity) return { warehouse, balance }; } throw fault("Insufficient available stock", "INSUFFICIENT_STOCK", 409); } async function reserveStock({ reservationKey, warehouseId, variantId, quantity, referenceType, referenceId, userId, expiresAt, requestId, transaction: externalTransaction }) { assertQuantity(quantity); if (!reservationKey) throw fault("reservationKey is required", "INVALID_RESERVATION"); const execute = async transaction => { const existing = await db.InventoryReservation.findOne({ where: { reservation_key: reservationKey }, transaction, lock: transaction.LOCK.UPDATE }); if (existing) return { reservation: existing, idempotent: true }; const variant = await db.ProductVariant.findOne({ where: { id: variantId, status: "ACTIVE" }, transaction }); if (!variant) throw fault("Active variant not found", "VARIANT_UNAVAILABLE", 404); let selected; if (warehouseId) { const warehouse = await db.Warehouse.findOne({ where: { id: warehouseId, status: "ACTIVE" }, transaction }); if (!warehouse) throw fault("Active warehouse not found", "WAREHOUSE_UNAVAILABLE", 404); const balance = await lockBalance(warehouseId, variantId, transaction); if (!balance || Number(balance.on_hand) - Number(balance.reserved) < quantity) throw fault("Insufficient available stock", "INSUFFICIENT_STOCK", 409); selected = { warehouse, balance }; } else selected = await selectWarehouse(variantId, quantity, transaction); await selected.balance.update({ reserved: Number(selected.balance.reserved) + quantity }, { transaction }); const reservation = await db.InventoryReservation.create({ id: crypto.randomUUID(), reservation_key: reservationKey, warehouse_id: selected.warehouse.id, variant_id: variantId, quantity, reference_type: referenceType, reference_id: referenceId, user_id: userId, expires_at: expiresAt || new Date(Date.now() + ttlMinutes() * 60000) }, { transaction }); await ledger({ event_id: `reserve:${reservationKey}`, warehouse_id: selected.warehouse.id, variant_id: variantId, type: "RESERVATION", reserved_delta: quantity, reference_type: "INVENTORY_RESERVATION", reference_id: reservation.id, request_id: requestId }, transaction); return { reservation, idempotent: false }; }; return externalTransaction ? execute(externalTransaction) : db.sequelize.transaction(execute); } async function finishReservation(reservationKey, status, requestId, externalTransaction) { const execute = async transaction => { const reservation = await db.InventoryReservation.findOne({ where: { reservation_key: reservationKey }, transaction, lock: transaction.LOCK.UPDATE }); if (!reservation) throw fault("Reservation not found", "RESERVATION_NOT_FOUND", 404); if (reservation.status === status) return { reservation, idempotent: true }; if (reservation.status !== "ACTIVE") throw fault(`Reservation is ${reservation.status}`, "RESERVATION_NOT_ACTIVE", 409); if (status === "CONSUMED" && new Date(reservation.expires_at) <= new Date()) throw fault("Reservation has expired", "RESERVATION_EXPIRED", 409); const balance = await lockBalance(reservation.warehouse_id, reservation.variant_id, transaction); if (!balance || Number(balance.reserved) < reservation.quantity) throw fault("Inventory invariant violated", "INVENTORY_INVARIANT", 409); const consume = status === "CONSUMED"; await balance.update({ reserved: Number(balance.reserved) - reservation.quantity, on_hand: Number(balance.on_hand) - (consume ? reservation.quantity : 0) }, { transaction }); await reservation.update({ status, released_at: consume ? null : new Date(), consumed_at: consume ? new Date() : null }, { transaction }); await ledger({ event_id: `${status.toLowerCase()}:${reservationKey}`, warehouse_id: reservation.warehouse_id, variant_id: reservation.variant_id, type: consume ? "RESERVATION_CONSUME" : "RESERVATION_RELEASE", quantity_delta: consume ? -reservation.quantity : 0, reserved_delta: -reservation.quantity, reference_type: "INVENTORY_RESERVATION", reference_id: reservation.id, request_id: requestId }, transaction); return { reservation, idempotent: false }; }; return externalTransaction ? execute(externalTransaction) : db.sequelize.transaction(execute); } const releaseReservation = args => finishReservation(args.reservationKey, args.expired ? "EXPIRED" : "RELEASED", args.requestId, args.transaction); const consumeReservation = args => finishReservation(args.reservationKey, "CONSUMED", args.requestId, args.transaction); async function transferStock({ eventId, sourceWarehouseId, destinationWarehouseId, variantId, quantity, actorUserId, requestId }) { assertQuantity(quantity); if (sourceWarehouseId === destinationWarehouseId) throw fault("Warehouses must differ", "INVALID_TRANSFER"); return db.sequelize.transaction(async transaction => { const prior = await db.InventoryTransfer.findOne({ where: { event_id: eventId }, transaction, lock: transaction.LOCK.UPDATE }); if (prior) return { transfer: prior, idempotent: true }; const ids = [sourceWarehouseId, destinationWarehouseId].sort(); const balances = {}; for (const id of ids) balances[id] = await lockBalance(id, variantId, transaction); const source = balances[sourceWarehouseId]; if (!source || Number(source.on_hand) - Number(source.reserved) < quantity) throw fault("Insufficient transferable stock", "INSUFFICIENT_STOCK", 409); let destination = balances[destinationWarehouseId]; if (!destination) destination = await db.InventoryBalance.create({ id: crypto.randomUUID(), warehouse_id: destinationWarehouseId, variant_id: variantId, on_hand: 0, reserved: 0 }, { transaction }); await source.update({ on_hand: Number(source.on_hand) - quantity }, { transaction }); await destination.update({ on_hand: Number(destination.on_hand) + quantity }, { transaction }); const transfer = await db.InventoryTransfer.create({ id: crypto.randomUUID(), event_id: eventId, transfer_number: eventId, source_warehouse_id: sourceWarehouseId, destination_warehouse_id: destinationWarehouseId, variant_id: variantId, quantity, created_by: actorUserId, completed_by: actorUserId, completed_at: new Date() }, { transaction }); await ledger({ event_id: `${eventId}:out`, warehouse_id: sourceWarehouseId, variant_id: variantId, type: "TRANSFER_OUT", quantity_delta: -quantity, reference_type: "INVENTORY_TRANSFER", reference_id: transfer.id, actor_user_id: actorUserId, request_id: requestId }, transaction); await ledger({ event_id: `${eventId}:in`, warehouse_id: destinationWarehouseId, variant_id: variantId, type: "TRANSFER_IN", quantity_delta: quantity, reference_type: "INVENTORY_TRANSFER", reference_id: transfer.id, actor_user_id: actorUserId, request_id: requestId }, transaction); return { transfer, idempotent: false }; }); } async function expireReservations(limit = 100) { const rows = await db.InventoryReservation.findAll({ where: { status: "ACTIVE", expires_at: { [Op.lte]: new Date() } }, order: [["expires_at", "ASC"]], limit }); for (const row of rows) { try { await releaseReservation({ reservationKey: row.reservation_key, expired: true }); } catch (error) { if (error.code !== "RESERVATION_NOT_ACTIVE") throw error; } } return rows.length; } module.exports = { adjustStock, reserveStock, releaseReservation, consumeReservation, transferStock, expireReservations, availabilityStatus, ttlMinutes };