-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #10 from lvchkn/refactoring-1.6
Refactoring - v1.6
- Loading branch information
Showing
21 changed files
with
1,190 additions
and
1,834 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,7 @@ | ||
export interface ParseRequest { | ||
profile: string; | ||
tag: string; | ||
fromPage: number; | ||
toPage: number; | ||
isTest?: boolean; | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,6 @@ | ||
export interface Release { | ||
artist: string; | ||
album: string; | ||
genres: string[]; | ||
year: number; | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,6 @@ | ||
export type TaskStatus = "Pending" | "Completed"; | ||
|
||
export interface Task { | ||
id: string; | ||
status: TaskStatus; | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,51 +1,50 @@ | ||
import { ConsumeMessage } from "amqplib"; | ||
import { Task, upsertTask } from "../db/tasksRepo.js"; | ||
import { parseAndSave } from "../parser.js"; | ||
import { Channel, Connection, ConsumeMessage } from "amqplib"; | ||
import { upsertTask } from "../db/tasksRepo.js"; | ||
import { connectToRabbitMq } from "./amqp.js"; | ||
import { save } from "../services/dataSaver.js"; | ||
import { Task } from "../models/task.js"; | ||
|
||
export async function startConsumer(): Promise<void> { | ||
const connection = await connectToRabbitMq(); | ||
const channel = await connection.createChannel(); | ||
let channel: Channel; | ||
let connection: Connection; | ||
|
||
process.once("SIGINT", async () => { | ||
await channel.close(); | ||
await connection.close(); | ||
}); | ||
process.once("SIGINT", async () => { | ||
await channel.close(); | ||
await connection.close(); | ||
}); | ||
|
||
export async function startConsumer(): Promise<void> { | ||
connection = await connectToRabbitMq(); | ||
channel = await connection.createChannel(); | ||
|
||
const { queue } = await channel.assertQueue("parse-tasks-queue", { | ||
durable: true, | ||
}); | ||
|
||
await channel.prefetch(2); | ||
await channel.consume(queue, onMessageHandler, { noAck: false }); | ||
} | ||
|
||
async function onMessageHandler(message: ConsumeMessage | null): Promise<void> { | ||
if (message === null) { | ||
console.log("Error while receiving the message"); | ||
return; | ||
} | ||
|
||
const json = message.content.toString(); | ||
console.log("Message received!"); | ||
|
||
const parseRequest = JSON.parse(json); | ||
const result = await save(parseRequest); | ||
|
||
console.log(JSON.stringify(result)); | ||
|
||
const task: Task = { | ||
id: message.properties.messageId, | ||
status: "Completed", | ||
}; | ||
|
||
const upsertTaskResult = await upsertTask(task); | ||
console.log("upsertTaskResult completed:", JSON.stringify(upsertTaskResult)); | ||
|
||
await channel.consume( | ||
queue, | ||
async (message: ConsumeMessage | null) => { | ||
if (message !== null) { | ||
const json = message.content.toString(); | ||
console.log("Message received!"); | ||
|
||
const parseRequest = JSON.parse(json); | ||
const result = await parseAndSave(parseRequest); | ||
|
||
console.log(JSON.stringify(result)); | ||
|
||
const task: Task = { | ||
id: message.properties.messageId, | ||
status: "Completed", | ||
}; | ||
|
||
const upsertTaskResult = await upsertTask(task); | ||
console.log( | ||
"upsertTaskResult completed:", | ||
JSON.stringify(upsertTaskResult) | ||
); | ||
|
||
channel.ack(message); | ||
} else { | ||
console.log("Error while receiving the message"); | ||
} | ||
}, | ||
{ noAck: false } | ||
); | ||
channel.ack(message); | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.