Add final result reduction
coordinator / test (push) Canceled after 0s

This commit is contained in:
Emil
2026-07-24 14:58:18 +03:00
parent 6e67daa9eb
commit a055473706
28 changed files with 849 additions and 54 deletions
+14 -4
View File
@@ -100,7 +100,7 @@ func (uc *CancelJob) Execute(ctx context.Context, jobID uuid.UUID) (int64, error
if job.Status == domain.JobCancelled {
return nil
}
if job.Status == domain.JobCompleted || job.Status == domain.JobFailed {
if job.Status == domain.JobReducing || job.Status == domain.JobCompleted || job.Status == domain.JobFailed {
return domain.ErrJobNotCancellable
}
// The lease reaper can be the transition that exhausted the final task.
@@ -111,7 +111,7 @@ func (uc *CancelJob) Execute(ctx context.Context, jobID uuid.UUID) (int64, error
return err
}
derived := progressFrom(*job, counts).DeriveStatus()
if derived == domain.JobCompleted || derived == domain.JobFailed {
if derived == domain.JobReducing || derived == domain.JobCompleted || derived == domain.JobFailed {
return domain.ErrJobNotCancellable
}
cancelled, err = uc.tasks.CancelByJob(ctx, jobID, now)
@@ -226,10 +226,20 @@ func syncJobStatus(ctx context.Context, jobs JobRepository, tasks TaskRepository
if err != nil {
return err
}
status := progressFrom(domain.Job{}, counts).DeriveStatus()
job, err := jobs.Get(ctx, jobID)
if err != nil {
return err
}
status := progressFrom(*job, counts).DeriveStatus()
// All worker shards being complete means scientific reduction is ready, not
// that the job's final artifact already exists. CTX-09 owns the transition
// from reducing to completed after it persists that artifact.
if status == domain.JobCompleted && job.Workload == "similarity-search" {
status = domain.JobReducing
}
var completedAt *time.Time
if status == domain.JobCompleted || status == domain.JobFailed {
if status == domain.JobFailed {
completedAt = &now
}
return jobs.UpdateStatus(ctx, jobID, status, completedAt)
+3
View File
@@ -70,6 +70,9 @@ type JobRepository interface {
Insert(ctx context.Context, j *domain.Job) error
Get(ctx context.Context, id uuid.UUID) (*domain.Job, error)
UpdateStatus(ctx context.Context, id uuid.UUID, status domain.JobStatus, completedAt *time.Time) error
ClaimReduction(ctx context.Context, id uuid.UUID, startedAt time.Time) (bool, error)
CompleteWithResult(ctx context.Context, id, resultArtifactID uuid.UUID, completedAt time.Time) error
FailReduction(ctx context.Context, id uuid.UUID, code, message string, completedAt time.Time) error
}
// WorkerRepository persists the worker registry.
+133
View File
@@ -0,0 +1,133 @@
package usecase
import (
"bytes"
"context"
"io"
"github.com/google/uuid"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
"github.com/emil28092005/SciMesh/coordinator/internal/reducer"
)
// ReduceJob turns completed coordinator-owned partial artifacts into one final
// artifact. It performs no worker I/O and never trusts a worker URI or path.
type ReduceJob struct {
jobs JobRepository
tasks TaskRepository
artifacts ArtifactRepository
blobs BlobStore
tx TxManager
clock Clock
}
func NewReduceJob(jobs JobRepository, tasks TaskRepository, artifacts ArtifactRepository,
blobs BlobStore, tx TxManager, clock Clock) *ReduceJob {
return &ReduceJob{jobs: jobs, tasks: tasks, artifacts: artifacts, blobs: blobs, tx: tx, clock: clock}
}
// Execute is idempotent for jobs that are not currently reducing. The worker
// completion path may call it after every result; only the last task changes a
// similarity-search job into reducing state.
func (uc *ReduceJob) Execute(ctx context.Context, jobID uuid.UUID) error {
claimed, err := uc.jobs.ClaimReduction(ctx, jobID, uc.clock.Now())
if err != nil || !claimed {
return err
}
job, err := uc.jobs.Get(ctx, jobID)
if err != nil {
return err
}
if job.Status != domain.JobReducing {
return nil
}
if job.Workload != "similarity-search" {
return uc.fail(ctx, jobID)
}
completed, err := uc.tasks.ListCompleted(ctx, jobID)
if err != nil {
return uc.fail(ctx, jobID)
}
if len(completed) == 0 {
return uc.fail(ctx, jobID)
}
readers := make([]io.Reader, 0, len(completed))
closers := make([]io.Closer, 0, len(completed))
for _, task := range completed {
if task.ResultArtifactID == nil {
closeAll(closers)
return uc.fail(ctx, jobID)
}
artifact, err := uc.artifacts.Get(ctx, *task.ResultArtifactID)
if err != nil || artifact.JobID != jobID || artifact.TaskID == nil || *artifact.TaskID != task.ID || artifact.Kind != domain.ArtifactPartialResult {
closeAll(closers)
return uc.fail(ctx, jobID)
}
body, err := uc.blobs.Open(ctx, artifact.StorageKey)
if err != nil {
closeAll(closers)
return uc.fail(ctx, jobID)
}
readers = append(readers, body)
closers = append(closers, body)
}
output, reduceErr := reducer.ReduceSimilaritySearch(readers, job.Parameters)
closeAll(closers)
if reduceErr != nil {
return uc.fail(ctx, jobID)
}
final, err := domain.NewArtifact(jobID, nil, domain.ArtifactFinalResult, "similarity-search.csv", "text/csv", uc.clock.Now())
if err != nil {
return uc.fail(ctx, jobID)
}
sum, size, err := uc.blobs.Put(ctx, final.StorageKey, bytes.NewReader(output))
if err != nil {
return uc.fail(ctx, jobID)
}
final.SetContent(sum, size)
if err := uc.tx.WithinTx(ctx, func(ctx context.Context) error {
if err := uc.artifacts.Insert(ctx, final); err != nil {
return err
}
return uc.jobs.CompleteWithResult(ctx, jobID, final.ID, uc.clock.Now())
}); err != nil {
_ = uc.blobs.Delete(ctx, final.StorageKey)
return err
}
return nil
}
func (uc *ReduceJob) fail(ctx context.Context, jobID uuid.UUID) error {
// The public state carries a stable sanitized failure, never parser/storage
// internals that may include local paths or implementation details.
return uc.jobs.FailReduction(ctx, jobID, "reducer_failed", "final result reduction failed", uc.clock.Now())
}
func closeAll(closers []io.Closer) {
for _, closer := range closers {
_ = closer.Close()
}
}
type GetJobResult struct {
jobs JobRepository
download *DownloadArtifact
}
func NewGetJobResult(jobs JobRepository, download *DownloadArtifact) *GetJobResult {
return &GetJobResult{jobs: jobs, download: download}
}
func (uc *GetJobResult) Execute(ctx context.Context, jobID uuid.UUID) (*domain.Artifact, io.ReadCloser, error) {
job, err := uc.jobs.Get(ctx, jobID)
if err != nil {
return nil, nil, err
}
if job.Status != domain.JobCompleted || job.ResultArtifactID == nil {
return nil, nil, domain.ErrArtifactNotFound
}
return uc.download.Execute(ctx, *job.ResultArtifactID)
}
@@ -55,6 +55,8 @@ type harness struct {
getInput *usecase.GetTaskInput
expire *usecase.ExpireLeases
cancel *usecase.CancelJob
reduce *usecase.ReduceJob
jobResult *usecase.GetJobResult
}
func newHarness() *harness {
@@ -81,9 +83,88 @@ func newHarness() *harness {
h.getInput = usecase.NewGetTaskInput(h.tasks, h.arts, h.blobs)
h.expire = usecase.NewExpireLeases(h.tasks, h.jobs, tx, h.clk)
h.cancel = usecase.NewCancelJob(h.jobs, h.tasks, tx, h.clk)
h.reduce = usecase.NewReduceJob(h.jobs, h.tasks, h.arts, h.blobs, tx, h.clk)
h.jobResult = usecase.NewGetJobResult(h.jobs, h.downloadArt)
return h
}
func TestSimilaritySearchReductionCreatesFinalArtifact(t *testing.T) {
h := newHarness()
jobID := h.seedJob(t, "similarity-search", 2)
_, err := h.jobs.Get(ctx, jobID)
if err != nil {
t.Fatal(err)
}
if err := h.jobs.UpdateStatus(ctx, jobID, domain.JobRunning, nil); err != nil {
t.Fatal(err)
}
partials := []string{
"rank,chembl_id,canonical_smiles,similarity\n1,B,CCC,0.50000048\n",
"rank,chembl_id,canonical_smiles,similarity\n1,A,CC,0.50000049\n",
}
for _, partial := range partials {
taskID, attempt := h.leaseOne(t, "w1", "similarity-search")
art, err := h.uploadArt.Execute(ctx, usecase.UploadArtifactInput{TaskID: taskID, WorkerID: "w1", Attempt: attempt, Filename: "partial.csv", ContentType: "text/csv", Body: strings.NewReader(partial)})
if err != nil {
t.Fatal(err)
}
if _, err := h.complete.Execute(ctx, usecase.CompleteTaskInput{TaskID: taskID, WorkerID: "w1", Attempt: attempt, ResultArtifactID: art.ID}); err != nil {
t.Fatal(err)
}
}
if err := h.reduce.Execute(ctx, jobID); err != nil {
t.Fatal(err)
}
progress, err := h.status.Execute(ctx, jobID)
if err != nil || progress.Job.Status != domain.JobCompleted {
t.Fatalf("status=%s err=%v", progress.Job.Status, err)
}
art, body, err := h.jobResult.Execute(ctx, jobID)
if err != nil {
t.Fatal(err)
}
defer body.Close()
bytes, _ := io.ReadAll(body)
if art.Kind != domain.ArtifactFinalResult || string(bytes) != "rank,chembl_id,canonical_smiles,similarity\n1,A,CC,0.500000\n2,B,CCC,0.500000\n" {
t.Fatalf("unexpected final %q", bytes)
}
}
func TestSimilaritySearchReductionFailureIsSanitized(t *testing.T) {
h := newHarness()
jobID := h.seedJob(t, "similarity-search", 1)
if err := h.jobs.UpdateStatus(ctx, jobID, domain.JobRunning, nil); err != nil {
t.Fatal(err)
}
taskID, attempt := h.leaseOne(t, "w1", "similarity-search")
art, err := h.uploadArt.Execute(ctx, usecase.UploadArtifactInput{
TaskID: taskID, WorkerID: "w1", Attempt: attempt, Filename: "partial.csv",
ContentType: "text/csv", Body: strings.NewReader("rank,chembl_id,canonical_smiles,similarity\n2,A,CC,0.9\n"),
})
if err != nil {
t.Fatal(err)
}
if _, err := h.complete.Execute(ctx, usecase.CompleteTaskInput{
TaskID: taskID, WorkerID: "w1", Attempt: attempt, ResultArtifactID: art.ID,
}); err != nil {
t.Fatal(err)
}
if err := h.reduce.Execute(ctx, jobID); err != nil {
t.Fatal(err)
}
job, err := h.jobs.Get(ctx, jobID)
if err != nil {
t.Fatal(err)
}
if job.Status != domain.JobFailed || job.ErrorCode == nil || *job.ErrorCode != "reducer_failed" ||
job.ErrorMessage == nil || *job.ErrorMessage != "final result reduction failed" {
t.Fatalf("unexpected failed job: %+v", job)
}
if job.ResultArtifactID != nil {
t.Fatal("failed reduction must not expose a final result")
}
}
// seedJob creates a URI-chunked job with n chunks and returns its id.
func (h *harness) seedJob(t *testing.T, workload string, n int) uuid.UUID {
t.Helper()