complete item 1 task substrate
This commit is contained in:
@@ -0,0 +1,141 @@
|
||||
package domain
|
||||
|
||||
import (
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"encoding/base32"
|
||||
"encoding/binary"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
var ErrConflict = errors.New("task version conflict")
|
||||
var ErrNotFound = errors.New("task not found")
|
||||
var ErrInvalid = errors.New("invalid event")
|
||||
|
||||
type TaskState string
|
||||
|
||||
const (
|
||||
StateQueued TaskState = "queued"
|
||||
StateLeased TaskState = "leased"
|
||||
StateCompleted TaskState = "completed"
|
||||
StateFailed TaskState = "failed"
|
||||
StateBlocked TaskState = "blocked"
|
||||
)
|
||||
|
||||
type Estimate struct {
|
||||
Value float64 `json:"value"`
|
||||
Who string `json:"who"`
|
||||
Confidence float64 `json:"confidence"`
|
||||
}
|
||||
type Lease struct {
|
||||
HarnessID string `json:"harness_id"`
|
||||
Until time.Time `json:"until"`
|
||||
}
|
||||
type Task struct {
|
||||
ID string `json:"id"`
|
||||
Source string `json:"source"`
|
||||
ExternalID string `json:"external_id"`
|
||||
Project string `json:"project"`
|
||||
Capability []string `json:"capability"`
|
||||
Parent string `json:"parent,omitempty"`
|
||||
InherentPriority int `json:"inherent_priority"`
|
||||
Due *time.Time `json:"due,omitempty"`
|
||||
Estimate *Estimate `json:"estimate,omitempty"`
|
||||
State TaskState `json:"state"`
|
||||
Lease *Lease `json:"lease,omitempty"`
|
||||
Version int `json:"version"`
|
||||
Title string `json:"title,omitempty"`
|
||||
}
|
||||
|
||||
type Event struct {
|
||||
Seq uint64 `json:"seq"`
|
||||
ID string `json:"id"`
|
||||
Type string `json:"type"`
|
||||
TaskID string `json:"task_id"`
|
||||
Version int `json:"version"`
|
||||
At time.Time `json:"at"`
|
||||
Payload json.RawMessage `json:"payload"`
|
||||
}
|
||||
|
||||
func Hash(v []byte) string { h := sha256.Sum256(v); return hex.EncodeToString(h[:]) }
|
||||
|
||||
// NewID returns a sortable, 128-bit ULID-like identifier using the canonical
|
||||
// 48-bit millisecond timestamp plus 80 bits of cryptographic randomness.
|
||||
|
||||
var ulidEncoding = base32.NewEncoding("0123456789ABCDEFGHJKMNPQRSTVWXYZ").WithPadding(base32.NoPadding)
|
||||
|
||||
func NewID() string {
|
||||
b := make([]byte, 16)
|
||||
binary.BigEndian.PutUint64(b[:8], uint64(time.Now().UnixMilli())<<16)
|
||||
_, _ = rand.Read(b[6:])
|
||||
return ulidEncoding.EncodeToString(b)
|
||||
}
|
||||
func ValidateEvent(e Event) error {
|
||||
if e.Type == "" || e.TaskID == "" || len(e.Payload) == 0 || len(e.Payload) > 64*1024 {
|
||||
return ErrInvalid
|
||||
}
|
||||
allowed := map[string]bool{"TaskCreated": true, "TaskLeased": true, "TaskReleased": true, "TaskCompleted": true, "TaskFailed": true, "TaskBlocked": true, "ApprovalRequested": true, "ApprovalGranted": true, "ApprovalDenied": true, "TaskAmended": true}
|
||||
if !allowed[e.Type] {
|
||||
return fmt.Errorf("%w: unknown type %q", ErrInvalid, e.Type)
|
||||
}
|
||||
var p map[string]any
|
||||
if err := json.Unmarshal(e.Payload, &p); err != nil {
|
||||
return fmt.Errorf("%w: payload is not JSON", ErrInvalid)
|
||||
}
|
||||
return ValidatePayload(e.Type, p)
|
||||
}
|
||||
func ValidateCreated(p map[string]any) error {
|
||||
for _, k := range []string{"source", "external_id", "project"} {
|
||||
if s, ok := p[k].(string); !ok || strings.TrimSpace(s) == "" {
|
||||
return fmt.Errorf("%w: %s required", ErrInvalid, k)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func ValidatePayload(typ string, p map[string]any) error {
|
||||
requiredString := func(key string) error {
|
||||
v, ok := p[key].(string)
|
||||
if !ok || strings.TrimSpace(v) == "" {
|
||||
return fmt.Errorf("%w: %s required", ErrInvalid, key)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
switch typ {
|
||||
case "TaskCreated":
|
||||
return ValidateCreated(p)
|
||||
case "TaskLeased":
|
||||
if err := requiredString("harness_id"); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, ok := p["until_ns"].(float64); !ok {
|
||||
return fmt.Errorf("%w: until_ns required", ErrInvalid)
|
||||
}
|
||||
case "TaskReleased":
|
||||
if err := requiredString("handoff_ref"); err != nil && p["reason"] == nil {
|
||||
return err
|
||||
}
|
||||
case "TaskCompleted":
|
||||
if err := requiredString("report_ref"); err != nil {
|
||||
return err
|
||||
}
|
||||
case "TaskFailed":
|
||||
if err := requiredString("reason"); err != nil {
|
||||
return err
|
||||
}
|
||||
case "TaskBlocked":
|
||||
if err := requiredString("blocker"); err != nil {
|
||||
return err
|
||||
}
|
||||
case "TaskAmended":
|
||||
if len(p) == 0 {
|
||||
return fmt.Errorf("%w: amendment cannot be empty", ErrInvalid)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
package provider
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"orchestra/internal/domain"
|
||||
)
|
||||
|
||||
type Sink interface{ Append(domain.Event) error }
|
||||
type Provider interface {
|
||||
Ingest(io.Reader, Sink) (int, error)
|
||||
}
|
||||
|
||||
// JSONL treats each line as an external task object. Replaying the same input
|
||||
// is safe because the store deduplicates the stable source/external_id key.
|
||||
type JSONL struct{}
|
||||
|
||||
func (JSONL) Ingest(r io.Reader, sink Sink) (int, error) {
|
||||
sc := bufio.NewScanner(r)
|
||||
count := 0
|
||||
line := 0
|
||||
for sc.Scan() {
|
||||
line++
|
||||
raw := sc.Bytes()
|
||||
if len(raw) == 0 {
|
||||
continue
|
||||
}
|
||||
var p map[string]any
|
||||
if err := json.Unmarshal(raw, &p); err != nil {
|
||||
return count, fmt.Errorf("line %d: %w", line, err)
|
||||
}
|
||||
if err := domain.ValidateCreated(p); err != nil {
|
||||
return count, fmt.Errorf("line %d: %w", line, err)
|
||||
}
|
||||
b, _ := json.Marshal(p)
|
||||
e := domain.Event{ID: domain.NewID(), Type: "TaskCreated", TaskID: domain.NewID(), Version: 1, Payload: b}
|
||||
if err := sink.Append(e); err != nil {
|
||||
return count, fmt.Errorf("line %d: %w", line, err)
|
||||
}
|
||||
count++
|
||||
}
|
||||
return count, sc.Err()
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
package provider
|
||||
|
||||
import (
|
||||
"orchestra/internal/domain"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
type sink struct{ events []domain.Event }
|
||||
|
||||
func (s *sink) Append(e domain.Event) error { s.events = append(s.events, e); return nil }
|
||||
func TestJSONLIngest(t *testing.T) {
|
||||
s := &sink{}
|
||||
n, err := (JSONL{}).Ingest(strings.NewReader("{\"source\":\"local\",\"external_id\":\"1\",\"project\":\"demo\"}\n"), s)
|
||||
if err != nil || n != 1 || len(s.events) != 1 {
|
||||
t.Fatalf("n=%d events=%d err=%v", n, len(s.events), err)
|
||||
}
|
||||
}
|
||||
func TestJSONLRejectsMalformedLine(t *testing.T) {
|
||||
n, err := (JSONL{}).Ingest(strings.NewReader("not-json\n"), &sink{})
|
||||
if err == nil || n != 0 {
|
||||
t.Fatalf("n=%d err=%v", n, err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,245 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"orchestra/internal/domain"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Store struct {
|
||||
mu sync.Mutex
|
||||
path string
|
||||
cas string
|
||||
events []domain.Event
|
||||
tasks map[string]domain.Task
|
||||
external map[string]string
|
||||
snapshot string
|
||||
}
|
||||
|
||||
func Open(dir string) (*Store, error) {
|
||||
if err := os.MkdirAll(dir, 0755); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s := &Store{path: filepath.Join(dir, "events.jsonl"), cas: filepath.Join(dir, "cas"), snapshot: filepath.Join(dir, "snapshot.json"), tasks: map[string]domain.Task{}, external: map[string]string{}}
|
||||
if err := os.MkdirAll(s.cas, 0755); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
f, err := os.Open(s.path)
|
||||
if os.IsNotExist(err) {
|
||||
return s, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer f.Close()
|
||||
sc := bufio.NewScanner(f)
|
||||
for sc.Scan() {
|
||||
var e domain.Event
|
||||
if err := json.Unmarshal(sc.Bytes(), &e); err == nil {
|
||||
s.events = append(s.events, e)
|
||||
if err := s.apply(e); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
} else {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return s, sc.Err()
|
||||
}
|
||||
func (s *Store) apply(e domain.Event) error {
|
||||
var p map[string]any
|
||||
if err := json.Unmarshal(e.Payload, &p); err != nil {
|
||||
return err
|
||||
}
|
||||
t := s.tasks[e.TaskID]
|
||||
switch e.Type {
|
||||
case "TaskCreated":
|
||||
if err := domain.ValidateCreated(p); err != nil {
|
||||
return err
|
||||
}
|
||||
t = domain.Task{ID: e.TaskID, Source: p["source"].(string), ExternalID: p["external_id"].(string), Project: p["project"].(string), State: domain.StateQueued}
|
||||
if v, ok := p["capability"].([]any); ok {
|
||||
for _, x := range v {
|
||||
if z, ok := x.(string); ok {
|
||||
t.Capability = append(t.Capability, z)
|
||||
}
|
||||
}
|
||||
}
|
||||
if v, ok := p["title"].(string); ok {
|
||||
t.Title = v
|
||||
}
|
||||
s.external[t.Source+"\x00"+t.ExternalID] = t.ID
|
||||
case "TaskLeased":
|
||||
t.State = domain.StateLeased
|
||||
t.Lease = &domain.Lease{HarnessID: p["harness_id"].(string), Until: time.Unix(0, int64(p["until_ns"].(float64)))}
|
||||
case "TaskReleased":
|
||||
t.State = domain.StateQueued
|
||||
t.Lease = nil
|
||||
case "TaskCompleted":
|
||||
t.State = domain.StateCompleted
|
||||
t.Lease = nil
|
||||
case "TaskFailed":
|
||||
t.State = domain.StateFailed
|
||||
t.Lease = nil
|
||||
case "TaskBlocked":
|
||||
t.State = domain.StateBlocked
|
||||
t.Lease = nil
|
||||
case "TaskAmended":
|
||||
for k, v := range p {
|
||||
if k == "title" {
|
||||
t.Title, v = v.(string)
|
||||
}
|
||||
}
|
||||
}
|
||||
t.Version = e.Version
|
||||
s.tasks[e.TaskID] = t
|
||||
return nil
|
||||
}
|
||||
func (s *Store) Append(e domain.Event) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if err := domain.ValidateEvent(e); err != nil {
|
||||
return err
|
||||
}
|
||||
if e.At.IsZero() {
|
||||
e.At = time.Now().UTC()
|
||||
}
|
||||
if e.Seq == 0 {
|
||||
e.Seq = uint64(len(s.events) + 1)
|
||||
}
|
||||
if e.Type == "TaskCreated" {
|
||||
var p map[string]any
|
||||
if err := json.Unmarshal(e.Payload, &p); err != nil {
|
||||
return err
|
||||
}
|
||||
if id := s.external[p["source"].(string)+"\x00"+p["external_id"].(string)]; id != "" {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
if t, ok := s.tasks[e.TaskID]; ok && e.Version != t.Version+1 {
|
||||
return domain.ErrConflict
|
||||
}
|
||||
if _, ok := s.tasks[e.TaskID]; !ok && e.Type != "TaskCreated" {
|
||||
return domain.ErrNotFound
|
||||
}
|
||||
if e.Type != "TaskCreated" && (e.Type == "TaskCompleted" || e.Type == "TaskBlocked" || e.Type == "TaskReleased") {
|
||||
var p map[string]any
|
||||
_ = json.Unmarshal(e.Payload, &p)
|
||||
for _, k := range []string{"handoff_ref", "report_ref"} {
|
||||
if ref, ok := p[k].(string); ok {
|
||||
if _, err := os.Stat(filepath.Join(s.cas, ref)); err != nil {
|
||||
return fmt.Errorf("%w: missing artifact %s", domain.ErrInvalid, ref)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if err := s.apply(e); err != nil {
|
||||
return err
|
||||
}
|
||||
f, err := os.OpenFile(s.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer f.Close()
|
||||
b, _ := json.Marshal(e)
|
||||
if _, err = f.Write(append(b, '\n')); err != nil {
|
||||
return err
|
||||
}
|
||||
if err = f.Sync(); err != nil {
|
||||
return err
|
||||
}
|
||||
s.events = append(s.events, e)
|
||||
if err := s.writeSnapshot(); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (s *Store) writeSnapshot() error {
|
||||
tasks := make([]domain.Task, 0, len(s.tasks))
|
||||
for _, t := range s.tasks {
|
||||
tasks = append(tasks, t)
|
||||
}
|
||||
b, err := json.Marshal(struct {
|
||||
Seq uint64 `json:"seq"`
|
||||
Tasks []domain.Task `json:"tasks"`
|
||||
}{uint64(len(s.events)), tasks})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
tmp := s.snapshot + ".tmp"
|
||||
if err = os.WriteFile(tmp, b, 0644); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.Rename(tmp, s.snapshot)
|
||||
}
|
||||
func (s *Store) Tasks() []domain.Task {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
out := make([]domain.Task, 0, len(s.tasks))
|
||||
for _, t := range s.tasks {
|
||||
out = append(out, t)
|
||||
}
|
||||
return out
|
||||
}
|
||||
func (s *Store) Events(since uint64) []domain.Event {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
var out []domain.Event
|
||||
for _, e := range s.events {
|
||||
if e.Seq > since {
|
||||
out = append(out, e)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
func (s *Store) PutArtifact(b []byte) (string, error) {
|
||||
h := domain.Hash(b)
|
||||
p := filepath.Join(s.cas, h)
|
||||
if _, err := os.Stat(p); errors.Is(err, os.ErrNotExist) {
|
||||
if err = os.WriteFile(p, b, 0644); err != nil {
|
||||
return "", err
|
||||
}
|
||||
}
|
||||
return h, nil
|
||||
}
|
||||
|
||||
func (s *Store) Task(id string) (domain.Task, bool) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
t, ok := s.tasks[id]
|
||||
return t, ok
|
||||
}
|
||||
|
||||
func (s *Store) Lease(id, harness string, ttl time.Duration) (domain.Event, error) {
|
||||
t, ok := s.Task(id)
|
||||
if !ok {
|
||||
return domain.Event{}, domain.ErrNotFound
|
||||
}
|
||||
if t.State != domain.StateQueued {
|
||||
return domain.Event{}, domain.ErrConflict
|
||||
}
|
||||
p, _ := json.Marshal(map[string]any{"harness_id": harness, "until_ns": time.Now().Add(ttl).UnixNano()})
|
||||
e := domain.Event{ID: id, Type: "TaskLeased", TaskID: id, Version: t.Version + 1, Payload: p}
|
||||
return e, s.Append(e)
|
||||
}
|
||||
|
||||
func (s *Store) ExpireLeases(now time.Time) ([]domain.Event, error) {
|
||||
var out []domain.Event
|
||||
for _, t := range s.Tasks() {
|
||||
if t.State == domain.StateLeased && t.Lease != nil && !t.Lease.Until.After(now) {
|
||||
p, _ := json.Marshal(map[string]any{"reason": "lease_expired", "harness_id": t.Lease.HarnessID})
|
||||
e := domain.Event{ID: t.ID, Type: "TaskReleased", TaskID: t.ID, Version: t.Version + 1, Payload: p}
|
||||
if err := s.Append(e); err != nil {
|
||||
return out, err
|
||||
}
|
||||
out = append(out, e)
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"orchestra/internal/domain"
|
||||
)
|
||||
|
||||
func created(id string) domain.Event {
|
||||
b, _ := json.Marshal(map[string]any{"source": "jsonl", "external_id": "42", "project": "demo", "capability": []string{"mechanical"}})
|
||||
return domain.Event{ID: id, Type: "TaskCreated", TaskID: "task-1", Version: 1, Payload: b}
|
||||
}
|
||||
|
||||
func TestAppendReplayAndDeduplicate(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
s, err := Open(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.Append(created("e1")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.Append(created("e2")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := len(s.Events(0)); got != 1 {
|
||||
t.Fatalf("duplicate ingest appended %d events", got)
|
||||
}
|
||||
ref, err := s.PutArtifact([]byte("report"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
completion, _ := json.Marshal(map[string]string{"report_ref": ref})
|
||||
if err := s.Append(domain.Event{Type: "TaskCompleted", TaskID: "task-1", Version: 2, Payload: completion}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.Append(domain.Event{Type: "TaskReleased", TaskID: "task-1", Version: 2, Payload: json.RawMessage(`{"handoff_ref":"x"}`)}); err != domain.ErrConflict {
|
||||
t.Fatalf("expected conflict, got %v", err)
|
||||
}
|
||||
s2, err := Open(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := s2.Tasks()[0].State; got != domain.StateCompleted {
|
||||
t.Fatalf("replay state = %s", got)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(dir, "snapshot.json")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestArtifactIsContentAddressed(t *testing.T) {
|
||||
s, err := Open(t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
h1, err := s.PutArtifact([]byte("proof"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
h2, err := s.PutArtifact([]byte("proof"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if h1 != h2 {
|
||||
t.Fatal("same artifact received different hashes")
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(s.cas, h1)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user