feat: workitempipeline 패키지와 HTTP 핸들러 개선

This commit is contained in:
toki 2026-05-23 21:19:26 +09:00
parent 0e1fb9cf27
commit 3e2cad7e03
5 changed files with 277 additions and 114 deletions

View file

@ -30,7 +30,7 @@ Work Item Provider Pipeline Design
- [x] 초기 entrypoint는 기존 `POST /api/integrations/plane/tasks`를 기준선으로 삼는다. - [x] 초기 entrypoint는 기존 `POST /api/integrations/plane/tasks`를 기준선으로 삼는다.
- [x] 생성 단계는 Plane work item 조회와 core task 저장까지만 담당한다. - [x] 생성 단계는 Plane work item 조회와 core task 저장까지만 담당한다.
- [x] 생성된 task는 `pending` 상태로 남기고 enqueue는 별도 단계에서 처리한다. - [x] 생성된 task는 `pending` 상태로 남기고 enqueue는 별도 단계에서 처리한다.
- [x] Plane 연결 정보는 provider-neutral external ref와 Plane metadata에 함께 저장한다. - [x] Plane 연결 정보는 provider-neutral external ref와 work item metadata에 함께 저장한다.
- [x] 자동 enqueue 여부와 사용자/운영 트리거 경계를 결정한다. - [x] 자동 enqueue 여부와 사용자/운영 트리거 경계를 결정한다.
- [x] `backlog`에서 `todo`로 옮기는 주체는 사용자로 둔다. - [x] `backlog`에서 `todo`로 옮기는 주체는 사용자로 둔다.
- [x] `todo` 상태여도 agent 작업자가 지정된 work item만 자동 실행 후보로 본다. - [x] `todo` 상태여도 agent 작업자가 지정된 work item만 자동 실행 후보로 본다.
@ -47,9 +47,9 @@ Work Item Provider Pipeline Design
- [x] provider별 상태/라벨/comment 매핑은 adapter 설정으로 분리한다. - [x] provider별 상태/라벨/comment 매핑은 adapter 설정으로 분리한다.
- [ ] 현재 Plane 고정 진입부를 provider-neutral 구조로 리팩토링한다. - [ ] 현재 Plane 고정 진입부를 provider-neutral 구조로 리팩토링한다.
- [ ] `POST /api/integrations/plane/tasks`의 Plane 전용 흐름을 generic pipeline service 아래로 옮긴다. - [ ] `POST /api/integrations/plane/tasks`의 Plane 전용 흐름을 generic pipeline service 아래로 옮긴다.
- [ ] HTTP handler가 Plane 타입과 직접 결합하지 않도록 요청 DTO와 변환 책임을 분리한다. - [x] HTTP handler가 Plane 타입과 직접 결합하지 않도록 요청 DTO와 변환 책임을 분리한다.
- [ ] `buildPlaneCreateTaskInput`의 core task 생성 로직을 provider-neutral mapper로 분리한다. - [x] `buildPlaneCreateTaskInput`의 core task 생성 로직을 provider-neutral mapper로 분리한다.
- [ ] 기존 Plane endpoint는 compatibility entrypoint로 유지하거나 generic endpoint로 대체할지 결정한다. - [x] 기존 Plane endpoint는 compatibility entrypoint로 유지하거나 generic endpoint로 대체할지 결정한다.
- [ ] provider-neutral 상태와 projection 계약을 확정한다. - [ ] provider-neutral 상태와 projection 계약을 확정한다.
- [x] board state는 `backlog`, `todo`, `in_progress`, `testing`, `complete`, `cancel`로 둔다. - [x] board state는 `backlog`, `todo`, `in_progress`, `testing`, `complete`, `cancel`로 둔다.
- [x] `in_progress` 내부 agent 상태는 core task metadata를 canonical source로 둔다. - [x] `in_progress` 내부 agent 상태는 core task metadata를 canonical source로 둔다.
@ -65,8 +65,8 @@ Work Item Provider Pipeline Design
## 완료 기준 ## 완료 기준
- [ ] Plane/Jira work item에서 core task로 이어지는 provider-neutral pipeline entrypoint와 계약이 문서화되어 있다. - [ ] Plane/Jira work item에서 core task로 이어지는 provider-neutral pipeline entrypoint와 계약이 문서화되어 있다.
- [ ] core workflow/pipeline 계층이 Plane package 타입에 직접 의존하지 않는다. - [x] core workflow/pipeline 계층이 Plane package 타입에 직접 의존하지 않는다.
- [ ] Plane 직접 접속부는 adapter 구현에 격리되고, Jira adapter 추가 지점이 명확하다. - [x] Plane 직접 접속부는 adapter 구현에 격리되고, Jira adapter 추가 지점이 명확하다.
- [ ] enqueue 조건, idempotency 기준, 실패 표시 방식이 결정되어 있다. - [ ] enqueue 조건, idempotency 기준, 실패 표시 방식이 결정되어 있다.
- [ ] provider board state, core canonical state, label/comment projection 계약이 문서화되어 있다. - [ ] provider board state, core canonical state, label/comment projection 계약이 문서화되어 있다.
- [ ] completed/failed/cancelled 결과를 provider에 반영하는 최소 정책이 결정되어 있다. - [ ] completed/failed/cancelled 결과를 provider에 반영하는 최소 정책이 결정되어 있다.
@ -91,21 +91,21 @@ Work Item Provider Pipeline Design
- 선행 작업: Plane Communication Foundation - 선행 작업: Plane Communication Foundation
- 후속 작업: Workflow Core - 후속 작업: Workflow Core
- 현재 코드 결합 상태: - 현재 코드 결합 상태:
- HTTP router는 `POST /api/integrations/plane/tasks`로 Plane 전용 entrypoint를 가진다. - HTTP router는 기존 compatibility route인 `POST /api/integrations/plane/tasks`를 유지한다.
- handler는 `PlaneWorkItemClient`, `plane.WorkItemRef`, `plane.WorkItem`에 직접 의존한다. - handler는 Plane adapter package나 Plane DTO를 직접 import하지 않고, 로컬 `WorkItemReader` interface와 `workitem.Ref`, `workitem.WorkItem`으로 조회 결과를 다룬다.
- `buildPlaneCreateTaskInput`이 Plane 조회 결과를 core task payload/external ref로 직접 변환한다. - `buildPlaneCreateTaskInput`은 HTTP package 안에서 core task payload/external ref를 직접 조립하지 않고 `workitem.BuildCreateTaskInput`에 위임한다.
- `services/core/internal/workitem`은 provider-neutral DTO, optional facet interface, projection mapping 제공한다. - `services/core/internal/workitem`은 provider-neutral DTO, optional facet interface, projection mapping, task create mapper를 제공한다.
- Plane adapter는 `workitem.Provider`, `Reader`, `Commenter`, `StatusProjector`를 구현하며 `LabelProjector`는 capability false로 남긴다. - Plane adapter는 `workitem.Provider`, `Reader`, `Commenter`, `StatusProjector`를 구현하며 `LabelProjector`는 capability false로 남긴다.
- DB schema는 `external_provider`, `external_id`, `external_url`, `external_metadata`로 provider-neutral 토대가 있으므로 유지한다. - DB schema는 `external_provider`, `external_id`, `external_url`, `external_metadata`로 provider-neutral 토대가 있으므로 유지한다.
- 리팩토링 목표는 Plane/Jira 직접 접속부를 adapter로 격리하고, 그 아래 pipeline/service 계층은 provider-neutral DTO와 interface만 보게 하는 것이다. - 아직 별도 generic pipeline service는 추출되지 않았으며, compatibility HTTP handler가 provider-neutral reader/mapper와 `workflow.CreateTask` 호출을 직접 조립한다.
- 초기 생성/연결 흐름: - 초기 생성/연결 흐름:
- 운영자 또는 상위 자동화가 `POST /api/integrations/plane/tasks`를 호출한다. - 운영자 또는 상위 자동화가 `POST /api/integrations/plane/tasks`를 호출한다.
- 요청 필수값은 `workspace_slug`, `project_id`, `work_item_id`이고, `state_id`, `external_url`, `comment`는 선택값으로 받는다. - 요청 필수값은 `workspace_slug`, `project_id`, `work_item_id`이고, `state_id`, `external_url`, `comment`는 선택값으로 받는다.
- core는 Plane adapter로 work item detail을 조회해 title과 description 후보를 얻는다. - core는 Plane adapter`workitem.Reader` 구현으로 work item detail을 조회해 provider-neutral title과 description 후보를 얻는다.
- task title은 Plane work item name을 우선하고, 비어 있으면 `work_item_id`를 사용한다. - task title은 provider work item title을 우선하고, 비어 있으면 `work_item_id`를 사용한다.
- task payload의 `message`는 요청 `comment`, Plane stripped description, plain description, HTML description, title 순서로 선택한다. - task payload의 `message`는 요청 `comment`, provider text description, HTML description, title 순서로 선택하며 공백-only 후보는 건너뛴다.
- task source는 `plane`, external provider는 `plane`, external id는 `work_item_id`로 저장한다. - task source는 `plane`, external provider는 `plane`, external id는 `work_item_id`로 저장한다.
- `workspace_slug`, `project_id`, `work_item_id`, `state_id`, `external_url``payload.plane`과 `external_metadata`에 저장한다. - provider-neutral metadata key `provider`, `tenant`, `project`, `id`, `state_id`, `external_url``payload.work_item`과 `external_metadata`에 저장한다.
- 생성 직후 task 상태는 `pending`이며, 이 entrypoint는 enqueue, Plane 상태 변경, 결과 comment 발행을 수행하지 않는다. - 생성 직후 task 상태는 `pending`이며, 이 entrypoint는 enqueue, Plane 상태 변경, 결과 comment 발행을 수행하지 않는다.
- 상태 설계 결정: - 상태 설계 결정:
- provider board state는 `backlog`, `todo`, `in_progress`, `testing`, `complete`, `cancel`의 소유권/검증 단계로 유지한다. - provider board state는 `backlog`, `todo`, `in_progress`, `testing`, `complete`, `cancel`의 소유권/검증 단계로 유지한다.
@ -144,10 +144,9 @@ Work Item Provider Pipeline Design
- `✅ VERIFY | <요약>`: 테스트/검증 완료 기록 - `✅ VERIFY | <요약>`: 테스트/검증 완료 기록
- 각 comment 본문은 1~2줄 요약을 기본으로 하며, 자세한 실행 로그나 상태 metadata는 core에 남긴다. - 각 comment 본문은 1~2줄 요약을 기본으로 하며, 자세한 실행 로그나 상태 metadata는 core에 남긴다.
- 현재 지점 / 착수 상태: - 현재 지점 / 착수 상태:
- 현재 상태는 `진행 중`이며, provider adapter interface 설계와 Plane adapter facet 검증은 완료됐다. - 현재 상태는 `진행 중`이며, provider adapter interface 설계, Plane adapter facet 검증, provider-neutral task mapper, HTTP compatibility endpoint의 Plane DTO 의존 제거가 완료됐다.
- 다음 진행 후보는 `현재 Plane 고정 진입부를 provider-neutral 구조로 리팩토링한다` 항목이다. - `agent-task/archive/2026/05/provider_neutral_plane_entrypoint/01_pipeline_mapper/complete.log`에서 `workitem.TaskCreateInput``workitem.BuildCreateTaskInput` 추가 및 whitespace fallback 회귀 수정이 PASS로 정리됐다.
- 정리된 내용은 provider-neutral 방향성, Plane/Jira adapter 추상화 필요성, label/comment 기반 상태 투영 샘플이다. - `agent-task/archive/2026/05/provider_neutral_plane_entrypoint/02+01_http_plane_entrypoint/complete.log`에서 HTTP handler의 provider-neutral reader/mapper 전환과 compatibility route 유지가 PASS로 정리됐다.
- 직전 항목인 `자동 enqueue 여부와 사용자/운영 트리거 경계를 결정한다`는 완료됐다. - 다음 진행 후보는 남은 `generic pipeline service` 경계 추출 여부를 결정하거나, core task metadata schema와 provider label mapping 확정으로 넘어가는 것이다.
- 이후 남은 큰 결정은 Plane 고정 진입부 전환, core metadata schema, provider label mapping이다. - 이후 남은 큰 결정은 core metadata schema, provider label mapping, comment prefix 작성 타이밍, idempotency/retry 기준이다.
- 다음에 이 마일스톤을 착수하면 Plane 고정 진입부 리팩토링에서 HTTP handler의 Plane 타입 의존을 끊는다.
- 확인 필요: trigger 방식은 현재 구현 상태와 운영 기대치를 보고 수동 endpoint 유지, webhook, polling 중 하나를 선택한다. - 확인 필요: trigger 방식은 현재 구현 상태와 운영 기대치를 보고 수동 endpoint 유지, webhook, polling 중 하나를 선택한다.

View file

@ -14,28 +14,33 @@ import (
"github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool" "github.com/jackc/pgx/v5/pgxpool"
"github.com/nomadcode/nomadcode-core/internal/storage"
"github.com/nomadcode/nomadcode-core/internal/workflow" "github.com/nomadcode/nomadcode-core/internal/workflow"
"github.com/nomadcode/nomadcode-core/internal/workitem" "github.com/nomadcode/nomadcode-core/internal/workitem"
"github.com/nomadcode/nomadcode-core/internal/workitempipeline"
) )
type WorkItemReader interface { type WorkItemTaskCreator interface {
FetchWorkItem(ctx context.Context, ref workitem.Ref) (workitem.WorkItem, error) CreateTaskFromWorkItem(ctx context.Context, input workitempipeline.CreateTaskInput) (storage.Task, error)
} }
type Handler struct { type Handler struct {
db *pgxpool.Pool db *pgxpool.Pool
workflow *workflow.Service workflow *workflow.Service
workItems WorkItemReader workItemTasks WorkItemTaskCreator
logger *slog.Logger logger *slog.Logger
} }
func NewHandler(pool *pgxpool.Pool, workflowService *workflow.Service, workItemReader WorkItemReader, logger *slog.Logger) *Handler { func NewHandler(pool *pgxpool.Pool, workflowService *workflow.Service, workItemReader workitem.Reader, logger *slog.Logger) *Handler {
return &Handler{ h := &Handler{
db: pool, db: pool,
workflow: workflowService, workflow: workflowService,
workItems: workItemReader, logger: logger,
logger: logger,
} }
if workItemReader != nil && workflowService != nil {
h.workItemTasks = workitempipeline.New(workItemReader, workflowService)
}
return h
} }
func (h *Handler) Healthz(w stdhttp.ResponseWriter, r *stdhttp.Request) { func (h *Handler) Healthz(w stdhttp.ResponseWriter, r *stdhttp.Request) {
@ -85,7 +90,7 @@ type createPlaneTaskRequest struct {
} }
func (h *Handler) CreatePlaneTask(w stdhttp.ResponseWriter, r *stdhttp.Request) { func (h *Handler) CreatePlaneTask(w stdhttp.ResponseWriter, r *stdhttp.Request) {
if h.workItems == nil { if h.workItemTasks == nil {
writeError(w, stdhttp.StatusServiceUnavailable, "plane client is not configured") writeError(w, stdhttp.StatusServiceUnavailable, "plane client is not configured")
return return
} }
@ -102,27 +107,17 @@ func (h *Handler) CreatePlaneTask(w stdhttp.ResponseWriter, r *stdhttp.Request)
return return
} }
item, err := h.workItems.FetchWorkItem(r.Context(), ref) task, err := h.workItemTasks.CreateTaskFromWorkItem(r.Context(), workitempipeline.CreateTaskInput{
Ref: ref,
StateID: input.StateID,
Comment: input.Comment,
})
if err != nil { if err != nil {
h.writeServiceError(w, err) h.writeServiceError(w, err)
return return
} }
createInput, err := buildPlaneCreateTaskInput(input, item) writeJSON(w, stdhttp.StatusCreated, map[string]string{
if err != nil {
writeError(w, stdhttp.StatusBadRequest, err.Error())
return
}
task, err := h.workflow.CreateTask(r.Context(), createInput)
if err != nil {
h.writeServiceError(w, err)
return
}
status := stdhttp.StatusCreated
writeJSON(w, status, map[string]string{
"id": task.ID, "id": task.ID,
"status": task.Status, "status": task.Status,
"external_provider": stringValue(task.ExternalProvider), "external_provider": stringValue(task.ExternalProvider),
@ -189,18 +184,6 @@ func (input createPlaneTaskRequest) workItemRef() (workitem.Ref, error) {
}) })
} }
func buildPlaneCreateTaskInput(input createPlaneTaskRequest, item workitem.WorkItem) (workflow.CreateTaskInput, error) {
ref, err := input.workItemRef()
if err != nil {
return workflow.CreateTaskInput{}, err
}
return workitem.BuildCreateTaskInput(workitem.TaskCreateInput{
Ref: ref,
StateID: input.StateID,
Comment: input.Comment,
}, item)
}
func stringValue(value *string) string { func stringValue(value *string) string {
if value == nil { if value == nil {
return "" return ""

View file

@ -1,12 +1,29 @@
package http package http
import ( import (
"context"
"encoding/json" "encoding/json"
"net/http"
"net/http/httptest"
"strings"
"testing" "testing"
"github.com/nomadcode/nomadcode-core/internal/storage"
"github.com/nomadcode/nomadcode-core/internal/workitem" "github.com/nomadcode/nomadcode-core/internal/workitem"
"github.com/nomadcode/nomadcode-core/internal/workitempipeline"
) )
type fakeWorkItemTaskCreator struct {
task storage.Task
err error
calledWith *workitempipeline.CreateTaskInput
}
func (f *fakeWorkItemTaskCreator) CreateTaskFromWorkItem(_ context.Context, input workitempipeline.CreateTaskInput) (storage.Task, error) {
f.calledWith = &input
return f.task, f.err
}
func TestPlaneTaskRequestBuildsProviderNeutralRef(t *testing.T) { func TestPlaneTaskRequestBuildsProviderNeutralRef(t *testing.T) {
req := createPlaneTaskRequest{ req := createPlaneTaskRequest{
WorkspaceSlug: " acme ", WorkspaceSlug: " acme ",
@ -35,67 +52,68 @@ func TestPlaneTaskRequestBuildsProviderNeutralRef(t *testing.T) {
} }
} }
func TestBuildPlaneCreateTaskInputUsesCommentBeforeDescription(t *testing.T) { func TestCreatePlaneTaskDelegatesToWorkItemPipeline(t *testing.T) {
input, err := buildPlaneCreateTaskInput(createPlaneTaskRequest{ creator := &fakeWorkItemTaskCreator{
WorkspaceSlug: "general", task: storage.Task{ID: "task-123", Status: "pending"},
ProjectID: "project-1",
WorkItemID: "work-1",
Comment: "operator instruction",
}, workitem.WorkItem{
Ref: workitem.Ref{Provider: "plane", Tenant: "general", Project: "project-1", ID: "work-1"},
Title: "Fix issue",
DescriptionText: "plane description",
})
if err != nil {
t.Fatalf("buildPlaneCreateTaskInput returned error: %v", err)
} }
if input.Title != "Fix issue" { h := &Handler{workItemTasks: creator}
t.Fatalf("unexpected title: %q", input.Title)
body := `{"workspace_slug":"acme","project_id":"proj-1","work_item_id":"work-1","state_id":"state-1","comment":"my note"}`
req := httptest.NewRequest(http.MethodPost, "/api/integrations/plane/tasks", strings.NewReader(body))
req.Header.Set("Content-Type", "application/json")
rec := httptest.NewRecorder()
h.CreatePlaneTask(rec, req)
if rec.Code != http.StatusCreated {
t.Fatalf("expected %d, got %d: %s", http.StatusCreated, rec.Code, rec.Body.String())
} }
if input.Source != "plane" { if creator.calledWith == nil {
t.Fatalf("unexpected source: %q", input.Source) t.Fatal("CreateTaskFromWorkItem was not called")
}
if creator.calledWith.Ref.Provider != "plane" {
t.Errorf("ref.provider: got %q", creator.calledWith.Ref.Provider)
}
if creator.calledWith.Ref.Tenant != "acme" {
t.Errorf("ref.tenant: got %q", creator.calledWith.Ref.Tenant)
}
if creator.calledWith.Ref.Project != "proj-1" {
t.Errorf("ref.project: got %q", creator.calledWith.Ref.Project)
}
if creator.calledWith.Ref.ID != "work-1" {
t.Errorf("ref.id: got %q", creator.calledWith.Ref.ID)
}
if creator.calledWith.StateID != "state-1" {
t.Errorf("state_id: got %q", creator.calledWith.StateID)
}
if creator.calledWith.Comment != "my note" {
t.Errorf("comment: got %q", creator.calledWith.Comment)
} }
var payload map[string]any var resp map[string]string
if err := json.Unmarshal(input.Payload, &payload); err != nil { if err := json.NewDecoder(rec.Body).Decode(&resp); err != nil {
t.Fatalf("decode payload: %v", err) t.Fatalf("decode response: %v", err)
} }
if payload["message"] != "operator instruction" { if resp["id"] != "task-123" || resp["status"] != "pending" {
t.Fatalf("unexpected message: %#v", payload) t.Errorf("response: %v", resp)
} }
} }
func TestBuildPlaneCreateTaskInputStoresExternalMetadata(t *testing.T) { func TestCreatePlaneTaskReturns503WhenPipelineNotConfigured(t *testing.T) {
input, err := buildPlaneCreateTaskInput(createPlaneTaskRequest{ h := &Handler{}
WorkspaceSlug: "general", req := httptest.NewRequest(http.MethodPost, "/api/integrations/plane/tasks", strings.NewReader(`{}`))
ProjectID: "project-1", rec := httptest.NewRecorder()
WorkItemID: "work-1",
StateID: "state-1",
ExternalURL: "https://plane.example/work-1",
}, workitem.WorkItem{
Ref: workitem.Ref{Provider: "plane", Tenant: "general", Project: "project-1", ID: "work-1"},
Title: "Fix issue",
})
if err != nil {
t.Fatalf("buildPlaneCreateTaskInput returned error: %v", err)
}
if input.External == nil {
t.Fatal("expected external ref")
}
if input.External.Provider != "plane" || input.External.ID != "work-1" {
t.Fatalf("unexpected external ref: %#v", input.External)
}
var metadata map[string]string h.CreatePlaneTask(rec, req)
if err := json.Unmarshal(input.External.Metadata, &metadata); err != nil {
t.Fatalf("decode metadata: %v", err) if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("expected 503, got %d", rec.Code)
} }
if metadata["provider"] != "plane" || var resp map[string]string
metadata["tenant"] != "general" || if err := json.NewDecoder(rec.Body).Decode(&resp); err != nil {
metadata["project"] != "project-1" || t.Fatalf("decode response: %v", err)
metadata["id"] != "work-1" || }
metadata["state_id"] != "state-1" || if resp["error"] != "plane client is not configured" {
metadata["external_url"] != "https://plane.example/work-1" { t.Errorf("error message: got %q", resp["error"])
t.Fatalf("unexpected metadata: %#v", metadata)
} }
} }

View file

@ -0,0 +1,44 @@
package workitempipeline
import (
"context"
"github.com/nomadcode/nomadcode-core/internal/storage"
"github.com/nomadcode/nomadcode-core/internal/workflow"
"github.com/nomadcode/nomadcode-core/internal/workitem"
)
type CreateTaskInput struct {
Ref workitem.Ref
StateID string
Comment string
}
type TaskCreator interface {
CreateTask(ctx context.Context, input workflow.CreateTaskInput) (storage.Task, error)
}
type Service struct {
reader workitem.Reader
tasks TaskCreator
}
func New(reader workitem.Reader, tasks TaskCreator) *Service {
return &Service{reader: reader, tasks: tasks}
}
func (s *Service) CreateTaskFromWorkItem(ctx context.Context, input CreateTaskInput) (storage.Task, error) {
item, err := s.reader.FetchWorkItem(ctx, input.Ref)
if err != nil {
return storage.Task{}, err
}
createInput, err := workitem.BuildCreateTaskInput(workitem.TaskCreateInput{
Ref: input.Ref,
StateID: input.StateID,
Comment: input.Comment,
}, item)
if err != nil {
return storage.Task{}, err
}
return s.tasks.CreateTask(ctx, createInput)
}

View file

@ -0,0 +1,119 @@
package workitempipeline
import (
"context"
"encoding/json"
"errors"
"testing"
"github.com/nomadcode/nomadcode-core/internal/storage"
"github.com/nomadcode/nomadcode-core/internal/workflow"
"github.com/nomadcode/nomadcode-core/internal/workitem"
)
type fakeReader struct {
item workitem.WorkItem
err error
}
func (f *fakeReader) FetchWorkItem(_ context.Context, _ workitem.Ref) (workitem.WorkItem, error) {
return f.item, f.err
}
type fakeTaskCreator struct {
task storage.Task
err error
calledWith *workflow.CreateTaskInput
}
func (f *fakeTaskCreator) CreateTask(_ context.Context, input workflow.CreateTaskInput) (storage.Task, error) {
f.calledWith = &input
return f.task, f.err
}
func TestCreateTaskFromWorkItemFetchesMapsAndCreatesTask(t *testing.T) {
ref := workitem.Ref{Provider: "plane", Tenant: "acme", Project: "proj-1", ID: "work-1", URL: "https://example.com/work-1"}
reader := &fakeReader{item: workitem.WorkItem{
Ref: ref,
Title: "Fix issue",
}}
creator := &fakeTaskCreator{task: storage.Task{ID: "task-abc", Status: "pending"}}
svc := New(reader, creator)
result, err := svc.CreateTaskFromWorkItem(context.Background(), CreateTaskInput{
Ref: ref,
StateID: "state-1",
Comment: "my comment",
})
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if result.ID != "task-abc" {
t.Errorf("id: got %q", result.ID)
}
if creator.calledWith == nil {
t.Fatal("CreateTask was not called")
}
if creator.calledWith.Source != "plane" {
t.Errorf("source: got %q", creator.calledWith.Source)
}
var payload map[string]any
if err := json.Unmarshal(creator.calledWith.Payload, &payload); err != nil {
t.Fatalf("payload unmarshal: %v", err)
}
if payload["message"] != "my comment" {
t.Errorf("message: got %v", payload["message"])
}
wi, ok := payload["work_item"].(map[string]any)
if !ok {
t.Fatalf("work_item missing: %v", payload)
}
checks := map[string]string{
"provider": "plane",
"tenant": "acme",
"project": "proj-1",
"id": "work-1",
}
for k, want := range checks {
if wi[k] != want {
t.Errorf("work_item[%q]: got %v, want %q", k, wi[k], want)
}
}
if creator.calledWith.External == nil ||
creator.calledWith.External.Provider != "plane" ||
creator.calledWith.External.ID != "work-1" {
t.Errorf("external ref: %v", creator.calledWith.External)
}
}
func TestCreateTaskFromWorkItemPropagatesReaderError(t *testing.T) {
readerErr := errors.New("not found")
reader := &fakeReader{err: readerErr}
creator := &fakeTaskCreator{}
svc := New(reader, creator)
_, err := svc.CreateTaskFromWorkItem(context.Background(), CreateTaskInput{
Ref: workitem.Ref{Provider: "plane", ID: "work-1"},
})
if !errors.Is(err, readerErr) {
t.Errorf("expected reader error, got %v", err)
}
if creator.calledWith != nil {
t.Error("CreateTask should not have been called")
}
}
func TestCreateTaskFromWorkItemPropagatesCreateTaskError(t *testing.T) {
ref := workitem.Ref{Provider: "plane", Tenant: "acme", Project: "proj-1", ID: "work-1"}
reader := &fakeReader{item: workitem.WorkItem{Ref: ref, Title: "T"}}
createErr := errors.New("db error")
creator := &fakeTaskCreator{err: createErr}
svc := New(reader, creator)
_, err := svc.CreateTaskFromWorkItem(context.Background(), CreateTaskInput{Ref: ref})
if !errors.Is(err, createErr) {
t.Errorf("expected create error, got %v", err)
}
}