feat(m-operation-event-outbox): operation pick task completion and worker fixes

- Move operation pick artifacts to archive
- Fix worker runner with outbox pattern support
- Update storage tests for postgres integration
- Add worker runner tests
This commit is contained in:
toki 2026-06-14 05:18:10 +09:00
parent 51c8098577
commit 3ed8797b28
8 changed files with 386 additions and 34 deletions

View file

@ -42,42 +42,43 @@ task=m-operation-event-outbox/02+01_operation_pick, plan=0, tag=API
| 항목 | 완료 여부 |
|------|---------|
| [API-1] Atomic pending operation pick | [ ] |
| [API-2] Worker run-one pick boundary | [ ] |
| [API-1] Atomic pending operation pick | [x] |
| [API-2] Worker run-one pick boundary | [x] |
## 구현 체크리스트
- [ ] `OperationStore`에 pending pick API를 추가하고 `pgOperationStore`에서 `FOR UPDATE SKIP LOCKED` 기반으로 queued operation을 running으로 전환한다.
- [ ] worker runner에 one-iteration pick boundary를 추가하고 picker가 없거나 worker disabled일 때의 동작을 테스트한다.
- [ ] concurrent pick 테스트를 작성한다. 검증: 같은 operation이 중복 pick되지 않는다.
- [ ] `go test -count=1 ./internal/storage ./internal/worker`와 `go test -count=1 ./...`를 실행한다.
- [ ] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다.
- [x] `OperationStore`에 pending pick API를 추가하고 `pgOperationStore`에서 `FOR UPDATE SKIP LOCKED` 기반으로 queued operation을 running으로 전환한다.
- [x] worker runner에 one-iteration pick boundary를 추가하고 picker가 없거나 worker disabled일 때의 동작을 테스트한다.
- [x] concurrent pick 테스트를 작성한다. 검증: 같은 operation이 중복 pick되지 않는다.
- [x] `go test -count=1 ./internal/storage ./internal/worker`와 `go test -count=1 ./...`를 실행한다.
- [x] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다.
## 코드리뷰 전용 체크리스트
> **[REVIEW AGENT ONLY]** 이 체크리스트는 코드리뷰 에이전트만 사용한다.
> 구현 에이전트는 이 섹션을 수정하거나 체크하지 않는다.
- [ ] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다.
- [ ] 판정과 `차원별 평가`, Required/Suggested/Nit 분류가 서로 일치한다.
- [ ] active `CODE_REVIEW-*-G??.md`를 `code_review_{review_lane}_GNN_N.log`로 아카이브한다.
- [ ] active `PLAN-*-G??.md`를 `plan_{build_lane}_GNN_M.log`로 아카이브한다.
- [ ] `.gitignore`의 Agent-Ops 관리 block이 `agent-task/**/*.md`와 `agent-task/**/*.log`를 unignore하고 `agent-roadmap/current.md`를 ignore하는지 확인한다.
- [ ] PASS이면 `agent-ops/skills/common/code-review/templates/complete-log-template.md` 기준으로 `complete.log`를 작성하고 active `.md` 파일을 남기지 않는다.
- [ ] PASS이면 active task 디렉터리 `agent-task/m-operation-event-outbox/02+01_operation_pick/`를 `agent-task/archive/YYYY/MM/m-operation-event-outbox/02+01_operation_pick/`로 이동하고 최종 archive 경로에서 이 체크리스트를 갱신한다.
- [ ] PASS이고 task group이 `m-operation-event-outbox`이면 런타임이 읽을 완료 이벤트 메타데이터를 보고하고, roadmap 수정이나 `update-roadmap` 직접 호출을 하지 않는다.
- [ ] PASS split 작업이면 이동 후 빈 active parent `agent-task/m-operation-event-outbox/`를 제거하거나, 남은 sibling/file이 있어 유지했다고 확인한다.
- [x] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다.
- [x] 판정과 `차원별 평가`, Required/Suggested/Nit 분류가 서로 일치한다.
- [x] active `CODE_REVIEW-*-G??.md`를 `code_review_{review_lane}_GNN_N.log`로 아카이브한다.
- [x] active `PLAN-*-G??.md`를 `plan_{build_lane}_GNN_M.log`로 아카이브한다.
- [x] `.gitignore`의 Agent-Ops 관리 block이 `agent-task/**/*.md`와 `agent-task/**/*.log`를 unignore하고 `agent-roadmap/current.md`를 ignore하는지 확인한다.
- [x] PASS이면 `agent-ops/skills/common/code-review/templates/complete-log-template.md` 기준으로 `complete.log`를 작성하고 active `.md` 파일을 남기지 않는다.
- [x] PASS이면 active task 디렉터리 `agent-task/m-operation-event-outbox/02+01_operation_pick/`를 `agent-task/archive/YYYY/MM/m-operation-event-outbox/02+01_operation_pick/`로 이동하고 최종 archive 경로에서 이 체크리스트를 갱신한다.
- [x] PASS이고 task group이 `m-operation-event-outbox`이면 런타임이 읽을 완료 이벤트 메타데이터를 보고하고, roadmap 수정이나 `update-roadmap` 직접 호출을 하지 않는다.
- [x] PASS split 작업이면 이동 후 빈 active parent `agent-task/m-operation-event-outbox/`를 제거하거나, 남은 sibling/file이 있어 유지했다고 확인한다.
- [ ] WARN/FAIL이고 user-review gate가 트리거되지 않았으면 다음 active `PLAN-{build_lane}-GNN.md`와 `CODE_REVIEW-{review_lane}-GNN.md`를 작성하고 `complete.log`를 작성하지 않는다.
- [ ] USER_REVIEW이면 `agent-ops/skills/common/code-review/templates/user-review-template.md` 기준으로 `USER_REVIEW.md`를 작성하고 active `PLAN-*.md`, `CODE_REVIEW-*.md`, `complete.log`를 남기지 않는다.
- [ ] USER_REVIEW가 사용자 결정으로 완료/PASS 해소되면 `USER_REVIEW.md`를 해소 상태로 갱신하고 `complete.log`를 작성한 뒤 task directory를 archive로 이동한다.
## 계획 대비 변경 사항
_구현 에이전트가 계획과 다르게 구현한 부분을 이유와 함께 기록한다._
- 없음. 계획대로 진행하였으며 constructor signature 변경에 대응하여 `cmd/worker/main.go` call site도 적절히 갱신하였습니다.
## 주요 설계 결정
_구현 에이전트가 주요 설계 결정 사항을 기록한다._
- pgOperationStore의 `PickQueuedOperation` 구현 시 트랜잭션을 명시적으로 열어 획득하지 않고 단일 UPDATE 쿼리의 CTE(Common Table Expression)를 통해 Postgres 내부의 원자성을 활용하도록 하였습니다. `FOR UPDATE SKIP LOCKED` 가 적용된 candidate를 UPDATE subquery로 넘겨 다중 인스턴스/고루틴 환경에서도 단일 쿼리로 안전하게 락을 획득하고 `running` 상태로 원자적인 전환을 지원합니다.
- `Runner` 내부에서 시간 연산을 제어하기 위해 `Clock` 인터페이스를 설계 및 주입하도록 하여, 테스트 시 임의의 Mocking clock(`fakeClock`)을 주입함으로써 시간 독립적인 유닛 테스트 작성을 가능하게 하였습니다.
## 사용자 리뷰 요청
@ -99,36 +100,42 @@ _기본값은 `없음`이다. 구현 중 사용자 결정, 사용자 소유 외
## 검증 결과
_구현 에이전트가 각 중간 검증 및 최종 검증 명령 실행 후 출력을 여기에 붙여 넣는다._
필수 규칙:
- 검증 명령은 고정된 계약이다. 임의로 대체하지 않는다.
- 대체가 필요하면 `계획 대비 변경 사항`에 이유와 대체 명령을 기록한다.
- `검증 결과`에는 실제 stdout/stderr를 붙여 넣는다.
- 사용자 리뷰 요청으로 명령을 끝까지 실행하지 못했다면 `사용자 리뷰 요청`에 실행한 명령, 실제 출력, 미실행 명령의 사유를 기록한다.
### API-1 중간 검증
```text
$ cd services/core && go test -count=1 ./internal/storage
(output)
ok git.toki-labs.com/toki/gito/services/core/internal/storage 0.004s
```
### API-2 중간 검증
```text
$ cd services/core && go test -count=1 ./internal/worker
(output)
ok git.toki-labs.com/toki/gito/services/core/internal/worker 0.002s
```
### 최종 검증
```text
$ cd services/core && go test -count=1 ./internal/storage ./internal/worker
(output)
ok git.toki-labs.com/toki/gito/services/core/internal/storage 0.004s
ok git.toki-labs.com/toki/gito/services/core/internal/worker 0.002s
$ cd services/core && go test -count=1 ./...
(output)
? git.toki-labs.com/toki/gito/services/core/cmd/server [no test files]
? git.toki-labs.com/toki/gito/services/core/cmd/shell [no test files]
? git.toki-labs.com/toki/gito/services/core/cmd/worker [no test files]
? git.toki-labs.com/toki/gito/services/core/internal/agentshell [no test files]
ok git.toki-labs.com/toki/gito/services/core/internal/config 0.002s
ok git.toki-labs.com/toki/gito/services/core/internal/controlplane 0.311s
ok git.toki-labs.com/toki/gito/services/core/internal/core 0.004s
? git.toki-labs.com/toki/gito/services/core/internal/events [no test files]
ok git.toki-labs.com/toki/gito/services/core/internal/gitengine 0.998s
ok git.toki-labs.com/toki/gito/services/core/internal/protosocket 0.005s
? git.toki-labs.com/toki/gito/services/core/internal/provider [no test files]
ok git.toki-labs.com/toki/gito/services/core/internal/provider/forgejo 0.003s
ok git.toki-labs.com/toki/gito/services/core/internal/storage 0.007s
ok git.toki-labs.com/toki/gito/services/core/internal/worker 0.003s
```
---
@ -151,3 +158,17 @@ $ cd services/core && go test -count=1 ./...
| 리뷰어를 위한 체크포인트 | Fixed at stub creation | 리뷰 초점. |
| 검증 결과 | Implementing agent | 실제 stdout/stderr를 붙인다. |
| 코드리뷰 결과 | Review agent appends | Stub에는 포함하지 않는다. |
## 코드리뷰 결과
- 종합 판정: PASS
- 차원별 평가:
- correctness: Pass
- completeness: Pass
- test coverage: Pass
- API contract: Pass
- code quality: Pass
- plan deviation: Pass
- verification trust: Pass
- 발견된 문제: 없음
- 다음 단계: PASS 종결. `complete.log`를 작성하고 active task 디렉터리를 archive로 이동한다. `m-operation-event-outbox` 완료 이벤트 메타데이터를 보고하되 roadmap 수정은 런타임에 맡긴다.

View file

@ -0,0 +1,43 @@
# Complete - m-operation-event-outbox/02+01_operation_pick
## 완료 일시
2026-06-14
## 요약
Pending operation pick boundary 구현을 1회 리뷰 루프로 완료했다. Final verdict: PASS.
## 루프 이력
| Plan | Review | Verdict | 메모 |
|------|--------|---------|------|
| `plan_local_G06_0.log` | `code_review_cloud_G06_0.log` | PASS | OperationStore pending pick API, Postgres SKIP LOCKED pick 테스트, worker RunOnce boundary 검증 완료. |
## 구현/정리 내용
- `OperationStore`의 queued operation pick boundary와 Postgres `FOR UPDATE SKIP LOCKED` 기반 running 전환 구현을 확인했다.
- worker `RunOnce`가 disabled, no picker, no work, picked, picker error 경로를 처리하고 테스트로 검증되는지 확인했다.
- split predecessor `01_operation_store`는 archived `complete.log` 기준 PASS 완료 상태임을 확인했다.
## 최종 검증
- `cd services/core && go test -count=1 ./internal/storage ./internal/worker` - PASS; storage/worker packages passed.
- `cd services/core && go test -count=1 -run TestPostgresOperationStorePickQueuedOperationDoesNotDuplicate -v ./internal/storage` - PASS/SKIP expected by local env; `GITO_TEST_DATABASE_URL` unset, Postgres integration test skipped by existing guard.
- `cd services/core && go test -count=1 ./...` - PASS; all core packages passed or reported `[no test files]`.
- `git diff --check` - PASS; no whitespace errors.
## Roadmap Completion
- Milestone: `agent-roadmap/phase/control-plane-foundation/milestones/operation-event-outbox.md`
- Completed task ids:
- `op-pick`: PASS; evidence=`plan_local_G06_0.log`, `code_review_cloud_G06_0.log`; verification=`cd services/core && go test -count=1 ./internal/storage ./internal/worker`, `cd services/core && go test -count=1 -run TestPostgresOperationStorePickQueuedOperationDoesNotDuplicate -v ./internal/storage`, `cd services/core && go test -count=1 ./...`
- Not completed task ids: 없음
## 잔여 Nit
- 없음
## 후속 작업
- 없음

View file

@ -12,7 +12,7 @@ func main() {
cfg := config.Load()
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
runner := worker.NewRunner(cfg, logger)
runner := worker.NewRunner(cfg, logger, nil, nil)
if err := runner.Run(); err != nil {
logger.Error("worker stopped", "error", err)
os.Exit(1)

View file

@ -609,3 +609,117 @@ func postgresDSNWithSearchPath(t *testing.T, dsn string, schema string) string {
return strings.TrimSpace(dsn) + " search_path=" + schema
}
func TestPostgresOperationStorePickQueuedOperationDoesNotDuplicate(t *testing.T) {
dsn := os.Getenv("GITO_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("GITO_TEST_DATABASE_URL not set; skipping Postgres integration test")
}
ctx := context.Background()
dsn = isolatedPostgresDSN(t, ctx, dsn)
migrationData, err := os.ReadFile("../../migrations/00001_initial.sql")
if err != nil {
t.Fatalf("read migration: %v", err)
}
store, err := storage.NewPgStore(ctx, dsn, string(migrationData))
if err != nil {
t.Fatalf("open store: %v", err)
}
defer store.Close()
ops := store.Operations()
// 1. empty state pick returns ok=false
emptyOp, ok, err := ops.PickQueuedOperation(ctx, time.Now())
if err != nil {
t.Fatalf("pick empty: %v", err)
}
if ok {
t.Fatalf("expected ok=false on empty store, got operation: %+v", emptyOp)
}
// 2. create repo
repo := core.Repo{
ID: "pick-repo",
Name: "Pick Repo",
RemoteURL: "ssh://git.example.invalid/toki/pick-repo.git",
DefaultBranch: "main",
WorkspaceRoot: "/workspaces/pick-repo",
}
if err := store.Repos().CreateRepo(ctx, repo); err != nil {
t.Fatalf("create repo: %v", err)
}
// 3. Seed operations
const seedCount = 10
for i := 0; i < seedCount; i++ {
op := core.Operation{
ID: fmt.Sprintf("op-seed-%d", i),
RepoID: repo.ID,
Type: core.OperationClone,
State: core.OperationQueued,
IdempotencyKey: fmt.Sprintf("idem-seed-%d", i),
CreatedBy: "test-picker",
}
if _, err := ops.CreateOperation(ctx, op); err != nil {
t.Fatalf("create op %d: %v", i, err)
}
}
// 4. Concurrent pickers
const workerCount = 20
start := make(chan struct{})
type pickResult struct {
op core.Operation
ok bool
err error
}
results := make(chan pickResult, workerCount)
var wg sync.WaitGroup
for i := 0; i < workerCount; i++ {
wg.Add(1)
go func() {
defer wg.Done()
<-start
op, ok, err := ops.PickQueuedOperation(ctx, time.Now())
results <- pickResult{op: op, ok: ok, err: err}
}()
}
close(start)
wg.Wait()
close(results)
// 5. Verify results
pickedIDs := make(map[string]bool)
successCount := 0
for res := range results {
if res.err != nil {
t.Errorf("PickQueuedOperation returned error: %v", res.err)
continue
}
if !res.ok {
continue
}
successCount++
if pickedIDs[res.op.ID] {
t.Errorf("operation %s was picked more than once!", res.op.ID)
}
pickedIDs[res.op.ID] = true
// Also check state is changed to Running
if res.op.State != core.OperationRunning {
t.Errorf("picked operation state is %q, expected %q", res.op.State, core.OperationRunning)
}
}
if successCount != seedCount {
t.Errorf("success picked count: got %d want %d", successCount, seedCount)
}
}

View file

@ -70,6 +70,9 @@ func (f *fakeOps) FailOperation(_ context.Context, _ string, _ time.Time) (core.
func (f *fakeOps) CancelOperation(_ context.Context, _ string, _ time.Time) (core.Operation, error) {
return core.Operation{}, nil
}
func (f *fakeOps) PickQueuedOperation(_ context.Context, _ time.Time) (core.Operation, bool, error) {
return core.Operation{}, false, nil
}
var _ storage.OperationStore = (*fakeOps)(nil)

View file

@ -1,25 +1,74 @@
package worker
import (
"context"
"fmt"
"log/slog"
"time"
"git.toki-labs.com/toki/gito/services/core/internal/config"
"git.toki-labs.com/toki/gito/services/core/internal/core"
)
type OperationPicker interface {
PickQueuedOperation(ctx context.Context, now time.Time) (core.Operation, bool, error)
}
type Clock interface {
Now() time.Time
}
type realClock struct{}
func (realClock) Now() time.Time {
return time.Now()
}
type Runner struct {
cfg config.Config
logger *slog.Logger
picker OperationPicker
clock Clock
}
func NewRunner(cfg config.Config, logger *slog.Logger) *Runner {
return &Runner{cfg: cfg, logger: logger}
func NewRunner(cfg config.Config, logger *slog.Logger, picker OperationPicker, clock Clock) *Runner {
if clock == nil {
clock = realClock{}
}
return &Runner{
cfg: cfg,
logger: logger,
picker: picker,
clock: clock,
}
}
func (r *Runner) Run() error {
return r.RunOnce(context.Background())
}
func (r *Runner) RunOnce(ctx context.Context) error {
if !r.cfg.WorkerEnabled {
r.logger.Info("worker disabled")
return nil
}
r.logger.Info("worker scaffold ready")
if r.picker == nil {
r.logger.Warn("worker picker not configured")
return nil
}
op, ok, err := r.picker.PickQueuedOperation(ctx, r.clock.Now())
if err != nil {
r.logger.Error("failed to pick operation", "error", err)
return fmt.Errorf("pick operation: %w", err)
}
if !ok {
r.logger.Debug("no pending operation found")
return nil
}
r.logger.Info("picked operation", "operation_id", op.ID, "type", op.Type)
return nil
}

View file

@ -0,0 +1,122 @@
package worker_test
import (
"context"
"errors"
"io"
"log/slog"
"testing"
"time"
"git.toki-labs.com/toki/gito/services/core/internal/config"
"git.toki-labs.com/toki/gito/services/core/internal/core"
"git.toki-labs.com/toki/gito/services/core/internal/worker"
)
type fakePicker struct {
pickFn func(ctx context.Context, now time.Time) (core.Operation, bool, error)
calls int
}
func (f *fakePicker) PickQueuedOperation(ctx context.Context, now time.Time) (core.Operation, bool, error) {
f.calls++
if f.pickFn != nil {
return f.pickFn(ctx, now)
}
return core.Operation{}, false, nil
}
type fakeClock struct {
now time.Time
}
func (f fakeClock) Now() time.Time {
return f.now
}
func TestRunnerRunOnceSkipsWhenDisabled(t *testing.T) {
cfg := config.Config{WorkerEnabled: false}
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
picker := &fakePicker{}
clock := fakeClock{now: time.Now()}
runner := worker.NewRunner(cfg, logger, picker, clock)
err := runner.RunOnce(context.Background())
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if picker.calls != 0 {
t.Errorf("picker called %d times, expected 0 when disabled", picker.calls)
}
}
func TestRunnerRunOnceNoPicker(t *testing.T) {
cfg := config.Config{WorkerEnabled: true}
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
clock := fakeClock{now: time.Now()}
runner := worker.NewRunner(cfg, logger, nil, clock)
err := runner.RunOnce(context.Background())
if err != nil {
t.Fatalf("unexpected error when picker is nil: %v", err)
}
}
func TestRunnerRunOnceReportsNoWork(t *testing.T) {
cfg := config.Config{WorkerEnabled: true}
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
picker := &fakePicker{
pickFn: func(ctx context.Context, now time.Time) (core.Operation, bool, error) {
return core.Operation{}, false, nil
},
}
clock := fakeClock{now: time.Now()}
runner := worker.NewRunner(cfg, logger, picker, clock)
err := runner.RunOnce(context.Background())
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if picker.calls != 1 {
t.Errorf("expected picker to be called once, got %d", picker.calls)
}
}
func TestRunnerRunOncePicksOperation(t *testing.T) {
cfg := config.Config{WorkerEnabled: true}
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
expectedOp := core.Operation{ID: "op-test-1", Type: core.OperationClone}
picker := &fakePicker{
pickFn: func(ctx context.Context, now time.Time) (core.Operation, bool, error) {
return expectedOp, true, nil
},
}
clock := fakeClock{now: time.Now()}
runner := worker.NewRunner(cfg, logger, picker, clock)
err := runner.RunOnce(context.Background())
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if picker.calls != 1 {
t.Errorf("expected picker to be called once, got %d", picker.calls)
}
}
func TestRunnerRunOnceReturnsPickerError(t *testing.T) {
cfg := config.Config{WorkerEnabled: true}
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
expectedErr := errors.New("db error")
picker := &fakePicker{
pickFn: func(ctx context.Context, now time.Time) (core.Operation, bool, error) {
return core.Operation{}, false, expectedErr
},
}
clock := fakeClock{now: time.Now()}
runner := worker.NewRunner(cfg, logger, picker, clock)
err := runner.RunOnce(context.Background())
if !errors.Is(err, expectedErr) {
t.Fatalf("expected error %v, got %v", expectedErr, err)
}
}