diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index fb73c3c..54b55f5 100755 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -1,18 +1,17 @@ - name: CI on: push: - branches: [ master ] + branches: [master] pull_request: - branches: [ master ] + branches: [master] jobs: build: runs-on: ubuntu-latest strategy: matrix: - go: [ '1.25' ] + go: ["1.26"] steps: - uses: actions/checkout@v3 diff --git a/.golangci.yml b/.golangci.yml index 2da416b..f2b16a7 100755 --- a/.golangci.yml +++ b/.golangci.yml @@ -1,7 +1,7 @@ version: "2" run: - go: "1.25" + go: "1.26" timeout: 5m tests: false issues-exit-code: 1 diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..8bd7e9a --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,28 @@ +# Agent instructions + +## Repository map + +- The root module is `go.osspkg.com/do`; its source and tests live at the repository root. +- `mapreduce`, `monad`, `pipline`, `words`, and `workerpool` are separate importable packages. +- `pipline` is the existing package and import path spelling. Preserve it unless a deliberate compatibility change is requested. +- `README.md` is the user-facing overview; keep examples and documented behavior aligned with the exported API. +- For usage questions or changes that use this module, read `skills/go-do-usage/SKILL.md` and the relevant linked reference. + +## Go version and commands + +Run commands from the repository root. `go.mod` requires Go 1.26.0. + +- Run all tests with `go test ./...`. +- For changes to the core helpers, worker pool, or map-reduce package, run `go test . ./workerpool ./mapreduce` first. +- Format changed Go files with `gofmt -w `. +- The GitHub Actions workflow is `.github/workflows/ci.yml`; it invokes `make ci`. +- `make ci` runs `make pre-commit`, which installs the latest `goppy` tool, runs `goppy setup-lib`, and then runs license, lint, test, and build targets. Review those setup and generated-file side effects before running it locally. +- `go.mod` requires Go 1.26.0. + +## Change guidance + +- Add regression tests for changed behavior. For concurrency changes, cover cancellation, shutdown, and channel ownership; use bounded waits in tests that could hang. +- Keep package boundaries intact. The root package contains generic slice/map helpers, conditionals, numeric helpers, panic recovery, and async helpers; the subpackages provide focused APIs. +- Update `README.md` when installation, public behavior, or usage examples change. +- Avoid changing exported names or import paths without considering downstream compatibility. +- Do not run publication, deployment, or release commands as part of local validation. diff --git a/LICENSE b/LICENSE index 2a1ad25..d6c3bf3 100644 --- a/LICENSE +++ b/LICENSE @@ -1,6 +1,6 @@ BSD 3-Clause License -Copyright (c) 2024-2025, Mikhail Knyazhev +Copyright (c) 2024-2026, Mikhail Knyazhev Redistribution and use in source and binary forms, with or without modification, are permitted provided that the following conditions are met: diff --git a/README.md b/README.md index fd3e276..597c75f 100644 --- a/README.md +++ b/README.md @@ -1 +1,138 @@ -# go-do \ No newline at end of file +# go-do + +[![Go version](https://img.shields.io/github/go-mod/go-version/osspkg/go-do)](https://go.dev/doc/install) +[![CI](https://github.com/osspkg/go-do/actions/workflows/ci.yml/badge.svg?branch=master)](https://github.com/osspkg/go-do/actions/workflows/ci.yml) +[![Go Reference](https://pkg.go.dev/badge/go.osspkg.com/do.svg)](https://pkg.go.dev/go.osspkg.com/do) +[![License](https://img.shields.io/github/license/osspkg/go-do)](LICENSE) + +`go-do` is a Go utility library with generic helpers for slices and maps, numeric operations, conditional expressions, panic recovery, and asynchronous work. It also provides focused packages for worker pools, map-reduce, state pipelines, result handling, and word tokenization. + +## Requirements + +- Go 1.26.0 or newer, as declared in [`go.mod`](go.mod). + +## Installation + +```sh +go get go.osspkg.com/do +``` + +Import the root package and any subpackage you need: + +```go +import ( + "go.osspkg.com/do" + "go.osspkg.com/do/mapreduce" + "go.osspkg.com/do/workerpool" +) +``` + +## Quick start + +```go +package main + +import ( + "fmt" + + "go.osspkg.com/do" +) + +func main() { + values := do.Filter([]int{1, 2, 3, 4, 5}, func(value, index int) bool { + return value%2 == 1 + }) + labels := do.Convert(values, func(value, index int) string { + return fmt.Sprintf("item-%d", value) + }) + + fmt.Println(labels) // [item-1 item-3 item-5] +} +``` + +## Packages and features + +### Root package: `go.osspkg.com/do` + +| Area | Functions and types | Notes | +| --- | --- | --- | +| Slices | `Each`, `Convert`, `Join`, `Chunk`, `Entries`, `Reduce`, `Filter`, `Treat`, `TreatValue`, `Diff`, `Unique`, `IndexOf`, `LastIndexOf`, `Include`, `Exclude`, `Copy`, `Splice`, `Pop`, `Push`, `Shift`, `Unshift`, `Reverse`, `ToMap` | `Splice`, `Push`, `Unshift`, `Pop`, `Shift`, and `Reverse` mutate the supplied slice or slice pointer. | +| Maps | `EachMap`, `ConvertMap`, `FilterMap`, `TreatMap`, `TreatMapValue`, `Keys`, `Values`, `JoinMap`, `FlipMap`, `DivideMap`, `CombineMap`, `ReduceMap`, `ToSlice` | `Keys`, `Values`, `DivideMap`, `ReduceMap`, and `ToSlice` use sorted keys. Other map transformations may follow Go's unspecified map iteration order. | +| Numeric helpers | `MinMax`, `MinMaxTime`, `Range`, `Sum`, `Average` | `Range` includes both endpoints when reached and requires a positive, progressing step; otherwise it returns an empty slice or stops when the value cannot advance. | +| Conditional helpers | `If`, `IfFunc`, `IfElse`, `IfElseFunc`, `DoIf` | `DoIf.ElseIf`, `Else`, `ElseIfFunc`, and `ElseFunc` form a conditional chain; function variants evaluate only the selected branch. | +| Panic and async helpers | `Recovery`, `Trace`, `Try`, `Async`, `AsyncGroup` | `AsyncGroup` waits for all supplied functions and returns their errors. | + +The `Summable` and `Comparable` type constraints define the numeric and ordered types accepted by the generic helpers. + +### `go.osspkg.com/do/workerpool` + +A bounded worker pool with generic task and result types. + +| API | Use | +| --- | --- | +| `Task[I, T]`, `Result[I, T]` | Carry a task ID and input data, or a task ID, result value, and error. | +| `Pool[I, T, R]` | Own the task and result channels and worker lifecycle. | +| `New[I, T, R]` | Create a pool with a fixed worker count and task handler. | +| `Pool.Start` | Start the workers; repeated calls have no effect. | +| `Pool.Send` | Submit a task; returns `false` after shutdown or when the pool is closing. | +| `Pool.Receive` | Get the result channel. | +| `Pool.Close` | Cancel workers and close the result channel; safe to call more than once. | + +Read results while work is running to avoid filling the bounded result buffer during normal processing. + +### `go.osspkg.com/do/mapreduce` + +| API | Use | +| --- | --- | +| `New[T, R, O]` | Map items concurrently and pass completed results to a reducer. | + +Set `workers` to a positive number. Mapping and reducing order is nondeterministic when more than one worker is used, so reducers should not depend on input order. A reducer error cancels the mapping context and waits for mapper goroutines to return; mapper functions should honor their context. + +### `go.osspkg.com/do/pipline` + +| API | Use | +| --- | --- | +| `Pipe[S, T]` | Define a state transition and its context-aware handler. | +| `Pipline[S, T]` | Store transitions keyed by their current state. | +| `New[S, T]` | Create an empty pipeline. | +| `Pipline.Set` | Register a transition; returns an error for a nil handler, duplicate current state, or identical current and next states. | +| `Pipline.Do` | Run handlers from a starting state until no transition exists, the context is canceled, or a handler returns an error. | + +### `go.osspkg.com/do/monad` + +| API | Use | +| --- | --- | +| `Result[T]` | Hold a value and an error for chained operations. | +| `Some[T]` | Create a successful result from a value. | +| `Bind[T, R]` | Apply a function to a successful result and propagate errors. | +| `Result.Return` | Retrieve the result's value and error. | + +### `go.osspkg.com/do/words` + +| API | Use | +| --- | --- | +| `Words` | Interface for tokenizing to strings or byte slices and configuring token categories. | +| `NewString`, `NewBytes` | Create a tokenizer with default rune detectors. | +| `Strings`, `Bytes` | Tokenize a string or byte slice with the default block and symbol detectors. | +| `Words.Strings`, `Words.Bytes` | Tokenize using the instance's current detector configuration. | +| `UseDefaultBlock`, `SetBlock` | Select default or custom runes that form word blocks. | +| `UseDefaultDigital`, `SetDigital` | Select default or custom digit runes. | +| `UseDefaultSymbol`, `SetSymbol` | Select default or custom symbol runes. | + +## Testing + +Run the complete test suite from the repository root: + +```sh +go test ./... +``` + +The project Makefile also provides `make tests`, `make lint`, `make build`, and `make ci`. `make ci` installs the latest `goppy` tool and runs setup, license, lint, test, and build targets; see the [Makefile](Makefile) before running it locally. + +## Contributing + +Pull requests are welcome. Include tests for behavior changes and run `go test ./...` before submitting. Keep public API and import-path compatibility in mind. + +## License + +This project is licensed under the [BSD 3-Clause License](LICENSE). diff --git a/async.go b/async.go index 46019ff..1b8a75f 100644 --- a/async.go +++ b/async.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/async_test.go b/async_test.go index 5bc2d68..f2966db 100644 --- a/async_test.go +++ b/async_test.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/common.go b/common.go index 590d4bf..4b90815 100644 --- a/common.go +++ b/common.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/go.mod b/go.mod index e4afa5d..72a2d88 100644 --- a/go.mod +++ b/go.mod @@ -1,8 +1,8 @@ module go.osspkg.com/do -go 1.25.0 +go 1.26.0 require ( - go.osspkg.com/casecheck v0.3.0 - golang.org/x/sync v0.21.0 + go.osspkg.com/casecheck v0.3.2 + golang.org/x/sync v0.23.0 ) diff --git a/go.sum b/go.sum index 4b03aa7..14afa29 100644 --- a/go.sum +++ b/go.sum @@ -1,4 +1,4 @@ -go.osspkg.com/casecheck v0.3.0 h1:x15blEszElbrHrEH5H02JIIhGIg/lGZzIt1kQlD3pwM= -go.osspkg.com/casecheck v0.3.0/go.mod h1:TRFXDMFJEOtnlp3ET2Hix3osbxwPWhvaiT/HfD3+gBA= -golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM= -golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +go.osspkg.com/casecheck v0.3.2 h1:KDdtEsEnGcDKjtg8FKL0hjnVZ7kusmXBD/VS/5IKSjE= +go.osspkg.com/casecheck v0.3.2/go.mod h1:nf1vimi3VPl1o0hV+bKZsy3+1Qi8wNRRWiNGqivFV1A= +golang.org/x/sync v0.23.0 h1:KameEIfc1IkluZyXWLn39Wd4tURc6GbCiISGiZm2bQk= +golang.org/x/sync v0.23.0/go.mod h1:sUUOizhqBxiL6pEWpqNLUiaJn1ShEbZ6BBqskPbjZm0= diff --git a/if.go b/if.go index ac02859..2fba9bd 100644 --- a/if.go +++ b/if.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/if_test.go b/if_test.go index 19923e8..b5b818e 100644 --- a/if_test.go +++ b/if_test.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/map.go b/map.go index 630bcaa..0ee2e04 100644 --- a/map.go +++ b/map.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/map_test.go b/map_test.go index 33d2490..696e63d 100644 --- a/map_test.go +++ b/map_test.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/mapreduce/mapreduce.go b/mapreduce/mapreduce.go index 95f6055..e2380ae 100644 --- a/mapreduce/mapreduce.go +++ b/mapreduce/mapreduce.go @@ -1,7 +1,13 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package mapreduce import ( "context" + "errors" "golang.org/x/sync/errgroup" ) @@ -14,7 +20,13 @@ func New[T, R, O any]( initial O, workers int, ) (O, error) { - g, ctx := errgroup.WithContext(ctx) + if workers <= 0 { + return initial, errors.New("workers must be greater than zero") + } + + workCtx, cancel := context.WithCancel(ctx) + defer cancel() + g, ctx := errgroup.WithContext(workCtx) g.SetLimit(workers) results := make(chan R, len(items)) @@ -45,6 +57,8 @@ func New[T, R, O any]( var err error acc, err = reducer(ctx, acc, r) if err != nil { + cancel() + _ = g.Wait() return acc, err } } diff --git a/mapreduce/mapreduce_test.go b/mapreduce/mapreduce_test.go index 6003425..662ca80 100644 --- a/mapreduce/mapreduce_test.go +++ b/mapreduce/mapreduce_test.go @@ -1,8 +1,14 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package mapreduce import ( "context" "errors" + "fmt" "sync" "testing" "time" @@ -256,3 +262,46 @@ func TestUnit_MapReduceWaitForAllGoroutines(t *testing.T) { t.Errorf("expected %d results, got %d", len(items), len(result)) } } + +func TestUnit_MapReduceRejectsNonPositiveWorkers(t *testing.T) { + mapper := func(ctx context.Context, value int) (int, error) { return value, nil } + reducer := func(ctx context.Context, acc, value int) (int, error) { return acc + value, nil } + + for _, workers := range []int{0, -1} { + t.Run(fmt.Sprintf("workers=%d", workers), func(t *testing.T) { + got, err := New(context.Background(), []int{1}, mapper, reducer, 7, workers) + if err == nil { + t.Fatal("expected an error for non-positive workers") + } + if got != 7 { + t.Fatalf("initial value: got %d, want 7", got) + } + }) + } +} + +func TestUnit_MapReduceReducerErrorCancelsMappers(t *testing.T) { + wantErr := errors.New("reducer error") + canceled := make(chan struct{}) + mapper := func(ctx context.Context, value int) (int, error) { + if value == 1 { + return value, nil + } + <-ctx.Done() + close(canceled) + return 0, ctx.Err() + } + reducer := func(ctx context.Context, acc, value int) (int, error) { + return acc, wantErr + } + + _, err := New(context.Background(), []int{1, 2}, mapper, reducer, 0, 2) + if !errors.Is(err, wantErr) { + t.Fatalf("got error %v, want %v", err, wantErr) + } + select { + case <-canceled: + case <-time.After(time.Second): + t.Fatal("mapper did not observe cancellation") + } +} diff --git a/math.go b/math.go index 9b1a682..3b47d7c 100644 --- a/math.go +++ b/math.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ @@ -29,10 +29,18 @@ func MinMaxTime(value, minimum, maximum time.Time) time.Time { func Range[T Summable](from, to, step T) (out []T) { out = make([]T, 0, 2) - for i := from; i <= to; i += step { + if step <= 0 { + return out + } + for i := from; i <= to; { out = append(out, i) + next := i + step + if next <= i { + return out + } + i = next } - return + return out } func Sum[T Comparable](elements ...T) (out T) { diff --git a/math_test.go b/math_test.go index ce57e36..cd9d66b 100644 --- a/math_test.go +++ b/math_test.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ @@ -37,6 +37,13 @@ func TestUnit_Range(t *testing.T) { casecheck.Equal(t, []byte{'a', 'b', 'c'}, do.Range[byte]('a', 'c', 1)) } +func TestUnit_RangeInvalidOrNonProgressingStep(t *testing.T) { + casecheck.Equal(t, []int{}, do.Range[int](1, 10, 0)) + casecheck.Equal(t, []int{}, do.Range[int](1, 10, -1)) + casecheck.Equal(t, []int8{126, 127}, do.Range[int8](126, 127, 1)) + casecheck.Equal(t, []float64{1e20}, do.Range[float64](1e20, 1e21, 1)) +} + func TestUnit_MinMaxTime(t *testing.T) { mi := time.Now() ma := time.Now().Add(1 * time.Hour) diff --git a/monad/monad.go b/monad/monad.go index 072e8cd..90268cd 100644 --- a/monad/monad.go +++ b/monad/monad.go @@ -1,3 +1,8 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package monad type Result[T any] struct { diff --git a/monad/monad_test.go b/monad/monad_test.go index c8fea7b..8f98511 100644 --- a/monad/monad_test.go +++ b/monad/monad_test.go @@ -1,3 +1,8 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package monad_test import ( @@ -5,6 +10,7 @@ import ( "testing" "go.osspkg.com/casecheck" + . "go.osspkg.com/do/monad" ) diff --git a/pipline/pipline.go b/pipline/pipline.go index 3a8d811..f3d2999 100644 --- a/pipline/pipline.go +++ b/pipline/pipline.go @@ -1,3 +1,8 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package pipline import ( diff --git a/pipline/pipline_test.go b/pipline/pipline_test.go index a5d3a6c..0c0adf0 100644 --- a/pipline/pipline_test.go +++ b/pipline/pipline_test.go @@ -1,3 +1,8 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package pipline import ( diff --git a/skills/go-do-usage/SKILL.md b/skills/go-do-usage/SKILL.md new file mode 100644 index 0000000..9c144e2 --- /dev/null +++ b/skills/go-do-usage/SKILL.md @@ -0,0 +1,31 @@ +--- +name: go-do-usage +description: Choose and use APIs from the go.osspkg.com/do Go utility library. Use when writing, reviewing, or explaining Go code that imports this module; skip for Go work that does not use go-do. +metadata: + short-description: Use the go-do Go library +--- + +# Using go-do + +Use the library's generic helpers and focused subpackages when they match the caller's needs. Preserve their actual mutation, ordering, error, and lifecycle behavior; do not infer semantics from familiar function names alone. + +## Workflow + +1. Confirm the importing module uses `go.osspkg.com/do` and inspect its `go.mod` for the required Go version. +2. Read [the API guide](references/api-guide.md) to choose a package and understand behavioral constraints. +3. Read [the examples](references/examples.md) for the relevant operation, then verify details against the current implementation or `go doc` if code and reference may have drifted. +4. Add or update tests for behavior changes using the consuming repository's conventions. + +## Important usage constraints + +- Import subpackages by their exact module paths. The state pipeline package is spelled `pipline` in both its package name and import path. +- Root helpers generally return new slices or maps, but slice operations such as `Reverse`, `Push`, `Unshift`, `Pop`, `Shift`, and `Splice` mutate their input. Check the API guide before assuming ownership or aliasing behavior. +- Go map iteration is unordered. Helpers that sort keys provide stable key order; do not assume that concurrent map-reduce preserves input order. +- Use a positive worker count for `mapreduce.New`, and make mapper functions honor their context. A reducer error cancels that context and waits for mapper functions to return. +- For `workerpool`, start the pool once, handle every expected result while tasks are running, and call `Close` to cancel workers and close the result channel. A handler that needs prompt shutdown must honor its context. +- `do.Range` is ascending: use a positive step that advances the value. Non-positive steps produce an empty slice; integer overflow or floating-point precision that prevents progress stops the range. + +## References + +- [API guide](references/api-guide.md): package map, selection guidance, and behavior notes. +- [Examples](references/examples.md): focused snippets for collection helpers, map-reduce, worker pools, pipelines, words, and result chaining. diff --git a/skills/go-do-usage/agents/openai.yaml b/skills/go-do-usage/agents/openai.yaml new file mode 100644 index 0000000..64a88c2 --- /dev/null +++ b/skills/go-do-usage/agents/openai.yaml @@ -0,0 +1,3 @@ +interface: + display_name: "Go Do Usage" + short_description: "Help with Go Do Usage tasks" diff --git a/skills/go-do-usage/references/api-guide.md b/skills/go-do-usage/references/api-guide.md new file mode 100644 index 0000000..89bd5b4 --- /dev/null +++ b/skills/go-do-usage/references/api-guide.md @@ -0,0 +1,75 @@ +# go-do API guide + +This guide describes the repository's current public API. Check the source or `go doc` when behavior is important and this reference may be stale. The module path is `go.osspkg.com/do`; `go.mod` declares Go 1.26.0. + +## Choose a package + +| Package | Use it for | +| --- | --- | +| `go.osspkg.com/do` | Generic slice and map operations, simple numeric helpers, conditional expressions, panic recovery, and asynchronous calls. | +| `go.osspkg.com/do/mapreduce` | Concurrent mapping followed by reduction of completed results. | +| `go.osspkg.com/do/workerpool` | Reusable bounded workers that accept identified tasks and publish results on a channel. | +| `go.osspkg.com/do/pipline` | Context-aware handlers connected by comparable state values. Preserve the existing spelling `pipline`. | +| `go.osspkg.com/do/monad` | A small generic value-plus-error `Result` with `Bind`. | +| `go.osspkg.com/do/words` | Split strings or byte slices into tokens using default or custom rune detectors. | + +## Root package + +### Slices + +- Iteration and conversion: `Each`, `Convert`, `Reduce`, `Filter`, `Treat`, `TreatValue`, `Entries`, `ToMap`. +- Composition and grouping: `Join`, `Chunk`, `Diff`, `Unique`, `Exclude`. +- Search: `IndexOf`, `LastIndexOf`, `Include`. +- Mutation: `Push`, `Pop`, `Shift`, `Unshift`, `Reverse`, `Splice`. These mutate the supplied slice or slice pointer. `Copy` returns a new slice; its `start` and `end` indexes are inclusive and clamped to the input bounds. + +`Reduce` starts with the element type's zero value; it does not accept a separate initial accumulator. Empty inputs return zero or an empty result according to the function. + +### Maps + +- Traversal and transformation: `EachMap`, `ConvertMap`, `FilterMap`, `TreatMap`, `TreatMapValue`, `ReduceMap`, `ToSlice`. +- Key/value operations: `Keys`, `Values`, `JoinMap`, `FlipMap`, `DivideMap`, `CombineMap`. + +`Keys`, `Values`, `DivideMap`, `ReduceMap`, and `ToSlice` use sorted keys, constrained by the library's `Comparable` type set. Other helpers that range directly over a map do not promise iteration order. Key collisions in conversion, joining, flipping, or combining use the last value assigned during iteration; where input map iteration is involved, which colliding entry wins may be nondeterministic. + +### Numbers and conditionals + +- `MinMax` and `MinMaxTime` clamp a value to a minimum and maximum. +- `Range` generates an ascending inclusive sequence using a positive step. It returns an empty slice for non-positive steps and stops if adding the step does not increase the current value. +- `Sum` adds values supported by `Comparable`; `Average` supports `Summable` numeric types and uses the type's integer or floating-point division semantics. Calling `Average` with no values divides by zero for integer types. +- `If`, `IfElse`, and their `Func` variants provide conditional selection. Function arguments are useful for lazy evaluation. `DoIf` supports an `ElseIf`/`Else` chain and lazy `ElseIfFunc`/`ElseFunc` branches. + +`Summable` and `Comparable` are generic type constraints intended for these APIs, not runtime validators. + +### Panic and asynchronous helpers + +- `Recovery` runs a function and returns a recovered panic as an error; it does not recover panics in other goroutines. +- `Try` runs optional try/catch/finally callbacks. Panics raised by callbacks are recovered by its internal wrappers. +- `Trace` formats runtime call frames. +- `Async` starts one function in a goroutine and optionally reports a panic through the error callback. It does not provide a join handle. +- `AsyncGroup` starts one goroutine per supplied function, waits for all of them, and returns their errors. It passes the supplied context but does not itself cancel sibling functions after one fails; use a bounded worker pool for large workloads. + +## `mapreduce` + +`mapreduce.New(ctx, items, mapper, reducer, initial, workers)` allocates a result channel sized to `len(items)`, runs mapper calls with an `errgroup` concurrency limit, then reduces results as they arrive. Large input slices therefore also require a potentially large result buffer. `workers` must be greater than zero. With multiple workers, completion and reduction order are nondeterministic; reducers that require input order should use one worker or an order-preserving design. Mapper functions should observe context cancellation. A reducer error cancels the mapper context and waits for mapper goroutines to return. + +## `workerpool` + +- `New(workers, handler)` creates a pool; non-positive worker counts are clamped to one. +- `Start` starts workers once; repeated calls do nothing. +- `Send` returns whether the task was accepted before shutdown. It returns `false` after closure or cancellation. +- `Receive` exposes a result channel buffered to the worker count. +- `Close` cancels workers, waits for senders and workers, and closes the result channel. It is safe to call repeatedly. + +The handler receives the pool context. `Close` waits for active handlers, so handlers that need responsive shutdown should return when the context is canceled. Continue receiving results during normal work because the result buffer is bounded; cancellation unblocks workers that cannot publish another result. + +## `pipline` + +Create a `Pipline[S, T]` with `New`, register `Pipe[S, T]` transitions using `Set`, and execute using `Do(ctx, state, value)`. `Set` rejects a nil handler, duplicate current states, and transitions whose current and next states are equal. `Do` follows transitions until a state has no registered handler, the context is canceled, or a handler returns an error. + +## `monad` + +`Some(value)` creates a successful `Result[T]`. `Bind(result, next)` skips `next` if the result already holds an error; otherwise it stores the returned value or error. `Result.Return()` retrieves the pair. The contained fields are private, so construct values with `Some` and `Bind`. + +## `words` + +`Strings(s)` and `Bytes(b)` use the default detectors. `NewString` and `NewBytes` return a `Words` interface that can be configured with `SetBlock`, `SetDigital`, and `SetSymbol`, or reset with the matching `UseDefault...` methods. `Words.Strings()` and `Words.Bytes()` tokenize from the beginning on each call. The default block detector groups letters and digits; digit and symbol detectors control how adjacent tokens are grouped. Tokenization uses `bufio.Scanner` with its default token limit and ignores scanner errors, so very long tokens can cause incomplete output without an error being returned. diff --git a/skills/go-do-usage/references/examples.md b/skills/go-do-usage/references/examples.md new file mode 100644 index 0000000..87c6eb1 --- /dev/null +++ b/skills/go-do-usage/references/examples.md @@ -0,0 +1,145 @@ +# go-do examples + +The examples use the module imports from `go.osspkg.com/do`. Each code block is a function intended to be copied into a Go package together with the listed imports. + +## Filter and convert a slice + +```go +func filterLabels() []string { + values := do.Filter([]int{1, 2, 3, 4, 5}, func(value, index int) bool { + return value%2 == 1 + }) + return do.Convert(values, func(value, index int) string { + return fmt.Sprintf("item-%d", value) + }) +} +``` + +Imports: `fmt` and `go.osspkg.com/do`. + +## Transform a map with stable key order + +```go +func formatScores() []string { + scores := map[string]int{"Ada": 10, "Lin": 8} + return do.ToSlice(scores, func(score int, name string) string { + return fmt.Sprintf("%s=%d", name, score) + }) +} +``` + +Imports: `fmt` and `go.osspkg.com/do`. `ToSlice` visits map entries in sorted key order. + +## Run map-reduce work + +```go +func sumSquares() (int, error) { + ctx := context.Background() + total, err := mapreduce.New( + ctx, + []int{1, 2, 3, 4}, + func(ctx context.Context, value int) (int, error) { + return value * value, nil + }, + func(ctx context.Context, total, value int) (int, error) { + return total + value, nil + }, + 0, + 4, + ) + if err != nil { + return 0, err + } + return total, nil // 30; completion order is unspecified with multiple workers. +} +``` + +Imports: `context` and `go.osspkg.com/do/mapreduce`. + +## Use a worker pool + +```go +func runPool() error { + pool := workerpool.New(2, func(ctx context.Context, task workerpool.Task[int, int]) (int, error) { + return task.Data * 2, nil + }) + pool.Start() + defer pool.Close() + + tasks := []int{10, 20, 30} + for id, value := range tasks { + if !pool.Send(workerpool.Task[int, int]{ID: id, Data: value}) { + return errors.New("worker pool is closed") + } + } + + for range tasks { + result := <-pool.Receive() + if result.Err != nil { + return result.Err + } + fmt.Println(result.TaskID, result.Result) + } + return nil +} +``` + +Imports: `context`, `errors`, `fmt`, and `go.osspkg.com/do/workerpool`. Read as many results as tasks submitted; result arrival order is not guaranteed. + +## Chain states through a pipeline + +```go +func runPipeline() (string, error) { + const ( + start = iota + done + ) + + pipe := pipline.New[int, string]() + err := pipe.Set(pipline.Pipe[int, string]{ + Current: start, + Next: done, + Handler: func(ctx context.Context, value string) (string, error) { + return value + "!", nil + }, + }) + if err != nil { + return "", err + } + + state, value, err := pipe.Do(context.Background(), start, "hello") + if err != nil { + return "", err + } + if state != done { + return "", errors.New("pipeline stopped before the done state") + } + return value, nil // "hello!" +} +``` + +Imports: `context`, `errors`, and `go.osspkg.com/do/pipline`. + +## Configure word tokenization + +```go +func tokenize() []string { + tokenizer := words.NewString("Hello, go-do!") + tokenizer.SetBlock(unicode.IsLetter) + return tokenizer.Strings() +} +``` + +Imports: `unicode` and `go.osspkg.com/do/words`. The interface also supports custom digit and symbol detectors. + +## Chain a value and error + +```go +func parseNumber() (int, error) { + result := monad.Some("42") + parsed := monad.Bind(result, strconv.Atoi) + return parsed.Return() +} +``` + +Imports: `strconv` and `go.osspkg.com/do/monad`. diff --git a/slice.go b/slice.go index 4bf9cde..aa0b131 100644 --- a/slice.go +++ b/slice.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/slice_test.go b/slice_test.go index abbc1f4..a7beee7 100644 --- a/slice_test.go +++ b/slice_test.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/trace.go b/trace.go index 253e495..c39c965 100644 --- a/trace.go +++ b/trace.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/try.go b/try.go index d7dba8b..a957405 100644 --- a/try.go +++ b/try.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/try_test.go b/try_test.go index cccdb5a..9b0f46e 100644 --- a/try_test.go +++ b/try_test.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/words/strings_words.go b/words/strings_words.go index 8aae3ff..c9c27a6 100644 --- a/words/strings_words.go +++ b/words/strings_words.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ diff --git a/words/strings_words_test.go b/words/strings_words_test.go index 8a910f2..7f4194e 100644 --- a/words/strings_words_test.go +++ b/words/strings_words_test.go @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Mikhail Knyazhev . All rights reserved. + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. */ @@ -10,6 +10,7 @@ import ( "unicode" "go.osspkg.com/casecheck" + "go.osspkg.com/do/words" ) diff --git a/workerpool/workerpool.go b/workerpool/workerpool.go index bf6cce2..25ffe6e 100644 --- a/workerpool/workerpool.go +++ b/workerpool/workerpool.go @@ -1,3 +1,8 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package workerpool import ( @@ -17,13 +22,18 @@ type Result[I comparable, T any] struct { } type Pool[I comparable, T, R any] struct { - workers int - tasksC chan Task[I, T] - resultsC chan Result[I, R] - handler func(ctx context.Context, task Task[I, T]) (R, error) - ctx context.Context - cancel context.CancelFunc - wg sync.WaitGroup + workers int + tasksC chan Task[I, T] + resultsC chan Result[I, R] + handler func(ctx context.Context, task Task[I, T]) (R, error) + ctx context.Context + cancel context.CancelFunc + wg sync.WaitGroup + sendWg sync.WaitGroup + mu sync.Mutex + closeOnce sync.Once + started bool + closed bool } func New[I comparable, T, R any]( @@ -45,13 +55,26 @@ func New[I comparable, T, R any]( } func (p *Pool[I, T, R]) Close() { - p.cancel() - p.wg.Wait() - close(p.tasksC) - close(p.resultsC) + p.closeOnce.Do(func() { + p.mu.Lock() + p.closed = true + p.cancel() + p.mu.Unlock() + + p.sendWg.Wait() + p.wg.Wait() + close(p.resultsC) + }) } func (p *Pool[I, T, R]) Start() { + p.mu.Lock() + defer p.mu.Unlock() + if p.started || p.closed { + return + } + p.started = true + p.wg.Add(p.workers) for i := 0; i < p.workers; i++ { go func() { @@ -61,10 +84,14 @@ func (p *Pool[I, T, R]) Start() { select { case task := <-p.tasksC: result, err := p.handler(p.ctx, task) - p.resultsC <- Result[I, R]{ + select { + case p.resultsC <- Result[I, R]{ TaskID: task.ID, Result: result, Err: err, + }: + case <-p.ctx.Done(): + return } case <-p.ctx.Done(): return @@ -75,15 +102,18 @@ func (p *Pool[I, T, R]) Start() { } func (p *Pool[I, T, R]) Send(task Task[I, T]) bool { - select { - case <-p.ctx.Done(): + p.mu.Lock() + if p.closed { + p.mu.Unlock() return false - default: } + p.sendWg.Add(1) + p.mu.Unlock() + defer p.sendWg.Done() select { case p.tasksC <- task: - return true + return p.ctx.Err() == nil case <-p.ctx.Done(): return false } diff --git a/workerpool/workerpool_test.go b/workerpool/workerpool_test.go index 66f5264..6e98e32 100644 --- a/workerpool/workerpool_test.go +++ b/workerpool/workerpool_test.go @@ -1,3 +1,8 @@ +/* + * Copyright (c) 2024-2026 Mikhail Knyazhev . All rights reserved. + * Use of this source code is governed by a BSD 3-Clause license that can be found in the LICENSE file. + */ + package workerpool import ( @@ -210,3 +215,73 @@ func TestUnit_Pool_Concurrency(t *testing.T) { p.Close() } + +func TestUnit_Pool_CloseWithFullResultsBuffer(t *testing.T) { + started := make(chan struct{}, 2) + p := New(1, func(ctx context.Context, task Task[int, int]) (int, error) { + started <- struct{}{} + return task.Data, nil + }) + p.Start() + + if !p.Send(Task[int, int]{ID: 1, Data: 1}) { + t.Fatal("first Send failed") + } + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("first handler did not start") + } + deadline := time.Now().Add(time.Second) + for len(p.resultsC) == 0 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + if len(p.resultsC) == 0 { + t.Fatal("first result was not buffered") + } + + if !p.Send(Task[int, int]{ID: 2, Data: 2}) { + t.Fatal("second Send failed") + } + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("second handler did not start") + } + + closed := make(chan struct{}) + go func() { + p.Close() + close(closed) + }() + select { + case <-closed: + case <-time.After(time.Second): + t.Fatal("Close blocked with a full results buffer") + } +} + +func TestUnit_Pool_SendConcurrentWithClose(t *testing.T) { + p := New(2, simpleHandler) + p.Start() + + var senders sync.WaitGroup + for i := 0; i < 8; i++ { + senders.Add(1) + go func(id int) { + defer senders.Done() + for j := 0; j < 100; j++ { + if !p.Send(Task[int, int]{ID: id*100 + j, Data: j + 1}) { + return + } + } + }(i) + } + + p.Close() + senders.Wait() + p.Close() + if p.Send(Task[int, int]{ID: 1000, Data: 1}) { + t.Fatal("Send succeeded after Close") + } +}