Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_text/tests/detect_corpus/go/goroutines.txt

1.2 KiB, 1 run

created by r1870400018:11976, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1package worker
2
3import (
4 "context"
5 "fmt"
6 "sync"
7 "time"
8)
9
10type Job struct {
11 ID int
12 Payload string
13}
14
15type Result struct {
16 JobID int
17 Output string
18 Duration time.Duration
19}
20
21func worker(ctx context.Context, id int, jobs <-chan Job, results chan<- Result) {
22 for {
23 select {
24 case <-ctx.Done():
25 return
26 case job, ok := <-jobs:
27 if !ok {
28 return
29 }
30 start := time.Now()
31 output := process(job.Payload)
32 results <- Result{
33 JobID: job.ID,
34 Output: output,
35 Duration: time.Since(start),
36 }
37 }
38 }
39}
40
41func process(payload string) string {
42 time.Sleep(10 * time.Millisecond)
43 return fmt.Sprintf("processed: %s", payload)
44}
45
46func RunPool(ctx context.Context, numWorkers int, jobs []Job) []Result {
47 jobCh := make(chan Job, len(jobs))
48 resultCh := make(chan Result, len(jobs))
49
50 var wg sync.WaitGroup
51 for i := 0; i < numWorkers; i++ {
52 wg.Add(1)
53 go func(id int) {
54 defer wg.Done()
55 worker(ctx, id, jobCh, resultCh)
56 }(i)
57 }
58
59 for _, job := range jobs {
60 jobCh <- job
61 }
62 close(jobCh)
63
64 go func() {
65 wg.Wait()
66 close(resultCh)
67 }()
68
69 var results []Result
70 for r := range resultCh {
71 results = append(results, r)
72 }
73 return results
74}