Skip to content

Commit 839602b

Browse files
committed
mongoose reset the session.
Avoid getting caught up on using a connection that perhaps was a transaction that has finised.
1 parent 24ea30e commit 839602b

1 file changed

Lines changed: 132 additions & 105 deletions

File tree

src/drivers/mongoose.ts

Lines changed: 132 additions & 105 deletions
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,31 @@
1-
import { Schema, model, Model, Document, Types } from 'mongoose';
2-
import type { UpdateQuery, FilterQuery, QueryOptions } from 'mongoose';
3-
import type { DatabaseAdapter, QueueJobRecord } from '../interfaces/database.ts';
4-
import type { JobMeta, JobStatus, BaseJobOptions, WithPriority, WithDelay } from '../interfaces/job.ts';
5-
import { DbQueue } from '../drivers/db.ts';
1+
import { Schema, model, Model, Document, Types } from "mongoose";
2+
import type { UpdateQuery, FilterQuery, QueryOptions } from "mongoose";
3+
import type {
4+
DatabaseAdapter,
5+
QueueJobRecord,
6+
} from "../interfaces/database.ts";
7+
import type {
8+
JobMeta,
9+
JobStatus,
10+
BaseJobOptions,
11+
WithPriority,
12+
WithDelay,
13+
} from "../interfaces/job.ts";
14+
import { DbQueue } from "../drivers/db.ts";
615

716
// Driver-specific job request interface
8-
export interface MongooseJobRequest<TPayload> extends BaseJobOptions, WithPriority, WithDelay {
17+
export interface MongooseJobRequest<TPayload>
18+
extends BaseJobOptions,
19+
WithPriority,
20+
WithDelay {
921
/** Job payload */
1022
payload: TPayload;
1123
// Mongoose/MongoDB queue supports both priority and delays
1224
}
1325

1426
// MongoDB document structure for queue jobs
1527
export interface IQueueJobDocument {
16-
payload: Buffer;
28+
payload: any;
1729
ttr: number;
1830
delaySeconds: number;
1931
priority: number;
@@ -22,7 +34,7 @@ export interface IQueueJobDocument {
2234
reserveTime: Date | null;
2335
doneTime: Date | null;
2436
expireTime: Date | null;
25-
status: 'waiting' | 'reserved' | 'done' | 'failed';
37+
status: "waiting" | "reserved" | "done" | "failed";
2638
attempt: number;
2739
errorMessage?: string;
2840
}
@@ -33,28 +45,31 @@ export interface IQueueJob extends IQueueJobDocument, Document {
3345
}
3446

3547
// Queue job schema
36-
export const QueueJobSchema = new Schema<IQueueJob>({
37-
payload: { type: Buffer, required: true },
38-
ttr: { type: Number, required: true, default: 300 },
39-
delaySeconds: { type: Number, required: true, default: 0 },
40-
priority: { type: Number, required: true, default: 0 },
41-
pushTime: { type: Date, required: true },
42-
delayTime: { type: Date, default: null },
43-
reserveTime: { type: Date, default: null },
44-
doneTime: { type: Date, default: null },
45-
expireTime: { type: Date, default: null },
46-
status: {
47-
type: String,
48-
required: true,
49-
enum: ['waiting', 'reserved', 'done', 'failed'],
50-
default: 'waiting'
48+
export const QueueJobSchema = new Schema<IQueueJob>(
49+
{
50+
payload: { type: Schema.Types.Mixed, required: true },
51+
ttr: { type: Number, required: true, default: 300 },
52+
delaySeconds: { type: Number, required: true, default: 0 },
53+
priority: { type: Number, required: true, default: 0 },
54+
pushTime: { type: Date, required: true },
55+
delayTime: { type: Date, default: null },
56+
reserveTime: { type: Date, default: null },
57+
doneTime: { type: Date, default: null },
58+
expireTime: { type: Date, default: null },
59+
status: {
60+
type: String,
61+
required: true,
62+
enum: ["waiting", "reserved", "done", "failed"],
63+
default: "waiting",
64+
},
65+
attempt: { type: Number, required: true, default: 0 },
66+
errorMessage: { type: String },
5167
},
52-
attempt: { type: Number, required: true, default: 0 },
53-
errorMessage: { type: String }
54-
}, {
55-
collection: 'queue_jobs',
56-
timestamps: false
57-
});
68+
{
69+
collection: "queue_jobs",
70+
timestamps: false,
71+
}
72+
);
5873

5974
// Add indexes
6075
QueueJobSchema.index({ status: 1, delayTime: 1, priority: -1, pushTime: 1 });
@@ -69,18 +84,20 @@ type MongoUpdate = UpdateQuery<IQueueJobDocument>;
6984
export class MongooseDatabaseAdapter implements DatabaseAdapter {
7085
constructor(private model: Model<IQueueJob>) {}
7186

72-
async insertJob(payload: Buffer, meta: JobMeta): Promise<string> {
87+
async insertJob(payload: any, meta: JobMeta): Promise<string> {
7388
const now = new Date();
74-
89+
7590
const doc = await this.model.create({
7691
payload,
7792
ttr: meta.ttr ?? 300,
7893
delaySeconds: meta.delaySeconds ?? 0,
7994
priority: meta.priority ?? 0,
8095
pushTime: now,
81-
delayTime: meta.delaySeconds ? new Date(now.getTime() + meta.delaySeconds * 1000) : null,
82-
status: 'waiting',
83-
attempt: 0
96+
delayTime: meta.delaySeconds
97+
? new Date(now.getTime() + meta.delaySeconds * 1000)
98+
: null,
99+
status: "waiting",
100+
attempt: 0,
84101
});
85102

86103
return doc._id.toHexString();
@@ -91,43 +108,42 @@ export class MongooseDatabaseAdapter implements DatabaseAdapter {
91108

92109
// First, recover timed-out jobs
93110
await this.model.updateMany(
94-
{
95-
status: 'reserved',
96-
expireTime: { $lt: now }
111+
{
112+
status: "reserved",
113+
expireTime: { $lt: now },
114+
},
115+
{
116+
$set: {
117+
status: "waiting",
118+
reserveTime: null,
119+
expireTime: null,
120+
},
121+
$inc: { attempt: 1 },
97122
},
98-
{
99-
$set: {
100-
status: 'waiting',
101-
reserveTime: null,
102-
expireTime: null
103-
},
104-
$inc: { attempt: 1 }
123+
{
124+
session: undefined,
105125
}
106126
);
107127

108-
// Atomically claim the next available job
109-
const filter: MongoFilter = {
110-
status: 'waiting',
111-
$or: [
112-
{ delayTime: null },
113-
{ delayTime: { $lte: now } }
114-
]
115-
};
116-
117-
const update: MongoUpdate = {
118-
$set: {
119-
status: 'reserved',
120-
reserveTime: now
128+
const doc = await this.model.findOneAndUpdate(
129+
// Atomically claim the next available job
130+
{
131+
status: "waiting",
132+
$or: [{ delayTime: null }, { delayTime: { $lte: now } }],
133+
},
134+
{
135+
$set: {
136+
status: "reserved",
137+
reserveTime: now,
138+
},
139+
},
140+
{
141+
sort: { priority: -1, pushTime: 1 },
142+
new: true,
143+
session: undefined,
121144
}
122-
};
123-
124-
const options: QueryOptions = {
125-
sort: { priority: -1, pushTime: 1 },
126-
new: true
127-
};
145+
);
128146

129-
const doc = await this.model.findOneAndUpdate(filter, update, options);
130-
131147
if (!doc) {
132148
return null;
133149
}
@@ -136,7 +152,8 @@ export class MongooseDatabaseAdapter implements DatabaseAdapter {
136152
const ttr = doc.ttr || 300;
137153
await this.model.updateOne(
138154
{ _id: doc._id },
139-
{ $set: { expireTime: new Date(now.getTime() + ttr * 1000) } }
155+
{ $set: { expireTime: new Date(now.getTime() + ttr * 1000) } },
156+
{ session: undefined }
140157
);
141158

142159
return {
@@ -147,90 +164,101 @@ export class MongooseDatabaseAdapter implements DatabaseAdapter {
147164
delaySeconds: doc.delaySeconds,
148165
priority: doc.priority,
149166
pushedAt: doc.pushTime,
150-
reservedAt: now
167+
reservedAt: now,
151168
},
152169
pushedAt: doc.pushTime,
153-
reservedAt: now
170+
reservedAt: now,
154171
};
155172
}
156173

157174
async completeJob(id: string): Promise<void> {
158175
await this.model.updateOne(
159176
{ _id: new Types.ObjectId(id) },
160-
{ $set: { status: 'done', doneTime: new Date() } }
177+
{ $set: { status: "done", doneTime: new Date() } },
178+
{ session: undefined }
161179
);
162180
}
163181

164182
async releaseJob(id: string): Promise<void> {
165183
await this.model.updateOne(
166184
{ _id: new Types.ObjectId(id) },
167-
{
168-
$set: {
169-
status: 'waiting',
170-
reserveTime: null,
171-
expireTime: null
172-
}
173-
}
185+
{
186+
$set: {
187+
status: "waiting",
188+
reserveTime: null,
189+
expireTime: null,
190+
},
191+
},
192+
{ session: undefined }
174193
);
175194
}
176195

177196
async failJob(id: string, error: string): Promise<void> {
178197
await this.model.updateOne(
179198
{ _id: new Types.ObjectId(id) },
180-
{
181-
$set: {
182-
status: 'failed',
183-
errorMessage: error,
184-
doneTime: new Date()
185-
}
186-
}
199+
{
200+
$set: {
201+
status: "failed",
202+
errorMessage: error,
203+
doneTime: new Date(),
204+
},
205+
},
206+
{ session: undefined }
187207
);
188208
}
189209

190210
async getJobStatus(id: string): Promise<JobStatus | null> {
191-
const doc = await this.model.findOne(
192-
{ _id: new Types.ObjectId(id) },
193-
{ status: 1, delayTime: 1 }
194-
).exec();
195-
211+
const doc = await this.model
212+
.findOne(
213+
{ _id: new Types.ObjectId(id) },
214+
{ status: 1, delayTime: 1 },
215+
{ session: undefined }
216+
)
217+
.exec();
218+
196219
if (!doc) {
197220
return null;
198221
}
199-
222+
200223
// Check if job is delayed
201-
if (doc.status === 'waiting' && doc.delayTime && doc.delayTime > new Date()) {
202-
return 'delayed';
224+
if (
225+
doc.status === "waiting" &&
226+
doc.delayTime &&
227+
doc.delayTime > new Date()
228+
) {
229+
return "delayed";
203230
}
204-
231+
205232
switch (doc.status) {
206-
case 'waiting':
207-
return 'waiting';
208-
case 'reserved':
209-
return 'reserved';
210-
case 'done':
211-
case 'failed':
212-
return 'done';
233+
case "waiting":
234+
return "waiting";
235+
case "reserved":
236+
return "reserved";
237+
case "done":
238+
case "failed":
239+
return "done";
213240
default:
214241
return null;
215242
}
216243
}
217244
}
218245

219246
// Mongoose-specific queue class
220-
export class MongooseQueue<TJobMap = Record<string, unknown>> extends DbQueue<TJobMap> {
247+
export class MongooseQueue<
248+
TJobMap = Record<string, unknown>
249+
> extends DbQueue<TJobMap> {
221250
mongooseAdapter: MongooseDatabaseAdapter;
222-
251+
223252
constructor(config: { model: Model<IQueueJob>; name: string }) {
224253
const adapter = new MongooseDatabaseAdapter(config.model);
225254
super(adapter, { name: config.name });
226255
this.mongooseAdapter = adapter;
227256
}
228257
}
229258

230-
231259
// Create a default queue model
232260
export function createQueueModel(
233-
modelName: string = 'QueueJob',
261+
modelName: string = "QueueJob",
234262
collectionName?: string
235263
): Model<IQueueJob> {
236264
// Check if model already exists
@@ -240,7 +268,7 @@ export function createQueueModel(
240268
// Create new model
241269
const schema = QueueJobSchema.clone();
242270
if (collectionName) {
243-
schema.set('collection', collectionName);
271+
schema.set("collection", collectionName);
244272
}
245273
return model<IQueueJob>(modelName, schema);
246274
}
@@ -251,4 +279,3 @@ export const QueueJob = createQueueModel();
251279

252280
// Re-export for convenience
253281
export { DbQueue };
254-

0 commit comments

Comments
 (0)