Skip to content
Open
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
3 changes: 3 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
module go-task-task

go 1.26.3
169 changes: 167 additions & 2 deletions main.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,172 @@
package main

import "fmt"
import (
"context"
"fmt"
"math/rand"
"strings"
"sync"
"time"

"go-task-task/workspace/task"
)

func main() {
fmt.Println("Hello, Bounty Hunter!")
fmt.Println("Race Condition Fix Demo - Concurrent Task Execution")
fmt.Println(strings.Repeat("=", 60))

demonstrateConcurrentExecution()
demonstrateTimeoutHandling()
demonstrateConcurrentStateAccess()
}

func demonstrateConcurrentExecution() {
fmt.Println("\n1. Concurrent Task Execution Demo")
fmt.Println(strings.Repeat("-", 40))

runner := task.NewRunner()

tasks := []*task.Task{
task.NewTask("fast-success", func(ctx context.Context) error {
time.Sleep(10 * time.Millisecond)
return nil
}),
task.NewTask("slow-success", func(ctx context.Context) error {
time.Sleep(50 * time.Millisecond)
return nil
}),
task.NewTask("fast-failure", func(ctx context.Context) error {
time.Sleep(5 * time.Millisecond)
return fmt.Errorf("simulated fast failure")
}),
task.NewTask("context-aware", func(ctx context.Context) error {
select {
case <-time.After(30 * time.Millisecond):
return nil
case <-ctx.Done():
return ctx.Err()
}
}),
}

for _, t := range tasks {
runner.AddTask(t)
}

start := time.Now()
runner.Run(context.Background())
duration := time.Since(start)

fmt.Printf("Executed %d tasks in %v\n", len(tasks), duration)
fmt.Printf("Active: %d, Completed: %d, Failed: %d\n",
runner.GetActiveCount(), runner.GetCompletedCount(), runner.GetFailedCount())

for _, t := range tasks {
fmt.Printf("Task %s: %s (started: %v, finished: %v)\n",
t.ID, t.GetState(), t.GetStartedAt().Format("15:04:05.000"), t.GetFinishedAt().Format("15:04:05.000"))
}
}

func demonstrateTimeoutHandling() {
fmt.Println("\n2. Timeout Handling Demo")
fmt.Println(strings.Repeat("-", 40))

runner := task.NewRunner()

longTask := task.NewTask("long-running", func(ctx context.Context) error {
select {
case <-time.After(200 * time.Millisecond):
return nil
case <-ctx.Done():
return ctx.Err()
}
})

runner.AddTask(longTask)

fmt.Println("Running task with 100ms timeout...")
err := runner.RunWithTimeout(context.Background(), 100*time.Millisecond)

if err == context.DeadlineExceeded {
fmt.Println("Task correctly timed out")
fmt.Printf("Task state: %s, Error: %v\n", longTask.GetState(), longTask.GetErr())
} else {
fmt.Printf("Unexpected result: %v\n", err)
}
}

func demonstrateConcurrentStateAccess() {
fmt.Println("\n3. Concurrent State Access Demo")
fmt.Println(strings.Repeat("-", 40))

runner := task.NewRunner()
numTasks := 100

for i := 0; i < numTasks; i++ {
task := task.NewTask(fmt.Sprintf("concurrent-%d", i), func(ctx context.Context) error {
sleepTime := time.Duration(rand.Intn(20)) * time.Millisecond
time.Sleep(sleepTime)
return nil
})
runner.AddTask(task)
}

var wg sync.WaitGroup
stopChan := make(chan struct{})

wg.Add(1)
go func() {
defer wg.Done()
readCount := 0
for {
select {
case <-stopChan:
fmt.Printf("State reader performed %d reads\n", readCount)
return
default:
_ = runner.GetActiveCount()
_ = runner.GetCompletedCount()
_ = runner.GetFailedCount()
tasks := runner.GetTasks()

for _, task := range tasks {
_ = task.GetState()
_ = task.GetHistory()
}
readCount++
}
}
}()

wg.Add(1)
go func() {
defer wg.Done()
writeCount := 0
for {
select {
case <-stopChan:
fmt.Printf("Metadata writer performed %d writes\n", writeCount)
return
default:
tasks := runner.GetTasks()
for _, task := range tasks {
task.SetMetadata("timestamp", time.Now())
}
writeCount++
time.Sleep(1 * time.Millisecond)
}
}
}()

fmt.Println("Running tasks with concurrent state access...")
start := time.Now()
runner.Run(context.Background())
duration := time.Since(start)

close(stopChan)
wg.Wait()

fmt.Printf("Successfully executed %d tasks in %v with concurrent state access\n", numTasks, duration)
fmt.Printf("Final stats - Active: %d, Completed: %d, Failed: %d\n",
runner.GetActiveCount(), runner.GetCompletedCount(), runner.GetFailedCount())
}
70 changes: 50 additions & 20 deletions workspace/task/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,16 @@ package task
import (
"context"
"sync"
"sync/atomic"
"time"
)

type Runner struct {
mu sync.RWMutex
tasks []*Task
mu sync.RWMutex
tasks []*Task
activeCount int64 // Use atomic operations for this counter
completedCount int64 // Use atomic operations for this counter
failedCount int64 // Use atomic operations for this counter
}

func NewRunner() *Runner {
Expand All @@ -23,6 +27,26 @@ func (r *Runner) AddTask(t *Task) {
r.tasks = append(r.tasks, t)
}

func (r *Runner) GetActiveCount() int64 {
return atomic.LoadInt64(&r.activeCount)
}

func (r *Runner) GetCompletedCount() int64 {
return atomic.LoadInt64(&r.completedCount)
}

func (r *Runner) GetFailedCount() int64 {
return atomic.LoadInt64(&r.failedCount)
}

func (r *Runner) GetTasks() []*Task {
r.mu.RLock()
defer r.mu.RUnlock()
tasks := make([]*Task, len(r.tasks))
copy(tasks, r.tasks)
return tasks
}

func (r *Runner) Run(ctx context.Context) {
r.mu.RLock()
tasks := make([]*Task, len(r.tasks))
Expand All @@ -36,34 +60,40 @@ func (r *Runner) Run(ctx context.Context) {
defer wg.Done()

if ctx.Err() != nil {
task.mu.Lock()
task.State = StateFailed
task.Err = ctx.Err()
task.History = append(task.History, StateFailed)
task.mu.Unlock()
task.SetStateFailed(ctx.Err())
atomic.AddInt64(&r.failedCount, 1)
return
}

task.mu.Lock()
task.State = StateRunning
task.StartedAt = time.Now()
task.History = append(task.History, StateRunning)
task.mu.Unlock()
// Increment active counter
atomic.AddInt64(&r.activeCount, 1)
defer atomic.AddInt64(&r.activeCount, -1)

// Atomically update state to running and set start time
task.SetStateRunning()

err := task.Action(ctx)

task.mu.Lock()
task.FinishedAt = time.Now()
// Atomically update final state and finish time
if err != nil {
task.State = StateFailed
task.Err = err
task.History = append(task.History, StateFailed)
task.SetStateFailed(err)
atomic.AddInt64(&r.failedCount, 1)
} else {
task.State = StateCompleted
task.History = append(task.History, StateCompleted)
task.SetStateCompleted()
atomic.AddInt64(&r.completedCount, 1)
}
task.mu.Unlock()
}(t)
}
wg.Wait()
}

func (r *Runner) RunWithTimeout(ctx context.Context, timeout time.Duration) error {
if timeout > 0 {
var cancel context.CancelFunc
ctx, cancel = context.WithTimeout(ctx, timeout)
defer cancel()
}

r.Run(ctx)
return ctx.Err()
}
Loading