Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/pre-commit.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ jobs:
go-version: stable

- name: Install dependency tools
run: go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@latest | go install golang.org/x/tools/cmd/goimports@latest
run: go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@v2.12.2 | go install golang.org/x/tools/cmd/goimports@latest

- name: Set up pre-commit Cache
uses: pre-commit/action@v3.0.1
2 changes: 2 additions & 0 deletions .pre-commit-config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -13,4 +13,6 @@ repos:
- id: go-fmt
- id: go-imports
- id: go-unit-tests
pass_filenames: false
- id: golangci-lint
pass_filenames: false
17 changes: 12 additions & 5 deletions docker-compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -20,15 +20,22 @@ services:
- RABBITMQ_USER=guest
- RABBITMQ_PASSWORD=guest
- RABBITMQ_PORT=5672
- FILESTORAGE_HOST=file-storage
- FILESTORAGE_PORT=8081
- STORAGE_HOST=file-storage
- STORAGE_PORT=8081
- DOCKER_HOST=unix:///var/run/docker.sock
- WORKER_QUEUE_NAME=worker_queue
- MAX_WORKERS=10
file-storage:
image: maxit/file-storage
pull_policy: never
image: ghcr.io/mini-maxit/file-storage:latest
container_name: file-storage
# Internal server (writes + /sign) is network-only — never host-published.
# Public server (signed GET/HEAD) is exposed for local testing.
ports:
- "8081:8081"
- "8888:8888"
environment:
- PUBLIC_SERVER_PORT=8888
- INTERNAL_SERVER_PORT=8081
- SIGNING_SECRET=dev-secret

volumes:
rabbitmq_data:
77 changes: 49 additions & 28 deletions internal/pipeline/pipeline.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package pipeline
import (
"fmt"
"os"
"sync"

"github.com/mini-maxit/worker/internal/logger"
"github.com/mini-maxit/worker/internal/rabbitmq/responder"
Expand Down Expand Up @@ -33,6 +34,7 @@ type WorkerState struct {

type worker struct {
id int
mu sync.RWMutex
state WorkerState
responseQueue string
packager packager.Packager
Expand Down Expand Up @@ -67,25 +69,48 @@ func (ws *worker) GetId() int {
}

func (ws *worker) GetState() WorkerState {
ws.mu.RLock()
defer ws.mu.RUnlock()
return ws.state
}

func (ws *worker) UpdateStatus(status constants.WorkerStatus) {
ws.mu.Lock()
defer ws.mu.Unlock()
ws.state.Status = status
}

func (ws *worker) GetProcessingMessageID() string {
ws.mu.RLock()
defer ws.mu.RUnlock()
return ws.state.ProcessingMessageID
}

func (ws *worker) setProcessing(messageID, responseQueue string) {
ws.mu.Lock()
defer ws.mu.Unlock()
ws.state.ProcessingMessageID = messageID
ws.responseQueue = responseQueue
}

func (ws *worker) clearProcessing() {
ws.mu.Lock()
defer ws.mu.Unlock()
ws.state.ProcessingMessageID = ""
ws.responseQueue = ""
}

func (ws *worker) ProcessTask(messageID, responseQueue string, task *messages.TaskQueueMessage) {
ws.logger.Infof("Processing task [MsgID: %s]", messageID)
ws.setProcessing(messageID, responseQueue)
defer func() {
ws.clearProcessing()
if r := recover(); r != nil {
if err, ok := r.(error); ok {
ws.responder.PublishTaskErrorToResponseQueue(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
err,
)
} else {
Expand All @@ -94,21 +119,13 @@ func (ws *worker) ProcessTask(messageID, responseQueue string, task *messages.Ta
}
}()

ws.logger.Infof("Processing task [MsgID: %s]", messageID)
ws.state.ProcessingMessageID = messageID
ws.responseQueue = responseQueue
defer func() {
ws.state.ProcessingMessageID = ""
ws.responseQueue = ""
}()

langType, err := languages.ParseLanguageType(task.LanguageType)
if err != nil {
ws.logger.Errorf("Invalid language type %s: %s", task.LanguageType, err)
ws.responder.PublishTaskErrorToResponseQueue(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
err,
)
return
Expand All @@ -118,8 +135,8 @@ func (ws *worker) ProcessTask(messageID, responseQueue string, task *messages.Ta
if err != nil {
ws.responder.PublishTaskErrorToResponseQueue(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
err,
)
return
Expand Down Expand Up @@ -155,8 +172,8 @@ func (ws *worker) ProcessTask(messageID, responseQueue string, task *messages.Ta
if err != nil {
ws.responder.PublishTaskErrorToResponseQueue(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
err,
)
return
Expand All @@ -167,7 +184,7 @@ func (ws *worker) ProcessTask(messageID, responseQueue string, task *messages.Ta
fileInfo, statErr := os.Stat(dc.CompileErrFilePath)

if statErr == nil && fileInfo.Size() > 0 {
ws.publishCompilationError(dc, task.TestCases)
ws.publishCompilationError(dc, task.TestCases, messageID, responseQueue)
return
}
}
Expand All @@ -178,30 +195,34 @@ func (ws *worker) ProcessTask(messageID, responseQueue string, task *messages.Ta
if err != nil {
ws.responder.PublishTaskErrorToResponseQueue(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
err,
)
return
}

ws.responder.PublishPayloadTaskRespond(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
solutionResult,
)
ws.logger.Infof("Finished processing task [MsgID: %s]", messageID)
}

func (ws *worker) publishCompilationError(dirConfig *packager.TaskDirConfig, testCases []messages.TestCase) {
ws.logger.Infof("Compilation error occurred for message ID: %s", ws.state.ProcessingMessageID)
sendErr := ws.packager.SendSolutionPackage(dirConfig, testCases, true, ws.state.ProcessingMessageID)
func (ws *worker) publishCompilationError(
dirConfig *packager.TaskDirConfig,
testCases []messages.TestCase,
messageID, responseQueue string,
) {
ws.logger.Infof("Compilation error occurred for message ID: %s", messageID)
sendErr := ws.packager.SendSolutionPackage(dirConfig, testCases, true, messageID)
if sendErr != nil {
ws.responder.PublishTaskErrorToResponseQueue(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
sendErr,
)
return
Expand All @@ -213,8 +234,8 @@ func (ws *worker) publishCompilationError(dirConfig *packager.TaskDirConfig, tes
}
ws.responder.PublishPayloadTaskRespond(
constants.QueueMessageTypeTask,
ws.state.ProcessingMessageID,
ws.responseQueue,
messageID,
responseQueue,
solutionResult,
)
}
Loading
Loading