| name | mssql |
| description | Comprehensive guide for Microsoft SQL Server connectivity using the mssql Node.js library. This skill should be used when connecting to SQL Server databases, configuring connection pools, executing queries (simple, parameterized, stored procedures), handling transactions, performing bulk operations, streaming large result sets, or integrating SQL Server with Node-RED nodes. |
Microsoft SQL Server (mssql)
Overview
The mssql package is a Microsoft SQL Server client for Node.js that provides both callback and Promise-based APIs. It supports connection pooling, parameterized queries, stored procedures, transactions, bulk operations, and streaming. This skill covers patterns for building robust SQL Server integrations with emphasis on Node-RED node development.
Quick Reference
| Task | Pattern |
|---|
| Create pool | new sql.ConnectionPool(config) |
| Connect | await pool.connect() |
| Simple query | await pool.request().query('SELECT ...') |
| Parameterized query | request.input('name', sql.VarChar, value).query('SELECT ... WHERE name = @name') |
| Stored procedure | await request.execute('sp_name') |
| Transaction | const tx = new sql.Transaction(pool) |
| Bulk insert | const table = new sql.Table('name') |
| Stream results | request.stream = true |
| Close pool | await pool.close() |
Connection Configuration
Basic Connection
const sql = require('mssql');
const config = {
user: 'username',
password: 'password',
server: 'localhost',
database: 'mydb',
options: {
encrypt: true,
trustServerCertificate: true,
},
};
const pool = new sql.ConnectionPool(config);
await pool.connect();
const result = await pool.request().query('SELECT * FROM Users');
console.log(result.recordset);
await pool.close();
Connection Configuration Options
const config = {
user: 'username',
password: 'password',
server: 'localhost',
database: 'mydb',
port: 1433,
pool: {
max: 10,
min: 0,
idleTimeoutMillis: 30000,
acquireTimeoutMillis: 15000,
},
options: {
encrypt: true,
trustServerCertificate: false,
enableArithAbort: true,
connectTimeout: 15000,
requestTimeout: 15000,
cancelTimeout: 5000,
cryptoCredentialsDetails: {
minVersion: 'TLSv1.2',
},
},
};
Authentication Methods
SQL Server Authentication
const config = {
user: 'sa',
password: 'YourPassword123!',
server: 'localhost',
database: 'mydb',
options: {
encrypt: true,
trustServerCertificate: true,
},
};
Windows Authentication (NTLM)
const config = {
server: 'localhost',
database: 'mydb',
domain: 'MYDOMAIN',
user: 'domainuser',
password: 'password',
options: {
encrypt: true,
trustServerCertificate: true,
},
};
const config = {
server: 'localhost',
database: 'mydb',
options: {
trustedConnection: true,
encrypt: true,
},
};
Azure Active Directory Authentication
const config = {
server: 'your-server.database.windows.net',
database: 'mydb',
user: 'user@yourdomain.com',
password: 'password',
options: {
encrypt: true,
},
authentication: {
type: 'azure-active-directory-password',
},
};
const config = {
server: 'your-server.database.windows.net',
database: 'mydb',
options: {
encrypt: true,
},
authentication: {
type: 'azure-active-directory-access-token',
options: {
token: 'your-access-token',
},
},
};
const config = {
server: 'your-server.database.windows.net',
database: 'mydb',
options: {
encrypt: true,
},
authentication: {
type: 'azure-active-directory-msi-vm',
},
};
Connection Pooling
const sql = require('mssql');
let pool = null;
async function getPool() {
if (!pool) {
pool = new sql.ConnectionPool(config);
pool.on('error', (err) => {
console.error('Pool error:', err);
pool = null;
});
await pool.connect();
}
return pool;
}
async function queryUsers() {
const pool = await getPool();
const result = await pool.request().query('SELECT * FROM Users');
return result.recordset;
}
async function closePool() {
if (pool) {
await pool.close();
pool = null;
}
}
Connection Events
const pool = new sql.ConnectionPool(config);
pool.on('connect', () => {
console.log('Pool connected');
});
pool.on('close', () => {
console.log('Pool closed');
});
pool.on('error', (err) => {
console.error('Pool error:', err);
});
pool.on('acquire', (connection) => {
console.log('Connection acquired');
});
pool.on('release', (connection) => {
console.log('Connection released');
});
await pool.connect();
Query Execution
Simple Queries
const pool = await sql.connect(config);
const result = await pool.request().query('SELECT * FROM Users WHERE active = 1');
console.log(result.recordset);
console.log(result.recordsets);
console.log(result.rowsAffected);
console.log(result.output);
const insertResult = await pool.request().query(`
INSERT INTO Users (name, email) VALUES ('John', 'john@example.com')
`);
console.log(insertResult.rowsAffected[0]);
Parameterized Queries (Preventing SQL Injection)
const pool = await sql.connect(config);
const request = pool.request();
request.input('name', sql.VarChar(100), 'John');
request.input('email', sql.VarChar(255), 'john@example.com');
request.input('age', sql.Int, 30);
request.input('active', sql.Bit, true);
request.input('created', sql.DateTime, new Date());
const result = await request.query(`
INSERT INTO Users (name, email, age, active, created_at)
VALUES (@name, @email, @age, @active, @created)
`);
const searchRequest = pool.request();
searchRequest.input('searchTerm', sql.VarChar(100), '%john%');
const searchResult = await searchRequest.query(`
SELECT * FROM Users WHERE name LIKE @searchTerm
`);
Prepared Statements
const pool = await sql.connect(config);
const ps = new sql.PreparedStatement(pool);
ps.input('userId', sql.Int);
await ps.prepare('SELECT * FROM Users WHERE id = @userId');
const result1 = await ps.execute({ userId: 1 });
const result2 = await ps.execute({ userId: 2 });
const result3 = await ps.execute({ userId: 3 });
await ps.unprepare();
Stored Procedures
const pool = await sql.connect(config);
const request = pool.request();
request.input('userId', sql.Int, 123);
request.input('newEmail', sql.VarChar(255), 'new@example.com');
request.output('success', sql.Bit);
request.output('message', sql.VarChar(500));
const result = await request.execute('UpdateUserEmail');
console.log(result.recordset);
console.log(result.output.success);
console.log(result.output.message);
console.log(result.returnValue);
Multiple Result Sets
const pool = await sql.connect(config);
const result = await pool.request().query(`
SELECT * FROM Users;
SELECT * FROM Orders;
SELECT COUNT(*) as total FROM Products;
`);
const users = result.recordsets[0];
const orders = result.recordsets[1];
const productCount = result.recordsets[2][0].total;
const spResult = await pool.request().execute('GetDashboardData');
const [summary, recentActivity, alerts] = spResult.recordsets;
Transaction Handling
Basic Transaction
const pool = await sql.connect(config);
const transaction = new sql.Transaction(pool);
try {
await transaction.begin();
const request = new sql.Request(transaction);
request.input('name', sql.VarChar(100), 'John');
await request.query('INSERT INTO Users (name) VALUES (@name)');
const request2 = new sql.Request(transaction);
request2.input('userId', sql.Int, 1);
request2.input('amount', sql.Decimal(10, 2), 100.00);
await request2.query('INSERT INTO Balances (user_id, amount) VALUES (@userId, @amount)');
await transaction.commit();
console.log('Transaction committed');
} catch (err) {
await transaction.rollback();
console.error('Transaction rolled back:', err);
throw err;
}
Transaction Isolation Levels
const transaction = new sql.Transaction(pool);
await transaction.begin(sql.ISOLATION_LEVEL.SERIALIZABLE);
await transaction.commit();
Savepoints
const pool = await sql.connect(config);
const transaction = new sql.Transaction(pool);
try {
await transaction.begin();
const request1 = new sql.Request(transaction);
await request1.query('INSERT INTO AuditLog (action) VALUES (\'started\')');
await new sql.Request(transaction).query('SAVE TRANSACTION SavePoint1');
try {
const request2 = new sql.Request(transaction);
await request2.query('UPDATE Inventory SET quantity = quantity - 1 WHERE id = 1');
const check = await new sql.Request(transaction).query(
'SELECT quantity FROM Inventory WHERE id = 1'
);
if (check.recordset[0].quantity < 0) {
throw new Error('Insufficient inventory');
}
} catch (innerErr) {
await new sql.Request(transaction).query('ROLLBACK TRANSACTION SavePoint1');
console.log('Rolled back to savepoint:', innerErr.message);
}
await transaction.commit();
} catch (err) {
await transaction.rollback();
throw err;
}
Transaction Events
const transaction = new sql.Transaction(pool);
transaction.on('begin', () => {
console.log('Transaction begun');
});
transaction.on('commit', () => {
console.log('Transaction committed');
});
transaction.on('rollback', (aborted) => {
console.log('Transaction rolled back', aborted ? '(aborted)' : '');
});
Bulk Operations
Bulk Insert
const pool = await sql.connect(config);
const table = new sql.Table('Users');
table.create = false;
table.columns.add('name', sql.VarChar(100), { nullable: false });
table.columns.add('email', sql.VarChar(255), { nullable: false });
table.columns.add('age', sql.Int, { nullable: true });
table.columns.add('created_at', sql.DateTime, { nullable: false });
const users = [
{ name: 'John', email: 'john@example.com', age: 30 },
{ name: 'Jane', email: 'jane@example.com', age: 25 },
{ name: 'Bob', email: 'bob@example.com', age: null },
];
users.forEach((user) => {
table.rows.add(user.name, user.email, user.age, new Date());
});
const request = pool.request();
const result = await request.bulk(table);
console.log(`Inserted ${result.rowsAffected} rows`);
Bulk Insert with Options
const table = new sql.Table('Products');
table.columns.add('id', sql.Int, { nullable: false, primary: true });
table.columns.add('name', sql.VarChar(200), { nullable: false });
table.columns.add('price', sql.Decimal(10, 2), { nullable: false });
table.columns.add('category_id', sql.Int, { nullable: true });
products.forEach((p) => {
table.rows.add(p.id, p.name, p.price, p.categoryId);
});
const request = pool.request();
request.bulk(table, {
keepNulls: true,
checkConstraints: true,
tableLock: true,
fireTriggers: false,
});
Bulk Insert with Transaction
const transaction = new sql.Transaction(pool);
await transaction.begin();
try {
const table = new sql.Table('Orders');
table.columns.add('product_id', sql.Int, { nullable: false });
table.columns.add('quantity', sql.Int, { nullable: false });
orderItems.forEach((item) => {
table.rows.add(item.productId, item.quantity);
});
const request = new sql.Request(transaction);
await request.bulk(table);
await transaction.commit();
} catch (err) {
await transaction.rollback();
throw err;
}
Streaming Large Result Sets
Basic Streaming
const pool = await sql.connect(config);
const request = pool.request();
request.stream = true;
request.on('recordset', (columns) => {
console.log('Columns:', Object.keys(columns));
});
request.on('row', (row) => {
console.log('Row:', row);
});
request.on('rowsaffected', (count) => {
console.log('Rows affected:', count);
});
request.on('error', (err) => {
console.error('Stream error:', err);
});
request.on('done', (result) => {
console.log('Stream complete');
});
request.query('SELECT * FROM LargeTable');
Streaming with Backpressure
const { Writable } = require('stream');
const pool = await sql.connect(config);
async function streamToFile(query, outputPath) {
const fs = require('fs');
const writeStream = fs.createWriteStream(outputPath);
const request = pool.request();
request.stream = true;
return new Promise((resolve, reject) => {
let paused = false;
request.on('row', (row) => {
const line = JSON.stringify(row) + '\n';
const canContinue = writeStream.write(line);
if (!canContinue && !paused) {
paused = true;
request.pause();
writeStream.once('drain', () => {
paused = false;
request.resume();
});
}
});
request.on('error', (err) => {
writeStream.end();
reject(err);
});
request.on('done', () => {
writeStream.end();
resolve();
});
request.query(query);
});
}
Streaming to Transform Pipeline
const { Transform, pipeline } = require('stream');
const pool = await sql.connect(config);
function createRowStream(query) {
const request = pool.request();
request.stream = true;
const transform = new Transform({
objectMode: true,
transform(row, encoding, callback) {
callback(null, {
...row,
processed_at: new Date().toISOString(),
});
},
});
request.on('row', (row) => {
if (!transform.write(row)) {
request.pause();
transform.once('drain', () => request.resume());
}
});
request.on('error', (err) => transform.destroy(err));
request.on('done', () => transform.end());
request.query(query);
return transform;
}
const rowStream = createRowStream('SELECT * FROM Events');
pipeline(
rowStream,
createJsonOutputStream(),
fs.createWriteStream('output.json'),
(err) => {
if (err) console.error('Pipeline failed:', err);
else console.log('Pipeline succeeded');
}
);
Error Handling
Error Types
const sql = require('mssql');
try {
await pool.request().query('SELECT * FROM NonExistentTable');
} catch (err) {
if (err instanceof sql.ConnectionError) {
console.error('Connection failed:', err.message);
}
if (err instanceof sql.RequestError) {
console.error('Query failed:', err.message);
console.error('Error number:', err.number);
console.error('State:', err.state);
console.error('Class:', err.class);
console.error('Line number:', err.lineNumber);
console.error('Procedure:', err.procName);
}
if (err instanceof sql.TransactionError) {
console.error('Transaction failed:', err.message);
}
if (err instanceof sql.PreparedStatementError) {
console.error('Prepared statement failed:', err.message);
}
if (err instanceof sql.MSSQLError) {
console.error('MSSQL error:', err.message);
console.error('Code:', err.code);
}
}
Common Error Codes
| Error Number | Description | Handling |
|---|
| 18456 | Login failed | Check credentials |
| 4060 | Cannot open database | Check database name |
| 53 | Network error | Check server/port |
| 547 | FK constraint violation | Check related records |
| 2601/2627 | Unique constraint | Handle duplicates |
| 1205 | Deadlock victim | Retry transaction |
| 8152 | String truncation | Check data length |
Retry Logic
async function executeWithRetry(fn, maxRetries = 3, delay = 1000) {
let lastError;
for (let attempt = 1; attempt <= maxRetries; attempt++) {
try {
return await fn();
} catch (err) {
lastError = err;
const isRetryable =
err.code === 'ESOCKET' ||
err.code === 'ECONNCLOSED' ||
err.number === 1205 ||
err.number === -2;
if (!isRetryable || attempt === maxRetries) {
throw err;
}
console.log(`Attempt ${attempt} failed, retrying in ${delay}ms...`);
await new Promise((resolve) => setTimeout(resolve, delay));
delay *= 2;
}
}
throw lastError;
}
const result = await executeWithRetry(async () => {
return pool.request().query('SELECT * FROM Users');
});
Error Handling Best Practices
const sql = require('mssql');
class DatabaseService {
constructor(config) {
this.config = config;
this.pool = null;
}
async connect() {
try {
this.pool = await sql.connect(this.config);
this.pool.on('error', (err) => {
console.error('Pool error:', err);
this.pool = null;
});
} catch (err) {
if (err.code === 'ESOCKET') {
throw new Error(`Cannot connect to SQL Server at ${this.config.server}`);
}
if (err.number === 18456) {
throw new Error('Invalid SQL Server credentials');
}
throw err;
}
}
async query(sql, params = {}) {
if (!this.pool) {
await this.connect();
}
const request = this.pool.request();
for (const [name, value] of Object.entries(params)) {
request.input(name, value);
}
try {
return await request.query(sql);
} catch (err) {
console.error('Query error:', {
message: err.message,
number: err.number,
state: err.state,
});
if (err.number === 547) {
throw new Error('Operation failed due to related records');
}
if (err.number === 2627 || err.number === 2601) {
throw new Error('Record already exists');
}
throw new Error('Database operation failed');
}
}
async close() {
if (this.pool) {
await this.pool.close();
this.pool = null;
}
}
}
Data Type Mappings
SQL Server to JavaScript Types
| SQL Server Type | mssql Constant | JavaScript Type |
|---|
INT | sql.Int | number |
BIGINT | sql.BigInt | string (precision loss) |
SMALLINT | sql.SmallInt | number |
TINYINT | sql.TinyInt | number |
BIT | sql.Bit | boolean |
DECIMAL(p,s) | sql.Decimal(p, s) | number |
NUMERIC(p,s) | sql.Numeric(p, s) | number |
FLOAT | sql.Float | number |
REAL | sql.Real | number |
MONEY | sql.Money | number |
SMALLMONEY | sql.SmallMoney | number |
VARCHAR(n) | sql.VarChar(n) | string |
NVARCHAR(n) | sql.NVarChar(n) | string |
CHAR(n) | sql.Char(n) | string |
NCHAR(n) | sql.NChar(n) | string |
TEXT | sql.Text | string |
NTEXT | sql.NText | string |
VARCHAR(MAX) | sql.VarChar(sql.MAX) | string |
NVARCHAR(MAX) | sql.NVarChar(sql.MAX) | string |
DATETIME | sql.DateTime | Date |
DATETIME2 | sql.DateTime2 | Date |
DATETIMEOFFSET | sql.DateTimeOffset | Date |
DATE | sql.Date | Date |
TIME | sql.Time | Date |
SMALLDATETIME | sql.SmallDateTime | Date |
BINARY(n) | sql.Binary(n) | Buffer |
VARBINARY(n) | sql.VarBinary(n) | Buffer |
VARBINARY(MAX) | sql.VarBinary(sql.MAX) | Buffer |
IMAGE | sql.Image | Buffer |
UNIQUEIDENTIFIER | sql.UniqueIdentifier | string |
XML | sql.Xml | string |
GEOGRAPHY | sql.Geography | object (GeoJSON) |
GEOMETRY | sql.Geometry | object |
Type Usage Examples
const sql = require('mssql');
const request = pool.request();
request.input('name', sql.VarChar(100), 'John');
request.input('description', sql.NVarChar(sql.MAX), 'Long text with unicode: 日本語');
request.input('count', sql.Int, 42);
request.input('price', sql.Decimal(10, 2), 99.99);
request.input('percentage', sql.Float, 0.15);
request.input('active', sql.Bit, true);
request.input('created', sql.DateTime2, new Date());
request.input('dateOnly', sql.Date, new Date('2024-01-15'));
request.input('fileData', sql.VarBinary(sql.MAX), Buffer.from('binary data'));
request.input('guid', sql.UniqueIdentifier, 'a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11');
const tvp = new sql.Table();
tvp.columns.add('id', sql.Int);
tvp.columns.add('name', sql.VarChar(100));
tvp.rows.add(1, 'First');
tvp.rows.add(2, 'Second');
request.input('items', tvp);
Node-RED Integration
Connection Configuration Node
module.exports = function (RED) {
const sql = require('mssql');
function MssqlConfigNode(config) {
RED.nodes.createNode(this, config);
const node = this;
this.connectionConfig = {
server: config.server,
port: parseInt(config.port) || 1433,
database: config.database,
user: this.credentials.username,
password: this.credentials.password,
options: {
encrypt: config.encrypt,
trustServerCertificate: config.trustServerCertificate,
connectTimeout: parseInt(config.connectTimeout) || 15000,
requestTimeout: parseInt(config.requestTimeout) || 15000,
},
pool: {
max: parseInt(config.poolMax) || 10,
min: parseInt(config.poolMin) || 0,
idleTimeoutMillis: parseInt(config.poolIdleTimeout) || 30000,
},
};
this.pool = null;
this.connecting = false;
this.users = new Set();
this.getPool = async function () {
if (this.pool && this.pool.connected) {
return this.pool;
}
if (this.connecting) {
await new Promise((resolve) => {
const check = setInterval(() => {
if (!this.connecting) {
clearInterval(check);
resolve();
}
}, 100);
});
return this.pool;
}
this.connecting = true;
try {
this.pool = new sql.ConnectionPool(this.connectionConfig);
this.pool.on('error', (err) => {
node.error(`Pool error: ${err.message}`);
this.pool = null;
});
await this.pool.connect();
node.log(`Connected to ${config.server}/${config.database}`);
return this.pool;
} catch (err) {
this.pool = null;
throw err;
} finally {
this.connecting = false;
}
};
this.register = function (queryNode) {
this.users.add(queryNode.id);
};
this.deregister = async function (queryNode, done) {
this.users.delete(queryNode.id);
if (this.users.size === 0 && this.pool) {
try {
await this.pool.close();
this.pool = null;
node.log('Pool closed (no active users)');
} catch (err) {
node.error(`Error closing pool: ${err.message}`);
}
}
done();
};
this.on('close', async function (done) {
if (this.pool) {
try {
await this.pool.close();
this.pool = null;
} catch (err) {
node.error(`Error closing pool: ${err.message}`);
}
}
done();
});
}
RED.nodes.registerType('mssql-config', MssqlConfigNode, {
credentials: {
username: { type: 'text' },
password: { type: 'password' },
},
});
};
Query Node
module.exports = function (RED) {
const sql = require('mssql');
function MssqlQueryNode(config) {
RED.nodes.createNode(this, config);
const node = this;
this.configNode = RED.nodes.getNode(config.mssqlConfig);
if (!this.configNode) {
node.error('No database configuration');
node.status({ fill: 'red', shape: 'ring', text: 'no config' });
return;
}
this.configNode.register(this);
node.on('input', async function (msg, send, done) {
const query = config.query || msg.query;
if (!query) {
node.status({ fill: 'red', shape: 'ring', text: 'no query' });
done(new Error('No query specified'));
return;
}
node.status({ fill: 'blue', shape: 'dot', text: 'querying' });
try {
const pool = await node.configNode.getPool();
const request = pool.request();
if (msg.params && typeof msg.params === 'object') {
for (const [name, param] of Object.entries(msg.params)) {
if (typeof param === 'object' && param.type && 'value' in param) {
request.input(name, param.type, param.value);
} else {
request.input(name, param);
}
}
}
const result = await request.query(query);
msg.payload = result.recordset || [];
msg.recordsets = result.recordsets;
msg.rowsAffected = result.rowsAffected;
node.status({
fill: 'green',
shape: 'dot',
text: `${result.recordset?.length || 0} rows`,
});
send(msg);
done();
} catch (err) {
node.status({ fill: 'red', shape: 'ring', text: err.message });
done(err);
}
});
node.on('close', function (done) {
node.configNode.deregister(node, done);
});
}
RED.nodes.registerType('mssql-query', MssqlQueryNode);
};
Stored Procedure Node
module.exports = function (RED) {
const sql = require('mssql');
function MssqlProcedureNode(config) {
RED.nodes.createNode(this, config);
const node = this;
this.configNode = RED.nodes.getNode(config.mssqlConfig);
if (!this.configNode) {
node.error('No database configuration');
return;
}
this.configNode.register(this);
node.on('input', async function (msg, send, done) {
const procedure = config.procedure || msg.procedure;
if (!procedure) {
done(new Error('No procedure specified'));
return;
}
node.status({ fill: 'blue', shape: 'dot', text: 'executing' });
try {
const pool = await node.configNode.getPool();
const request = pool.request();
if (msg.params) {
for (const [name, param] of Object.entries(msg.params)) {
if (typeof param === 'object' && param.type) {
request.input(name, param.type, param.value);
} else {
request.input(name, param);
}
}
}
if (msg.outputParams) {
for (const [name, param] of Object.entries(msg.outputParams)) {
request.output(name, param.type, param.value);
}
}
const result = await request.execute(procedure);
msg.payload = result.recordset || [];
msg.recordsets = result.recordsets;
msg.output = result.output;
msg.returnValue = result.returnValue;
msg.rowsAffected = result.rowsAffected;
node.status({ fill: 'green', shape: 'dot', text: 'done' });
send(msg);
done();
} catch (err) {
node.status({ fill: 'red', shape: 'ring', text: err.message });
done(err);
}
});
node.on('close', function (done) {
node.configNode.deregister(node, done);
});
}
RED.nodes.registerType('mssql-procedure', MssqlProcedureNode);
};
Transaction Node Pattern
module.exports = function (RED) {
const sql = require('mssql');
function MssqlTransactionNode(config) {
RED.nodes.createNode(this, config);
const node = this;
this.configNode = RED.nodes.getNode(config.mssqlConfig);
if (!this.configNode) {
node.error('No database configuration');
return;
}
this.configNode.register(this);
node.on('input', async function (msg, send, done) {
const action = msg.action || config.action || 'begin';
try {
const pool = await node.configNode.getPool();
switch (action) {
case 'begin': {
const transaction = new sql.Transaction(pool);
await transaction.begin();
msg._transaction = transaction;
msg._transactionId = Date.now().toString();
node.status({ fill: 'yellow', shape: 'dot', text: 'transaction active' });
break;
}
case 'commit': {
if (!msg._transaction) {
throw new Error('No active transaction');
}
await msg._transaction.commit();
delete msg._transaction;
delete msg._transactionId;
node.status({ fill: 'green', shape: 'dot', text: 'committed' });
break;
}
case 'rollback': {
if (!msg._transaction) {
throw new Error('No active transaction');
}
await msg._transaction.rollback();
delete msg._transaction;
delete msg._transactionId;
node.status({ fill: 'yellow', shape: 'ring', text: 'rolled back' });
break;
}
case 'query': {
if (!msg._transaction) {
throw new Error('No active transaction');
}
const request = new sql.Request(msg._transaction);
if (msg.params) {
for (const [name, param] of Object.entries(msg.params)) {
if (typeof param === 'object' && param.type) {
request.input(name, param.type, param.value);
} else {
request.input(name, param);
}
}
}
const result = await request.query(msg.query);
msg.payload = result.recordset || [];
msg.rowsAffected = result.rowsAffected;
break;
}
default:
throw new Error(`Unknown action: ${action}`);
}
send(msg);
done();
} catch (err) {
if (msg._transaction) {
try {
await msg._transaction.rollback();
} catch (rollbackErr) {
node.error(`Rollback failed: ${rollbackErr.message}`);
}
delete msg._transaction;
delete msg._transactionId;
}
node.status({ fill: 'red', shape: 'ring', text: err.message });
done(err);
}
});
node.on('close', function (done) {
node.configNode.deregister(node, done);
});
}
RED.nodes.registerType('mssql-transaction', MssqlTransactionNode);
};
Testing Patterns
Mocking the mssql Package
const sinon = require('sinon');
function createMockPool() {
const mockRequest = {
input: sinon.stub().returnsThis(),
output: sinon.stub().returnsThis(),
query: sinon.stub(),
execute: sinon.stub(),
bulk: sinon.stub(),
};
const mockTransaction = {
begin: sinon.stub().resolves(),
commit: sinon.stub().resolves(),
rollback: sinon.stub().resolves(),
};
const mockPool = {
connect: sinon.stub().resolves(),
close: sinon.stub().resolves(),
request: sinon.stub().returns(mockRequest),
connected: true,
on: sinon.stub(),
};
return { mockPool, mockRequest, mockTransaction };
}
describe('DatabaseService', function () {
let sql;
let mockPool;
let mockRequest;
beforeEach(function () {
const mocks = createMockPool();
mockPool = mocks.mockPool;
mockRequest = mocks.mockRequest;
sql = {
ConnectionPool: sinon.stub().returns(mockPool),
Request: sinon.stub().returns(mockRequest),
Int: 'int',
VarChar: sinon.stub().returns('varchar'),
};
});
it('should execute query with parameters', async function () {
mockRequest.query.resolves({
recordset: [{ id: 1, name: 'Test' }],
rowsAffected: [1],
});
const service = new DatabaseService(sql, config);
await service.connect();
const result = await service.query('SELECT * FROM Users WHERE id = @id', {
id: 1,
});
expect(mockRequest.input.calledWith('id', 1)).to.be.true;
expect(result.recordset).to.have.length(1);
});
});
Testing Node-RED Nodes
const helper = require('node-red-node-test-helper');
const mssqlConfigNode = require('../nodes/mssql-config');
const mssqlQueryNode = require('../nodes/mssql-query');
const sinon = require('sinon');
describe('MSSQL Query Node', function () {
beforeEach(function (done) {
helper.startServer(done);
});
afterEach(function (done) {
helper.unload().then(() => helper.stopServer(done));
});
it('should query database and output results', async function () {
const mockResult = {
recordset: [{ id: 1, name: 'Test User' }],
rowsAffected: [1],
};
const mockRequest = {
input: sinon.stub().returnsThis(),
query: sinon.stub().resolves(mockResult),
};
const mockPool = {
connect: sinon.stub().resolves(),
close: sinon.stub().resolves(),
request: sinon.stub().returns(mockRequest),
connected: true,
on: sinon.stub(),
};
const proxyquire = require('proxyquire');
const mssqlStub = {
ConnectionPool: sinon.stub().returns(mockPool),
};
const flow = [
{
id: 'config1',
type: 'mssql-config',
server: 'localhost',
database: 'testdb',
},
{
id: 'query1',
type: 'mssql-query',
mssqlConfig: 'config1',
query: 'SELECT * FROM Users',
wires: [['helper1']],
},
{ id: 'helper1', type: 'helper' },
];
await helper.load([mssqlConfigNode, mssqlQueryNode], flow, {
config1: { username: 'user', password: 'pass' },
});
const queryNode = helper.getNode('query1');
const helperNode = helper.getNode('helper1');
const msgPromise = new Promise((resolve) => {
helperNode.on('input', resolve);
});
queryNode.receive({ payload: {} });
const msg = await msgPromise;
expect(msg.payload).to.deep.equal([{ id: 1, name: 'Test User' }]);
});
});
Integration Testing
const sql = require('mssql');
describe('Database Integration Tests', function () {
const config = {
server: process.env.TEST_MSSQL_SERVER || 'localhost',
database: process.env.TEST_MSSQL_DATABASE || 'testdb',
user: process.env.TEST_MSSQL_USER || 'sa',
password: process.env.TEST_MSSQL_PASSWORD || 'YourPassword123!',
options: {
encrypt: false,
trustServerCertificate: true,
},
};
let pool;
before(async function () {
this.timeout(10000);
pool = await sql.connect(config);
await pool.request().query(`
IF OBJECT_ID('TestUsers', 'U') IS NOT NULL DROP TABLE TestUsers;
CREATE TABLE TestUsers (
id INT PRIMARY KEY IDENTITY,
name NVARCHAR(100) NOT NULL,
email NVARCHAR(255) NOT NULL
);
`);
});
after(async function () {
await pool.request().query('DROP TABLE IF EXISTS TestUsers');
await pool.close();
});
beforeEach(async function () {
await pool.request().query('DELETE FROM TestUsers');
});
it('should insert and retrieve records', async function () {
const request = pool.request();
request.input('name', sql.NVarChar(100), 'Test User');
request.input('email', sql.NVarChar(255), 'test@example.com');
await request.query(`
INSERT INTO TestUsers (name, email) VALUES (@name, @email)
`);
const result = await pool.request().query('SELECT * FROM TestUsers');
expect(result.recordset).to.have.length(1);
expect(result.recordset[0].name).to.equal('Test User');
});
it('should handle transactions', async function () {
const transaction = new sql.Transaction(pool);
await transaction.begin();
try {
const request = new sql.Request(transaction);
request.input('name', sql.NVarChar(100), 'Transaction User');
request.input('email', sql.NVarChar(255), 'tx@example.com');
await request.query(`
INSERT INTO TestUsers (name, email) VALUES (@name, @email)
`);
await transaction.rollback();
} catch (err) {
await transaction.rollback();
throw err;
}
const result = await pool.request().query('SELECT * FROM TestUsers');
expect(result.recordset).to.have.length(0);
});
});
Graceful Shutdown
const sql = require('mssql');
class DatabaseManager {
constructor(config) {
this.config = config;
this.pool = null;
this.shuttingDown = false;
}
async connect() {
this.pool = new sql.ConnectionPool(this.config);
this.pool.on('error', (err) => {
console.error('Pool error:', err);
if (!this.shuttingDown) {
this.reconnect();
}
});
await this.pool.connect();
console.log('Database connected');
}
async reconnect() {
console.log('Attempting to reconnect...');
try {
await this.connect();
} catch (err) {
console.error('Reconnection failed:', err);
setTimeout(() => this.reconnect(), 5000);
}
}
async shutdown(timeout = 10000) {
this.shuttingDown = true;
console.log('Initiating database shutdown...');
if (!this.pool) {
return;
}
return new Promise((resolve) => {
const timer = setTimeout(() => {
console.warn('Shutdown timeout, forcing close');
resolve();
}, timeout);
this.pool
.close()
.then(() => {
clearTimeout(timer);
console.log('Database connection closed gracefully');
resolve();
})
.catch((err) => {
clearTimeout(timer);
console.error('Error during shutdown:', err);
resolve();
});
});
}
}
const db = new DatabaseManager(config);
process.on('SIGTERM', async () => {
console.log('SIGTERM received');
await db.shutdown();
process.exit(0);
});
process.on('SIGINT', async () => {
console.log('SIGINT received');
await db.shutdown();
process.exit(0);
});
References
For detailed information on specific topics:
references/data-type-mappings.md - Complete SQL Server to JavaScript type mappings with examples
references/connection-strings.md - Connection string formats and authentication examples
references/error-codes.md - SQL Server error codes and handling strategies
references/performance-tuning.md - Query optimization and connection pool tuning
External Documentation