From e7ddbcdfed9b448f20cb14cd29686fcdc3d3a1c0 Mon Sep 17 00:00:00 2001 From: toki Date: Sun, 24 May 2026 18:10:40 +0900 Subject: [PATCH] feat: add task metadata migration and update workflow service - Add migration for task metadata column (00003) - Update workflow model and service with metadata support - Update database layer (models, queries) for metadata field - Update storage store to handle metadata - Update roadmap documents and milestones --- agent-ops/roadmap/ROADMAP.md | 5 +- agent-ops/roadmap/current.md | 2 +- .../milestones/external-integration.md | 16 ++--- .../milestones/plane-task-pipeline-design.md | 39 +++++++------ .../project-workspace-management-ux.md | 16 ++--- agent-ops/roadmap/milestones/workflow-core.md | 14 ++--- services/core/internal/db/models.go | 1 + services/core/internal/db/tasks.sql.go | 58 ++++++++++++++++--- services/core/internal/storage/store.go | 17 ++++++ services/core/internal/workflow/model.go | 11 ++++ services/core/internal/workflow/service.go | 24 ++++++++ .../core/internal/workflow/service_test.go | 25 ++++++++ .../migrations/00003_add_task_metadata.sql | 7 +++ services/core/queries/tasks.sql | 22 ++++--- 14 files changed, 189 insertions(+), 68 deletions(-) create mode 100644 services/core/migrations/00003_add_task_metadata.sql diff --git a/agent-ops/roadmap/ROADMAP.md b/agent-ops/roadmap/ROADMAP.md index 20f7649..63ceb60 100644 --- a/agent-ops/roadmap/ROADMAP.md +++ b/agent-ops/roadmap/ROADMAP.md @@ -52,5 +52,6 @@ NomadCode는 모바일 앱, 웹 콘솔, core 서비스, 공유 계약, agent-ope - 완료 또는 폐기되어 아카이브된 Milestone은 이 문서의 `아카이브 Milestone 요약`에 당시 요약만 남기고, 아카이브 문서 링크나 상세 경로는 남기지 않는다. - 상세 문서가 있는 `agent-ops/roadmap/archive/**`는 사용자가 명시적으로 요청한 경우에만 읽는다. - 아카이브된 Milestone 문서는 최신 템플릿이나 스킬 규약에 맞춰 재포맷하지 않는다. -- 선택된 Milestone의 `구현 잠금` 섹션이 없거나 상태가 `잠금`이면 코드 구현, `agent-task` 구현 계획 생성, 세부 API/파일 구조 확정을 시작하지 않는다. -- `구현 잠금`이 없거나 잠긴 Milestone은 사용자가 "진행"을 요청해도 우회하지 않고, 먼저 Milestone 문서의 구현 구체화와 잠금 해제를 사용자에게 요청한다. +- 선택된 Milestone의 `구현 잠금` 섹션이 없거나 상태가 `잠금`이면 코드 구현, `agent-task` 구현 계획 생성, 세부 API/파일 구조 확정을 시작하기 전에 현재 요청에 직접 영향을 주는 `결정 필요` 항목만 확인한다. +- 현재 요청과 직접 관련 없는 미정 항목은 잠금 상태로 남겨도 되며, 기존 구조/도메인 rule/플랫폼 관례로 정할 수 있는 작업은 표준선으로 기록하고 진행할 수 있다. +- Milestone 전체에서 사용자만 결정할 항목이 더 이상 없고 에이전트가 표준선에 따라 실행하면 되는 상태라면 `구현 잠금` 상태를 `해제`로 둔다. diff --git a/agent-ops/roadmap/current.md b/agent-ops/roadmap/current.md index 1ddea15..0d89c2f 100644 --- a/agent-ops/roadmap/current.md +++ b/agent-ops/roadmap/current.md @@ -11,4 +11,4 @@ - 요청 내용, 현재 브랜치, 변경 파일, 관련 코드 경로를 보고 가장 관련 있는 Milestone을 선택하고 같은 세션에서 1회 읽는다. - 활성 Milestone 둘 이상에 걸치면 필요한 Milestone 문서를 모두 읽고 작업 범위를 좁힌다. - 활성 Milestone 밖의 작업이면 `agent-ops/roadmap/ROADMAP.md`의 Milestone 목록을 확인하고 사용자에게 진행 또는 전환 여부를 확인한다. -- 선택된 Milestone의 `구현 잠금` 섹션이 없거나 상태가 `잠금`이면 구현이나 구현 계획을 시작하지 않고, 먼저 Milestone 구체화 업데이트와 잠금 해제를 사용자에게 요청한다. +- 선택된 Milestone의 `구현 잠금` 섹션이 없거나 상태가 `잠금`이면 구현이나 구현 계획을 시작하기 전에 현재 요청에 직접 영향을 주는 `결정 필요` 항목만 확인한다. 관련 결정이 없고 표준선으로 처리 가능하면 잠금을 유지한 채 진행할 수 있으며, Milestone 전체에서 사용자만 결정할 항목이 더 이상 없을 때만 `구현 잠금` 상태를 `해제`로 둔다. diff --git a/agent-ops/roadmap/milestones/external-integration.md b/agent-ops/roadmap/milestones/external-integration.md index 0f9502d..28a9042 100644 --- a/agent-ops/roadmap/milestones/external-integration.md +++ b/agent-ops/roadmap/milestones/external-integration.md @@ -15,16 +15,10 @@ External Integration ## 구현 잠금 - 상태: 잠금 -- 이유: IOP Responses 호출 요청/설정 경계와 direct model endpoint 재분류는 README와 현재 core 경로에 반영됐지만, Agent Integrator 대체 여부와 provider 결과 발행 정책은 아직 구현 가능한 수준으로 확정되지 않았다. -- 해제 조건: - - [x] IOP OpenAI API Responses-compatible 호출 요청/응답 계약과 설정 키가 문서화됨 - - [x] NomadCode가 직접 모델 런타임을 호출하지 않는 전환 기준이 문서화됨 - - [ ] Mattermost, Plane/Jira 결과 발행, Agent Integrator 대체 여부의 책임 경계가 정리됨 - - [ ] 사용자가 이 Milestone의 구현 구체화와 잠금 해제를 명시적으로 승인함 -- 잠금 중 금지: - - 코드 구현 또는 `agent-task` 구현 계획 생성 - - API/DB/package/file 구조를 추측해 확정 - - 세부 구현 체크리스트를 완료 기준처럼 작성 +- 결정 필요: + - [ ] Mattermost와 Plane/Jira 결과 발행의 책임 경계를 결정한다. + - [ ] Agent Integrator를 유지할지 IOP/A2A 또는 다른 연결 지점으로 대체할지 결정한다. + - [ ] A2A 도입 시점을 이 Milestone 범위로 둘지 후속으로 미룰지 결정한다. ## 범위 @@ -77,4 +71,4 @@ External Integration - `services/core/internal/scheduler/jobs.go`는 model client 결과를 task completion 결과로 저장한다. - `README.md`와 `services/core/README.md`는 NomadCode의 기본 실행 호출을 IOP Edge OpenAI-compatible Responses 경로로 정리하고, direct Ollama/model endpoint는 local development compatibility로 재분류한다. - sibling IOP repository의 로드맵은 OpenAI-compatible API를 외부 모델 기반 호출 표면으로, A2A를 외부 agent 작업 위임 표면으로, IOP native protocol을 운영 제어 표면으로 분리한다. -- 확인 필요: Agent Integrator의 유지/대체 여부, Mattermost 및 Plane/Jira 결과 발행 책임 경계, A2A 도입 시점은 구현 전 별도 확인이 필요하다. +- 확인 필요: `구현 잠금`의 결정 필요 항목 참고. diff --git a/agent-ops/roadmap/milestones/plane-task-pipeline-design.md b/agent-ops/roadmap/milestones/plane-task-pipeline-design.md index ad7fc4b..b8ae0ee 100644 --- a/agent-ops/roadmap/milestones/plane-task-pipeline-design.md +++ b/agent-ops/roadmap/milestones/plane-task-pipeline-design.md @@ -14,14 +14,9 @@ Work Item Provider Pipeline Design ## 구현 잠금 -- 상태: 해제 -- 이유: 이 Milestone은 이미 진행 중이며 provider-neutral pipeline 경계, 상태/projection 결정, 구현 검증 근거가 작업 컨텍스트에 누적되어 있다. 남은 작업은 같은 범위 안의 계약 보강과 구현 정리로 이어진다. -- 해제 조건: - - [x] provider-neutral 생성, adapter interface, board state와 agent 실행 상태 분리 방향이 문서화되어 있다. - - [x] 현재 구현 상태와 남은 결정 항목이 작업 컨텍스트에 정리되어 있다. - - [x] 사용자가 현재 활성 Milestone의 잠금 보정을 승인했다. -- 잠금 중 금지: - - 해당 없음 +- 상태: 잠금 +- 결정 필요: + - [ ] generic pipeline service 이후 provider registry/generic endpoint를 별도 작업으로 둘지, compatibility endpoint 유지 상태로 이번 Milestone 범위를 닫을지 결정한다. ## 범위 @@ -61,12 +56,17 @@ Work Item Provider Pipeline Design - [x] [handler-dto] HTTP handler가 Plane 타입과 직접 결합하지 않도록 요청 DTO와 변환 책임을 분리한다. - [x] [task-mapper] `buildPlaneCreateTaskInput`의 core task 생성 로직을 provider-neutral mapper로 분리한다. - [x] [compat-route] 기존 Plane endpoint는 compatibility entrypoint로 유지하거나 generic endpoint로 대체할지 결정한다. -- [ ] [projection-contract] provider-neutral 상태와 projection 계약을 확정한다. +- [x] [projection-contract] provider-neutral 상태와 projection 계약을 확정한다. - [x] [board-states] board state는 `backlog`, `todo`, `in_progress`, `testing`, `complete`, `cancel`로 둔다. - [x] [agent-metadata] `in_progress` 내부 agent 상태는 core task metadata를 canonical source로 둔다. - [x] [label-first] Plane/Jira provider projection은 label-first로 둔다. - [x] [desc-readonly] provider 본문(description)은 agent 실행 상태 저장소로 쓰지 않는다. - - [ ] [metadata-schema] core task metadata schema와 provider label mapping을 구현 대상으로 확정한다. + - [x] [metadata-schema] core task metadata schema와 provider label mapping을 구현 대상으로 확정한다. + - [x] [metadata-keys] canonical metadata key는 `agent_run_state`, `agent_phase`, `wait_type`, `status_reason`, `last_heartbeat_at`, `plan_ref`, `attempt`를 기본값으로 둔다. + - [x] [metadata-storage] canonical metadata 저장 위치는 신규 `tasks.metadata jsonb`로 확정한다. + - [x] [label-map] provider label mapping은 `agent:*`, `phase:*`, `worker:*` 공통 prefix를 우선 사용한다. + - [x] [label-index] phase label은 작업별 phase index와 맞춰 `phase:0`부터 증가시키고, 최신 phase는 core metadata의 `agent_phase`를 canonical 값으로 둔다. + - [x] [provider-delta] provider 공통 projection 계약을 먼저 맞추고 Plane/Jira 차이는 adapter별 capability와 mapping 설정으로 분리한다. - [ ] [result-policy] 완료/실패 결과의 provider comment/status update 정책을 정한다. - [x] [comment-prefix-sample] agent 단계 완료 기록은 prefix와 이모지가 있는 comment로 남기는 방향을 샘플 검증한다. - [ ] [comment-timing] comment prefix 세트와 작성 타이밍을 확정한다. @@ -107,7 +107,7 @@ Work Item Provider Pipeline Design - `buildPlaneCreateTaskInput`은 HTTP package 안에서 core task payload/external ref를 직접 조립하지 않고 `workitem.BuildCreateTaskInput`에 위임한다. - `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로 남긴다. - - DB schema는 `external_provider`, `external_id`, `external_url`, `external_metadata`로 provider-neutral 토대가 있으므로 유지한다. + - DB schema는 `external_provider`, `external_id`, `external_url`, `external_metadata`로 provider-neutral 토대를 유지하고, agent 실행 상태 snapshot은 `metadata jsonb`로 분리한다. - `services/core/internal/workitempipeline`은 provider-neutral reader, task mapper, `workflow.CreateTask` 호출을 묶는 generic pipeline service를 제공한다. - compatibility HTTP handler는 Plane request를 `workitem.Ref`로 변환한 뒤 `workitempipeline.Service`에 생성 흐름을 위임한다. - 초기 생성/연결 흐름: @@ -122,13 +122,17 @@ Work Item Provider Pipeline Design - 상태 설계 결정: - provider board state는 `backlog`, `todo`, `in_progress`, `testing`, `complete`, `cancel`의 소유권/검증 단계로 유지한다. - `in_progress`를 planning, implementing, review 같은 board state로 쪼개지 않고, agent 내부 실행 상태는 core task metadata에 canonical 값으로 저장한다. - - agent 내부 실행 상태 후보는 `agent_run_state`, `agent_phase`, `wait_type`, `status_reason`, `last_heartbeat_at`, `plan_ref`, `attempt`다. + - agent 내부 실행 상태의 기본 canonical metadata key는 `agent_run_state`, `agent_phase`, `wait_type`, `status_reason`, `last_heartbeat_at`, `plan_ref`, `attempt`다. - agent가 작업 중 사용자 판단을 기다릴 때는 board state를 `in_progress`로 유지하고 `agent_run_state=waiting_for_user`와 `wait_type`으로 멈춤 이유를 표현한다. - `testing`은 agent가 plan, 구현, 자체 테스트, 자체 리뷰 루프를 끝낸 뒤 사용자가 직접 테스트하는 단계다. - - provider projection은 라벨을 우선 사용한다. 예: `agent:waiting-user`, `phase:planning`, `agent:blocked`, `agent:failed`. + - provider projection은 라벨을 우선 사용한다. 예: `agent:waiting-user`, `agent:blocked`, `agent:failed`, `phase:0`, `phase:1`. + - phase label은 `agent-task/archive`의 작업별 loop index처럼 0부터 증가시키며, provider label은 외부 표시/검색용이고 최신 phase의 canonical source는 core metadata의 `agent_phase`로 둔다. - provider 본문(description)은 작업 요구사항과 맥락의 원본으로 보고, agent 실행 상태를 매번 갱신하는 저장소로 쓰지 않는다. - Plane custom property는 현재 NomadCode dev project에서 `is_issue_type_enabled=False`라 기본 경로로 전제하지 않는다. - - Plane 샘플 work item `NOMAD-13`은 `In Progress` state와 `agent:waiting-user`, `phase:planning` 라벨로 board state와 agent 내부 실행 상태 분리 방식을 보여준다. + - Plane 샘플 work item `NOMAD-13`은 `In Progress` state와 `agent:waiting-user`, `phase:planning` 라벨로 board state와 agent 내부 실행 상태 분리 방식을 보여준다. 신규 projection 계약에서는 phase label을 `phase:0`부터의 index형 label로 전환한다. + - canonical metadata 저장 위치는 신규 `tasks.metadata jsonb`로 확정한다. `payload`는 task 생성 입력/요구사항, `external_metadata`는 provider 연결 정보, `metadata`는 agent 실행 상태의 최신 snapshot으로 분리한다. + - `services/core/migrations/00003_add_task_metadata.sql`은 `tasks.metadata jsonb NOT NULL DEFAULT '{}'`를 추가한다. + - workflow create/update 경계는 invalid metadata JSON을 거부하고 빈 값 또는 `null`은 `{}`로 정규화한다. - trigger/auth 결정: - 사용자가 `backlog`에서 `todo`로 이동시키는 행위가 AI 작업 위임 의사로 간주될 수 있는 첫 관문이다. - core는 `todo` 상태와 agent 작업자 지정 신호가 함께 있을 때만 자동 실행 후보로 본다. @@ -149,6 +153,7 @@ Work Item Provider Pipeline Design - provider-neutral adapter contract는 `workitem.Ref`, `WorkItem`, `CommentInput`, `StatusProjection`, `LabelProjection`, `Mapping`으로 둔다. - provider별 API 인증, URL, workspace/project/issue 식별자, status id, label id는 adapter 설정과 metadata로만 다룬다. - core의 canonical 상태와 agent 실행 metadata는 provider에서 읽어온 값이 아니라 NomadCode task 상태의 진실로 둔다. + - provider별 차이는 공통 `workitem` projection DTO와 capability를 먼저 통과시킨 뒤, Plane/Jira adapter 내부 mapping과 optional facet 구현으로 분리한다. - comment 기록 규칙 후보: - `🧭 PLAN | <요약>`: plan 작성 또는 계획 검토 단계 완료 기록 - `🛠️ WORK | <요약>`: 구현 또는 문서 반영 단계 완료 기록 @@ -159,6 +164,6 @@ Work Item Provider Pipeline Design - 현재 상태는 `진행 중`이며, provider adapter interface 설계, Plane adapter facet 검증, provider-neutral task mapper, generic pipeline service 추출, HTTP compatibility endpoint의 Plane DTO 의존 제거가 완료됐다. - `agent-task/archive/2026/05/provider_neutral_plane_entrypoint/01_pipeline_mapper/complete.log`에서 `workitem.TaskCreateInput`과 `workitem.BuildCreateTaskInput` 추가 및 whitespace fallback 회귀 수정이 PASS로 정리됐다. - `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로 정리됐다. - - 다음 진행 후보는 core task metadata schema와 provider label mapping 확정, comment prefix 작성 타이밍, idempotency/retry 기준 정리다. - - 이후 남은 큰 결정은 core metadata schema, provider label mapping, comment prefix 작성 타이밍, idempotency/retry 기준이다. -- 확인 필요: generic pipeline service 이후 provider registry/generic endpoint를 별도 작업으로 볼지, compatibility endpoint 유지 상태로 이번 Milestone 범위를 닫을지 결정한다. + - 다음 진행 후보는 comment prefix 작성 타이밍, idempotency/retry 기준 정리다. + - 이후 남은 큰 결정은 comment prefix 작성 타이밍, idempotency/retry 기준이다. +- 확인 필요: `구현 잠금`의 결정 필요 항목 참고. diff --git a/agent-ops/roadmap/milestones/project-workspace-management-ux.md b/agent-ops/roadmap/milestones/project-workspace-management-ux.md index baeb5ec..7df972f 100644 --- a/agent-ops/roadmap/milestones/project-workspace-management-ux.md +++ b/agent-ops/roadmap/milestones/project-workspace-management-ux.md @@ -15,16 +15,10 @@ Project Workspace Management UX ## 구현 잠금 - 상태: 잠금 -- 이유: 프로젝트 제어 화면의 정보 구조, desktop 새창 범위, mobile 상단탭 구조, IOP 호출 결과 표시 방식이 아직 구현 가능한 수준으로 확정되지 않았다. -- 해제 조건: - - [ ] 프로젝트 단위 제어 화면의 공통 기능 목록이 문서화됨 - - [ ] desktop tab, desktop project window, mobile top-tab의 역할 경계가 문서화됨 - - [ ] IOP OpenAI API Responses-compatible 호출 결과와 task 상태를 UX에 투영하는 기준이 문서화됨 - - [ ] 사용자가 이 Milestone의 구현 구체화와 잠금 해제를 명시적으로 승인함 -- 잠금 중 금지: - - 코드 구현 또는 `agent-task` 구현 계획 생성 - - API/DB/package/file 구조를 추측해 확정 - - 세부 구현 체크리스트를 완료 기준처럼 작성 +- 결정 필요: + - [ ] 프로젝트 단위 제어 화면의 공통 기능 목록과 우선순위를 결정한다. + - [ ] desktop tab, desktop project window, mobile top-tab의 역할 경계를 결정한다. + - [ ] IOP OpenAI API Responses-compatible 호출 결과, task 상태, 파일 변경 결과를 UX에 투영하는 최소 기준을 결정한다. ## 범위 @@ -64,4 +58,4 @@ Project Workspace Management UX - 관련 경로: `apps/web/`, `apps/mobile/`, `packages/contracts/`, `services/core/` - 선행 작업: External Integration - 후속 작업: 확인 필요 -- 확인 필요: 프로젝트 제어 화면의 공통 기능 목록, desktop 새창의 포함 범위, mobile 상단탭의 탭 단위 +- 확인 필요: `구현 잠금`의 결정 필요 항목 참고. diff --git a/agent-ops/roadmap/milestones/workflow-core.md b/agent-ops/roadmap/milestones/workflow-core.md index 67dfa89..90d75ad 100644 --- a/agent-ops/roadmap/milestones/workflow-core.md +++ b/agent-ops/roadmap/milestones/workflow-core.md @@ -15,16 +15,9 @@ Workflow Core ## 구현 잠금 - 상태: 잠금 -- 이유: Workflow Core는 후속 계획 Milestone이며, task lifecycle, retry, timeout, notification event의 책임 경계와 검증 기준이 아직 구현 가능한 수준으로 구체화되지 않았다. -- 해제 조건: - - [ ] Work Item Provider Pipeline Design에서 넘길 상태 변화와 실패 케이스가 확정됨 - - [ ] NomadCode workflow retry/timeout과 IOP 실행 정책의 경계가 문서화됨 - - [ ] task lifecycle 상태 모델과 notification event 발행 기준이 문서화됨 - - [ ] 사용자가 이 Milestone의 구현 구체화와 잠금 해제를 명시적으로 승인함 -- 잠금 중 금지: - - 코드 구현 또는 `agent-task` 구현 계획 생성 - - API/DB/package/file 구조를 추측해 확정 - - 세부 구현 체크리스트를 완료 기준처럼 작성 +- 결정 필요: + - [ ] NomadCode workflow retry/timeout과 IOP 실행 정책의 책임 경계를 결정한다. + - [ ] notification event의 소비자와 최소 발행 범위를 결정한다. ## 범위 @@ -63,3 +56,4 @@ Workflow Core - 선행 작업: Plane Communication Foundation, Work Item Provider Pipeline Design - 후속 작업: External Integration - 검증 기준: scheduler, persistence, adapter behavior를 크게 바꾸면 관련 테스트를 함께 갱신한다. +- 확인 필요: `구현 잠금`의 결정 필요 항목 참고. diff --git a/services/core/internal/db/models.go b/services/core/internal/db/models.go index e94008b..de39154 100644 --- a/services/core/internal/db/models.go +++ b/services/core/internal/db/models.go @@ -23,4 +23,5 @@ type Task struct { ExternalID *string `json:"external_id"` ExternalUrl *string `json:"external_url"` ExternalMetadata json.RawMessage `json:"external_metadata"` + Metadata json.RawMessage `json:"metadata"` } diff --git a/services/core/internal/db/tasks.sql.go b/services/core/internal/db/tasks.sql.go index b3f9925..559a50b 100644 --- a/services/core/internal/db/tasks.sql.go +++ b/services/core/internal/db/tasks.sql.go @@ -14,7 +14,7 @@ const completeTask = `-- name: CompleteTask :one UPDATE tasks SET status = 'completed', result = $2, error = NULL, updated_at = now() WHERE id::text = $1 -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata ` type CompleteTaskParams struct { @@ -39,20 +39,22 @@ func (q *Queries) CompleteTask(ctx context.Context, arg CompleteTaskParams) (Tas &i.ExternalID, &i.ExternalUrl, &i.ExternalMetadata, + &i.Metadata, ) return i, err } const createTask = `-- name: CreateTask :one -INSERT INTO tasks (title, source, payload, external_provider, external_id, external_url, external_metadata) -VALUES ($1, $2, $3, $4, $5, $6, $7) -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +INSERT INTO tasks (title, source, payload, metadata, external_provider, external_id, external_url, external_metadata) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8) +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata ` type CreateTaskParams struct { Title string `json:"title"` Source string `json:"source"` Payload json.RawMessage `json:"payload"` + Metadata json.RawMessage `json:"metadata"` ExternalProvider *string `json:"external_provider"` ExternalID *string `json:"external_id"` ExternalUrl *string `json:"external_url"` @@ -64,6 +66,7 @@ func (q *Queries) CreateTask(ctx context.Context, arg CreateTaskParams) (Task, e arg.Title, arg.Source, arg.Payload, + arg.Metadata, arg.ExternalProvider, arg.ExternalID, arg.ExternalUrl, @@ -84,6 +87,7 @@ func (q *Queries) CreateTask(ctx context.Context, arg CreateTaskParams) (Task, e &i.ExternalID, &i.ExternalUrl, &i.ExternalMetadata, + &i.Metadata, ) return i, err } @@ -92,7 +96,7 @@ const failTask = `-- name: FailTask :one UPDATE tasks SET status = 'failed', error = $2, updated_at = now() WHERE id::text = $1 -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata ` type FailTaskParams struct { @@ -117,12 +121,13 @@ func (q *Queries) FailTask(ctx context.Context, arg FailTaskParams) (Task, error &i.ExternalID, &i.ExternalUrl, &i.ExternalMetadata, + &i.Metadata, ) return i, err } const getTask = `-- name: GetTask :one -SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata FROM tasks WHERE id::text = $1 ` @@ -144,12 +149,13 @@ func (q *Queries) GetTask(ctx context.Context, id string) (Task, error) { &i.ExternalID, &i.ExternalUrl, &i.ExternalMetadata, + &i.Metadata, ) return i, err } const listTasks = `-- name: ListTasks :many -SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata FROM tasks ORDER BY created_at DESC LIMIT $1 @@ -178,6 +184,7 @@ func (q *Queries) ListTasks(ctx context.Context, limit int32) ([]Task, error) { &i.ExternalID, &i.ExternalUrl, &i.ExternalMetadata, + &i.Metadata, ); err != nil { return nil, err } @@ -189,11 +196,45 @@ func (q *Queries) ListTasks(ctx context.Context, limit int32) ([]Task, error) { return items, nil } +const updateTaskMetadata = `-- name: UpdateTaskMetadata :one +UPDATE tasks +SET metadata = $2, updated_at = now() +WHERE id::text = $1 +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata +` + +type UpdateTaskMetadataParams struct { + ID string `json:"id"` + Metadata json.RawMessage `json:"metadata"` +} + +func (q *Queries) UpdateTaskMetadata(ctx context.Context, arg UpdateTaskMetadataParams) (Task, error) { + row := q.db.QueryRow(ctx, updateTaskMetadata, arg.ID, arg.Metadata) + var i Task + err := row.Scan( + &i.ID, + &i.Title, + &i.Source, + &i.Status, + &i.Payload, + &i.Result, + &i.Error, + &i.CreatedAt, + &i.UpdatedAt, + &i.ExternalProvider, + &i.ExternalID, + &i.ExternalUrl, + &i.ExternalMetadata, + &i.Metadata, + ) + return i, err +} + const updateTaskStatus = `-- name: UpdateTaskStatus :one UPDATE tasks SET status = $2, updated_at = now() WHERE id::text = $1 -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata ` type UpdateTaskStatusParams struct { @@ -218,6 +259,7 @@ func (q *Queries) UpdateTaskStatus(ctx context.Context, arg UpdateTaskStatusPara &i.ExternalID, &i.ExternalUrl, &i.ExternalMetadata, + &i.Metadata, ) return i, err } diff --git a/services/core/internal/storage/store.go b/services/core/internal/storage/store.go index 0c5feb5..02a763e 100644 --- a/services/core/internal/storage/store.go +++ b/services/core/internal/storage/store.go @@ -14,6 +14,7 @@ type CreateTaskInput struct { Title string Source string Payload json.RawMessage + Metadata json.RawMessage ExternalProvider *string ExternalID *string ExternalURL *string @@ -29,6 +30,11 @@ func NewStore(pool *pgxpool.Pool) *Store { } func (s *Store) CreateTask(ctx context.Context, input CreateTaskInput) (db.Task, error) { + metadata := input.Metadata + if len(metadata) == 0 || string(metadata) == "null" { + metadata = json.RawMessage(`{}`) + } + externalMetadata := input.ExternalMetadata if len(externalMetadata) == 0 || string(externalMetadata) == "null" { externalMetadata = json.RawMessage(`{}`) @@ -38,6 +44,7 @@ func (s *Store) CreateTask(ctx context.Context, input CreateTaskInput) (db.Task, Title: input.Title, Source: input.Source, Payload: input.Payload, + Metadata: metadata, ExternalProvider: input.ExternalProvider, ExternalID: input.ExternalID, ExternalUrl: input.ExternalURL, @@ -53,6 +60,16 @@ func (s *Store) ListTasks(ctx context.Context, limit int32) ([]db.Task, error) { return s.queries.ListTasks(ctx, limit) } +func (s *Store) UpdateMetadata(ctx context.Context, id string, metadata json.RawMessage) (db.Task, error) { + if len(metadata) == 0 || string(metadata) == "null" { + metadata = json.RawMessage(`{}`) + } + return s.queries.UpdateTaskMetadata(ctx, db.UpdateTaskMetadataParams{ + ID: id, + Metadata: metadata, + }) +} + func (s *Store) UpdateStatus(ctx context.Context, id, status string) (db.Task, error) { return s.queries.UpdateTaskStatus(ctx, db.UpdateTaskStatusParams{ ID: id, diff --git a/services/core/internal/workflow/model.go b/services/core/internal/workflow/model.go index da27fba..e101410 100644 --- a/services/core/internal/workflow/model.go +++ b/services/core/internal/workflow/model.go @@ -17,6 +17,7 @@ type CreateTaskInput struct { Title string `json:"title"` Source string `json:"source"` Payload json.RawMessage `json:"payload"` + Metadata json.RawMessage `json:"metadata,omitempty"` External *ExternalRefInput `json:"external,omitempty"` } @@ -26,3 +27,13 @@ type ExternalRefInput struct { URL string `json:"url"` Metadata json.RawMessage `json:"metadata"` } + +const ( + MetadataKeyAgentRunState = "agent_run_state" + MetadataKeyAgentPhase = "agent_phase" + MetadataKeyWaitType = "wait_type" + MetadataKeyStatusReason = "status_reason" + MetadataKeyLastHeartbeat = "last_heartbeat_at" + MetadataKeyPlanRef = "plan_ref" + MetadataKeyAttempt = "attempt" +) diff --git a/services/core/internal/workflow/service.go b/services/core/internal/workflow/service.go index e4e0f32..7561b73 100644 --- a/services/core/internal/workflow/service.go +++ b/services/core/internal/workflow/service.go @@ -48,6 +48,11 @@ func (s *Service) CreateTask(ctx context.Context, input CreateTaskInput) (storag return storage.Task{}, ErrInvalidTaskInput } + metadata, err := NormalizeTaskMetadata(input.Metadata) + if err != nil { + return storage.Task{}, err + } + external, err := NormalizeExternalRef(input.External) if err != nil { return storage.Task{}, err @@ -57,6 +62,7 @@ func (s *Service) CreateTask(ctx context.Context, input CreateTaskInput) (storag Title: title, Source: source, Payload: payload, + Metadata: metadata, ExternalProvider: external.Provider, ExternalID: external.ID, ExternalURL: external.URL, @@ -110,6 +116,16 @@ func optionalString(value string) *string { return &value } +func NormalizeTaskMetadata(metadata json.RawMessage) (json.RawMessage, error) { + if len(metadata) == 0 || string(metadata) == "null" { + return json.RawMessage(`{}`), nil + } + if !json.Valid(metadata) { + return nil, ErrInvalidTaskInput + } + return metadata, nil +} + func (s *Service) ListTasks(ctx context.Context, limit int32) ([]storage.Task, error) { if limit <= 0 { limit = 20 @@ -120,6 +136,14 @@ func (s *Service) ListTasks(ctx context.Context, limit int32) ([]storage.Task, e return s.store.ListTasks(ctx, limit) } +func (s *Service) UpdateTaskMetadata(ctx context.Context, id string, metadata json.RawMessage) (storage.Task, error) { + normalized, err := NormalizeTaskMetadata(metadata) + if err != nil { + return storage.Task{}, err + } + return s.store.UpdateMetadata(ctx, id, normalized) +} + func (s *Service) EnqueueTask(ctx context.Context, id string) (storage.Task, error) { task, err := s.store.GetTask(ctx, id) if err != nil { diff --git a/services/core/internal/workflow/service_test.go b/services/core/internal/workflow/service_test.go index 9ef4097..e7b54b8 100644 --- a/services/core/internal/workflow/service_test.go +++ b/services/core/internal/workflow/service_test.go @@ -46,3 +46,28 @@ func TestNormalizeExternalRefRequiresProvider(t *testing.T) { t.Fatal("expected provider error") } } + +func TestNormalizeTaskMetadataDefaultsAndKeepsValidJSON(t *testing.T) { + metadata, err := NormalizeTaskMetadata(nil) + if err != nil { + t.Fatalf("NormalizeTaskMetadata returned error: %v", err) + } + if string(metadata) != "{}" { + t.Fatalf("unexpected default metadata: %s", metadata) + } + + metadata, err = NormalizeTaskMetadata(json.RawMessage(`{"agent_phase":0}`)) + if err != nil { + t.Fatalf("NormalizeTaskMetadata returned error: %v", err) + } + if string(metadata) != `{"agent_phase":0}` { + t.Fatalf("unexpected metadata: %s", metadata) + } +} + +func TestNormalizeTaskMetadataRejectsInvalidJSON(t *testing.T) { + _, err := NormalizeTaskMetadata(json.RawMessage(`{"agent_phase"`)) + if err == nil { + t.Fatal("expected invalid metadata error") + } +} diff --git a/services/core/migrations/00003_add_task_metadata.sql b/services/core/migrations/00003_add_task_metadata.sql new file mode 100644 index 0000000..9786bb4 --- /dev/null +++ b/services/core/migrations/00003_add_task_metadata.sql @@ -0,0 +1,7 @@ +-- +goose Up +ALTER TABLE tasks + ADD COLUMN IF NOT EXISTS metadata jsonb NOT NULL DEFAULT '{}'; + +-- +goose Down +ALTER TABLE tasks + DROP COLUMN IF EXISTS metadata; diff --git a/services/core/queries/tasks.sql b/services/core/queries/tasks.sql index 4769a4a..f775b0c 100644 --- a/services/core/queries/tasks.sql +++ b/services/core/queries/tasks.sql @@ -1,15 +1,15 @@ -- name: CreateTask :one -INSERT INTO tasks (title, source, payload, external_provider, external_id, external_url, external_metadata) -VALUES ($1, $2, $3, $4, $5, $6, $7) -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata; +INSERT INTO tasks (title, source, payload, metadata, external_provider, external_id, external_url, external_metadata) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8) +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata; -- name: GetTask :one -SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata FROM tasks WHERE id::text = $1; -- name: ListTasks :many -SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata +SELECT id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata FROM tasks ORDER BY created_at DESC LIMIT $1; @@ -18,16 +18,22 @@ LIMIT $1; UPDATE tasks SET status = $2, updated_at = now() WHERE id::text = $1 -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata; +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata; + +-- name: UpdateTaskMetadata :one +UPDATE tasks +SET metadata = $2, updated_at = now() +WHERE id::text = $1 +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata; -- name: CompleteTask :one UPDATE tasks SET status = 'completed', result = $2, error = NULL, updated_at = now() WHERE id::text = $1 -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata; +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata; -- name: FailTask :one UPDATE tasks SET status = 'failed', error = $2, updated_at = now() WHERE id::text = $1 -RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata; +RETURNING id, title, source, status, payload, result, error, created_at, updated_at, external_provider, external_id, external_url, external_metadata, metadata;