Skip to content

Commit d897cc8

Browse files
committed
refactor: change job payload handling from Buffer to string across queue implementations
1 parent fb07331 commit d897cc8

6 files changed

Lines changed: 18 additions & 18 deletions

File tree

src/core/queue.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
8686
this.emit('beforePush', event);
8787

8888
const jobData: JobData = { name: name as string, payload };
89-
const serializedPayload = Buffer.from(JSON.stringify(jobData));
89+
const serializedPayload = JSON.stringify(jobData);
9090
const id = await this.pushMessage(serializedPayload, meta);
9191

9292
const afterEvent: QueueEvent = { type: 'afterPush', id, name: name as string, payload, meta };
@@ -176,7 +176,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
176176
*/
177177
protected async handleMessage(message: QueueMessage): Promise<boolean> {
178178
try {
179-
const jobData: JobData = JSON.parse(message.payload.toString());
179+
const jobData: JobData = JSON.parse(message.payload);
180180
const { name, payload } = jobData;
181181

182182
const beforeEvent: QueueEvent = { type: 'beforeExec', id: message.id, name, payload, meta: message.meta };
@@ -214,7 +214,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
214214
*/
215215
protected async handleError(message: QueueMessage, error: unknown): Promise<boolean> {
216216
try {
217-
const jobData: JobData = JSON.parse(message.payload.toString());
217+
const jobData: JobData = JSON.parse(message.payload);
218218
const { name, payload } = jobData;
219219

220220
const errorEvent: QueueEvent = { type: 'afterError', id: message.id, name, payload, meta: message.meta, error };
@@ -237,13 +237,13 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
237237
/**
238238
* Pushes a new message to the queue storage backend.
239239
*
240-
* @param payload - Serialized job data as Buffer
240+
* @param payload - Serialized job data as string
241241
* @param meta - Job metadata including TTR, delay, priority
242242
* @returns Promise resolving to unique job ID
243243
* @protected
244244
* @abstract
245245
*/
246-
protected abstract pushMessage(payload: Buffer, meta: JobMeta): Promise<string>;
246+
protected abstract pushMessage(payload: string, meta: JobMeta): Promise<string>;
247247

248248
/**
249249
* Reserves the next available job from the queue for processing.

src/drivers/db.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,8 @@ export class DbQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, DbJob
1010
super(options);
1111
}
1212

13-
protected async pushMessage(payload: Buffer, meta: JobMeta): Promise<string> {
14-
return await this.db.insertJob(payload, meta);
13+
protected async pushMessage(payload: string, meta: JobMeta): Promise<string> {
14+
return await this.db.insertJob(Buffer.from(payload), meta);
1515
}
1616

1717
protected async reserve(timeout: number): Promise<QueueMessage | null> {
@@ -23,7 +23,7 @@ export class DbQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, DbJob
2323

2424
return {
2525
id: record.id,
26-
payload: record.payload,
26+
payload: record.payload.toString(),
2727
meta: record.meta
2828
};
2929
}

src/drivers/file.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -54,12 +54,12 @@ export class FileQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Fil
5454
}
5555
}
5656

57-
protected async pushMessage(payload: Buffer, meta: JobMeta): Promise<string> {
57+
protected async pushMessage(payload: string, meta: JobMeta): Promise<string> {
5858
const id = await this.touchIndex(async (data) => {
5959
const jobId = String(++data.lastId);
6060
const jobPath = path.join(this.path, `job${jobId}.data`);
6161

62-
await fs.writeFile(jobPath, payload);
62+
await fs.writeFile(jobPath, payload, 'utf8');
6363
if (this.fileMode !== undefined) {
6464
await fs.chmod(jobPath, this.fileMode);
6565
}
@@ -129,7 +129,7 @@ export class FileQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Fil
129129

130130
if (reserved) {
131131
const jobPath = path.join(this.path, `job${reserved.id}.data`);
132-
const payload = await fs.readFile(jobPath);
132+
const payload = await fs.readFile(jobPath, 'utf8');
133133
return {
134134
id: reserved.id,
135135
payload,

src/drivers/redis.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -136,13 +136,13 @@ export class RedisQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Re
136136
this.redis = new RedisAdapter(redisClient);
137137
}
138138

139-
protected async pushMessage(payload: Buffer, meta: JobMeta): Promise<string> {
139+
protected async pushMessage(payload: string, meta: JobMeta): Promise<string> {
140140
const id = (await this.redis.incr(this.idKey)).toString();
141141
const ttr = meta.ttr || this.ttrDefault;
142142
const now = Math.floor(Date.now() / 1000);
143143

144144
// Store message in Yii2 format: "ttr;jsonPayload"
145-
const message = `${ttr};${payload.toString('utf8')}`;
145+
const message = `${ttr};${payload}`;
146146
await this.redis.hset(this.messagesKey, id, message);
147147

148148
if (meta.delay && meta.delay > 0) {
@@ -208,7 +208,7 @@ export class RedisQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Re
208208

209209
return {
210210
id,
211-
payload: Buffer.from(payloadStr, 'utf8'),
211+
payload: payloadStr,
212212
meta: {
213213
ttr,
214214
pushedAt: new Date()

src/drivers/sqs.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,7 @@ export class SqsQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, SqsJ
4242
super(options);
4343
}
4444

45-
protected async pushMessage(payload: Buffer, meta: JobMeta): Promise<string> {
45+
protected async pushMessage(payload: string, meta: JobMeta): Promise<string> {
4646
const messageAttributes: Record<string, { StringValue: string; DataType: string }> = {};
4747

4848
if (meta.ttr) {
@@ -54,7 +54,7 @@ export class SqsQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, SqsJ
5454

5555
const result = await this.client.sendMessage({
5656
QueueUrl: this.queueUrl,
57-
MessageBody: payload.toString('utf8'),
57+
MessageBody: payload,
5858
DelaySeconds: meta.delay || 0,
5959
MessageAttributes: messageAttributes
6060
});
@@ -81,7 +81,7 @@ export class SqsQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, SqsJ
8181
if (!message || !message.Body || !message.MessageId || !message.ReceiptHandle) {
8282
return null;
8383
}
84-
const payload = Buffer.from(message.Body, 'utf8');
84+
const payload = message.Body;
8585

8686
const meta: JobMeta = {};
8787
if (message.MessageAttributes?.ttr?.StringValue) {

src/interfaces/job.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ export interface JobMeta {
1212

1313
export interface QueueMessage {
1414
id: string;
15-
payload: Buffer;
15+
payload: string;
1616
meta: JobMeta;
1717
}
1818

0 commit comments

Comments
 (0)