diff --git a/README.md b/README.md index f780bee9d..75d9a94a1 100644 --- a/README.md +++ b/README.md @@ -489,6 +489,8 @@ pad item deps Show dependencies pad item unblock Remove dependency pad item related Show direct relationships for an item pad item implemented-by Show incoming implementers for an item +pad item claim Atomically claim an item for execution (--holder, --ttl; 409 names the live holder) +pad item release Release your execution lease (idempotent) pad item bulk-update --status X Batch update multiple items pad collection list List collections with item counts diff --git a/cmd/pad/cmd_item_lease.go b/cmd/pad/cmd_item_lease.go new file mode 100644 index 000000000..f96184156 --- /dev/null +++ b/cmd/pad/cmd_item_lease.go @@ -0,0 +1,116 @@ +package main + +import ( + "fmt" + "time" + + "github.com/fatih/color" + "github.com/spf13/cobra" + + "github.com/PerpetualSoftware/pad/internal/cli" +) + +// `pad item claim` / `pad item release` (#1221): the execution lease. A +// claim answers "may I be the one executing this right now" atomically — +// two pollers that both read "unclaimed" get one winner and one +// structured refusal naming the live holder, instead of two winners. + +func itemClaimCmd() *cobra.Command { + var ( + holderFlag string + ttlFlag string + ) + + cmd := &cobra.Command{ + Use: "claim ", + Short: "Atomically claim an item for execution (lease with expiry)", + Long: `Acquire the execution lease on an item, or refresh it if you already +hold it (a re-claim by the live holder extends the expiry — heartbeat). + +Fails with the live holder and expiry when someone else holds the item, so +the loser can log who won and skip instead of double-working. The lease +expires on its own; a crashed holder blocks nobody past the TTL. + +Examples: + pad item claim TASK-5 + pad item claim TASK-5 --holder sweep-runner --ttl 30m`, + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + ttlSeconds := 0 + if ttlFlag != "" { + d, err := time.ParseDuration(ttlFlag) + if err != nil { + return fmt.Errorf("invalid --ttl %q: %w (use Go durations, e.g. 15m, 1h)", ttlFlag, err) + } + if d <= 0 { + return fmt.Errorf("--ttl must be positive, got %s", d) + } + ttlSeconds = int(d.Seconds()) + } + + client, _ := getClient() + ws := getWorkspace() + + lease, err := client.ClaimItem(ws, args[0], holderFlag, ttlSeconds) + if err != nil { + return err + } + + if formatFlag == "json" { + return cli.PrintJSON(lease) + } + + green := color.New(color.FgGreen) + fmt.Printf("%s Lease acquired on %s\n", green.Sprint("✓"), args[0]) + fmt.Printf(" Holder: %s\n", lease.Holder) + fmt.Printf(" Expires: %s (%s)\n", + lease.ExpiresAt.Local().Format("2006-01-02 15:04:05"), + cli.LeaseCountdown(lease.ExpiresAt)) + return nil + }, + } + + cmd.Flags().StringVar(&holderFlag, "holder", "", "lease holder identity (default: the authenticated user)") + cmd.Flags().StringVar(&ttlFlag, "ttl", "", "lease duration, e.g. 15m, 1h (default: 15m; max 24h)") + + return cmd +} + +func itemReleaseCmd() *cobra.Command { + var holderFlag string + + cmd := &cobra.Command{ + Use: "release ", + Short: "Release an item's execution lease (idempotent)", + Long: `Release the execution lease you hold on an item. Releasing an absent or +already-expired lease is a no-op, not an error — cleanup code never needs +to check whether it still holds the lease first. Releasing another +holder's LIVE lease is refused.`, + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + client, _ := getClient() + ws := getWorkspace() + + released, err := client.ReleaseItem(ws, args[0], holderFlag) + if err != nil { + return err + } + + if formatFlag == "json" { + return cli.PrintJSON(map[string]any{"ref": args[0], "released": released}) + } + + green := color.New(color.FgGreen) + if released { + fmt.Printf("%s Lease released on %s\n", green.Sprint("✓"), args[0]) + } else { + fmt.Printf("%s No live lease to release on %s (already expired or never held)\n", green.Sprint("✓"), args[0]) + } + return nil + }, + } + + cmd.Flags().StringVar(&holderFlag, "holder", "", "lease holder identity (default: the authenticated user)") + + return cmd +} diff --git a/cmd/pad/cmd_item_lease_test.go b/cmd/pad/cmd_item_lease_test.go new file mode 100644 index 000000000..230e164c5 --- /dev/null +++ b/cmd/pad/cmd_item_lease_test.go @@ -0,0 +1,199 @@ +package main + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" +) + +// Tests for `pad item claim` / `pad item release` (#1221) — the CLI face +// of the execution lease. The stub answers the two endpoints; the tests +// drive the real commands and pin what the automation-facing output and +// exit behaviour look like, because sweep scripts branch on them. + +type stubLeaseServer struct { + *httptest.Server + claimBodies []map[string]any + releaseBodies []map[string]any + holdConflict bool // when true, claim answers 409 lease_held +} + +func newStubLeaseServer(t *testing.T) *stubLeaseServer { + t.Helper() + s := &stubLeaseServer{} + mux := http.NewServeMux() + mux.HandleFunc("/api/v1/workspaces/ws1/items/TASK-5/claim", func(w http.ResponseWriter, r *http.Request) { + var body map[string]any + _ = json.NewDecoder(r.Body).Decode(&body) + s.claimBodies = append(s.claimBodies, body) + if s.holdConflict { + w.WriteHeader(http.StatusConflict) + _ = json.NewEncoder(w).Encode(map[string]any{ + "error": map[string]any{ + "code": "lease_held", + "message": `item is leased to "other-runner" until 2026-09-02T22:00:00Z`, + "details": map[string]any{ + "ref": "TASK-5", "holder": "other-runner", + "acquired_at": "2026-09-02T21:45:00Z", + "expires_at": "2026-09-02T22:00:00Z", + }, + }, + }) + return + } + _ = json.NewEncoder(w).Encode(map[string]any{ + "ref": "TASK-5", + "lease": map[string]any{ + "holder": "sweep-runner", + "acquired_at": "2026-09-02T21:30:00Z", + "expires_at": "2026-09-02T21:45:00Z", + }, + }) + }) + mux.HandleFunc("/api/v1/workspaces/ws1/items/TASK-5/release", func(w http.ResponseWriter, r *http.Request) { + var body map[string]any + _ = json.NewDecoder(r.Body).Decode(&body) + s.releaseBodies = append(s.releaseBodies, body) + released := len(s.claimBodies) > 0 // released=true only if something was claimed + _ = json.NewEncoder(w).Encode(map[string]any{"ref": "TASK-5", "released": released}) + }) + s.Server = httptest.NewServer(mux) + t.Cleanup(s.Close) + return s +} + +func setupLeaseCLI(t *testing.T) *stubLeaseServer { + t.Helper() + setTempHomeMain(t) + srv := newStubLeaseServer(t) + t.Setenv("PAD_URL", srv.URL) + t.Setenv("PAD_TOKEN", "pad_envtoken") + // Workspace detection walks up from CWD for a .pad.toml — + // setTempHomeMain already chdir'd to a fresh temp dir, so a minimal + // pin there routes getWorkspace() to the stub's workspace. + cwd, err := os.Getwd() + if err != nil { + t.Fatalf("getwd: %v", err) + } + if err := os.WriteFile(filepath.Join(cwd, ".pad.toml"), []byte("workspace = \"ws1\"\n"), 0644); err != nil { + t.Fatalf("write .pad.toml: %v", err) + } + return srv +} + +// claim sends holder + ttl (converted to ttl_seconds) and confirms the +// acquired lease with holder and expiry — the sweep's success signal. +func TestItemClaim_SendsTTLAndPrintsLease(t *testing.T) { + srv := setupLeaseCLI(t) + + cmd := itemClaimCmd() + cmd.SetArgs([]string{"TASK-5", "--holder", "sweep-runner", "--ttl", "10m"}) + var runErr error + out := captureStdout(t, func() { + runErr = cmd.Execute() + }) + if runErr != nil { + t.Fatalf("claim: %v", runErr) + } + if len(srv.claimBodies) != 1 { + t.Fatalf("expected one claim POST, got %d", len(srv.claimBodies)) + } + body := srv.claimBodies[0] + if body["holder"] != "sweep-runner" { + t.Errorf("posted holder = %v, want sweep-runner", body["holder"]) + } + if n, _ := body["ttl_seconds"].(float64); int(n) != 600 { + t.Errorf("posted ttl_seconds = %v, want 600", body["ttl_seconds"]) + } + if !strings.Contains(out, "sweep-runner") || !strings.Contains(strings.ToLower(out), "lease") { + t.Errorf("output should confirm the lease and holder:\n%s", out) + } +} + +// A contended claim exits non-zero with the holder and expiry in the +// error — the loser must be able to log WHO holds the item without a +// second call. +func TestItemClaim_ContendedExitsWithHolder(t *testing.T) { + srv := setupLeaseCLI(t) + srv.holdConflict = true + + cmd := itemClaimCmd() + cmd.SetArgs([]string{"TASK-5"}) + cmd.SilenceUsage = true + cmd.SilenceErrors = true + var runErr error + captureStdout(t, func() { + runErr = cmd.Execute() + }) + if runErr == nil { + t.Fatal("contended claim must exit with an error") + } + if !strings.Contains(runErr.Error(), "other-runner") { + t.Errorf("error should name the live holder, got %q", runErr.Error()) + } +} + +// An unparseable --ttl fails locally; the server is never asked. +func TestItemClaim_BadTTLFailsLocally(t *testing.T) { + srv := setupLeaseCLI(t) + + cmd := itemClaimCmd() + cmd.SetArgs([]string{"TASK-5", "--ttl", "banana"}) + cmd.SilenceUsage = true + cmd.SilenceErrors = true + var runErr error + captureStdout(t, func() { + runErr = cmd.Execute() + }) + if runErr == nil { + t.Fatal("expected a parse error for --ttl banana") + } + if len(srv.claimBodies) != 0 { + t.Errorf("no claim should reach the server on a bad ttl, got %d", len(srv.claimBodies)) + } +} + +// release distinguishes a real release from the idempotent no-op — both +// succeed (exit 0), but the words differ so a human reading sweep logs +// can tell them apart. +func TestItemRelease_ReportsReleasedVsNoop(t *testing.T) { + srv := setupLeaseCLI(t) + + // No claim yet: the stub answers released=false. + cmd := itemReleaseCmd() + cmd.SetArgs([]string{"TASK-5"}) + var runErr error + out := captureStdout(t, func() { + runErr = cmd.Execute() + }) + if runErr != nil { + t.Fatalf("no-op release must still exit 0: %v", runErr) + } + if !strings.Contains(strings.ToLower(out), "no live lease") { + t.Errorf("no-op release should say nothing was held:\n%s", out) + } + + // After a claim: released=true. + claim := itemClaimCmd() + claim.SetArgs([]string{"TASK-5"}) + captureStdout(t, func() { _ = claim.Execute() }) + + cmd = itemReleaseCmd() + cmd.SetArgs([]string{"TASK-5", "--holder", "sweep-runner"}) + out = captureStdout(t, func() { + runErr = cmd.Execute() + }) + if runErr != nil { + t.Fatalf("release: %v", runErr) + } + if len(srv.releaseBodies) != 2 || srv.releaseBodies[1]["holder"] != "sweep-runner" { + t.Errorf("release bodies = %v, want second with holder sweep-runner", srv.releaseBodies) + } + if !strings.Contains(strings.ToLower(out), "released") { + t.Errorf("real release should say released:\n%s", out) + } +} diff --git a/cmd/pad/groups.go b/cmd/pad/groups.go index 9dd9907df..73f555214 100644 --- a/cmd/pad/groups.go +++ b/cmd/pad/groups.go @@ -126,6 +126,8 @@ func itemCmd() *cobra.Command { commentsCmd(), noteCmd(), decideCmd(), + itemClaimCmd(), + itemReleaseCmd(), blocksCmd(), blockedByCmd(), depsCmd(), diff --git a/internal/cli/client_lease.go b/internal/cli/client_lease.go new file mode 100644 index 000000000..a391c0f00 --- /dev/null +++ b/internal/cli/client_lease.go @@ -0,0 +1,53 @@ +package cli + +import ( + "net/url" + + "github.com/PerpetualSoftware/pad/internal/models" +) + +// Item execution lease endpoints (#1221): +// POST /workspaces/{ws}/items/{ref}/claim and .../release. A 409 with +// code "lease_held" comes back as *APIError whose message names the live +// holder and expiry; callers surface it rather than retrying blindly. + +// leaseRequest is the wire body for claim and release. Zero values are +// omitted so the server applies its defaults (authenticated identity as +// holder; 15-minute TTL). +type leaseRequest struct { + Holder string `json:"holder,omitempty"` + TTLSeconds int `json:"ttl_seconds,omitempty"` +} + +// claimResponse mirrors the structured success bodies: ref plus the lease +// (claim) or the released verdict (release). +type leaseResponse struct { + Ref string `json:"ref"` + Lease *models.ItemLease `json:"lease,omitempty"` + Released *bool `json:"released,omitempty"` +} + +// ClaimItem acquires (or refreshes, for the live holder) the execution +// lease on an item. holder and ttlSeconds may be zero for the server +// defaults. +func (c *Client) ClaimItem(wsSlug, ref, holder string, ttlSeconds int) (*models.ItemLease, error) { + var out leaseResponse + err := c.post("/workspaces/"+url.PathEscape(wsSlug)+"/items/"+url.PathEscape(ref)+"/claim", + leaseRequest{Holder: holder, TTLSeconds: ttlSeconds}, &out) + if err != nil { + return nil, err + } + return out.Lease, nil +} + +// ReleaseItem clears the holder's lease on an item. released=false means +// there was nothing live to release — a success, not an error. +func (c *Client) ReleaseItem(wsSlug, ref, holder string) (bool, error) { + var out leaseResponse + err := c.post("/workspaces/"+url.PathEscape(wsSlug)+"/items/"+url.PathEscape(ref)+"/release", + leaseRequest{Holder: holder}, &out) + if err != nil { + return false, err + } + return out.Released != nil && *out.Released, nil +} diff --git a/internal/cli/format.go b/internal/cli/format.go index c34d13d9b..4640ff5fb 100644 --- a/internal/cli/format.go +++ b/internal/cli/format.go @@ -105,6 +105,21 @@ func StatusIcon(status string) string { } // RelativeTime returns a human-readable relative time string. +// LeaseCountdown renders a lease expiry as a short human countdown +// ("expires in 12m"). Sub-minute leases show seconds; anything longer +// rounds to whole minutes — a lease TTL is coarse coordination state, so +// second-level precision past the first minute is noise. +func LeaseCountdown(expires time.Time) string { + d := time.Until(expires) + if d <= 0 { + return "expired" + } + if d < time.Minute { + return fmt.Sprintf("expires in %ds", int(d.Seconds())) + } + return fmt.Sprintf("expires in %dm", int(d.Round(time.Minute).Minutes())) +} + func RelativeTime(t time.Time) string { d := time.Since(t) switch { @@ -252,6 +267,7 @@ func renderItemTable(w io.Writer, items []models.Item, maxWidth int) { // title budget), so we stash its raw pieces. type row struct { pin string // rendered pin marker ("" or "* " in yellow) + lease string // rendered lease marker ("" or "» " in cyan) — live execution lease (#1221) ref string // raw ref, e.g. "TASK-5" — never truncated title string // raw title — the only thing we truncate archived bool @@ -267,6 +283,13 @@ func renderItemTable(w io.Writer, items []models.Item, maxWidth int) { if item.Pinned { r.pin = color.YellowString("* ") } + // A leased row gets a glyph on the ref/title cell rather than its + // own column: leases are absent on most rows, and a permanent + // column would tax every row's title budget for the empty case. + // Holder + expiry detail lives on `item show` and in JSON. + if item.Lease != nil { + r.lease = color.CyanString("» ") + } r.ref = ItemRef(item) // Strip any SGR escapes a title might carry (we apply our own Bold // styling). This keeps displayWidth == rune count for the title, so @@ -347,7 +370,7 @@ func renderItemTable(w io.Writer, items []models.Item, maxWidth int) { for _, r := range rows { writeRow( - buildRefCell(r.pin, r.ref, r.title, r.archived, titleBudget), + buildRefCell(r.pin+r.lease, r.ref, r.title, r.archived, titleBudget), r.status, r.priority, r.collection, r.updated, ) } @@ -508,6 +531,14 @@ func PrintItemMeta(item *models.Item) { fmt.Printf("%s %s\n", label.Sprint("Assigned: "), assignStr) } + // Live execution lease (#1221) — who is actively executing this item + // right now, distinct from Assigned's longer-term ownership. The + // server omits expired leases, so presence here means live. + if item.Lease != nil { + fmt.Printf("%s %s\n", label.Sprint("Lease: "), + fmt.Sprintf("%s (%s)", item.Lease.Holder, LeaseCountdown(item.Lease.ExpiresAt))) + } + tags := item.Tags if tags == "[]" || tags == "" || tags == "null" { tags = Dim.Sprint("(none)") diff --git a/internal/cli/format_lease_test.go b/internal/cli/format_lease_test.go new file mode 100644 index 000000000..eaa2eb8fd --- /dev/null +++ b/internal/cli/format_lease_test.go @@ -0,0 +1,110 @@ +package cli + +import ( + "io" + "os" + "strings" + "testing" + "time" + + "github.com/PerpetualSoftware/pad/internal/models" +) + +// Display tests for the execution lease (#1221): `item show` renders a +// Lease: line only when a live lease is present, and the list table marks +// leased rows with a glyph on the ref/title cell instead of a column. + +// captureMetaStdout captures os.Stdout around fn — PrintItemMeta writes +// via fmt.Printf directly. +func captureMetaStdout(t *testing.T, fn func()) string { + t.Helper() + old := os.Stdout + r, w, err := os.Pipe() + if err != nil { + t.Fatalf("pipe: %v", err) + } + os.Stdout = w + defer func() { os.Stdout = old }() + fn() + _ = w.Close() + data, err := io.ReadAll(r) + if err != nil { + t.Fatalf("read pipe: %v", err) + } + return string(data) +} + +func leaseDisplayItem(withLease bool) *models.Item { + item := &models.Item{ + Title: "Leased item", + Fields: `{"status":"open"}`, + Tags: "[]", + } + if withLease { + item.Lease = &models.ItemLease{ + Holder: "sweep-runner", + AcquiredAt: time.Now().UTC().Add(-time.Minute), + ExpiresAt: time.Now().UTC().Add(12 * time.Minute), + } + } + return item +} + +func TestPrintItemMeta_LeaseLineOnlyWhenLive(t *testing.T) { + out := captureMetaStdout(t, func() { PrintItemMeta(leaseDisplayItem(true)) }) + if !strings.Contains(out, "Lease:") { + t.Errorf("live lease missing from meta header:\n%s", out) + } + if !strings.Contains(out, "sweep-runner") { + t.Errorf("lease line must name the holder:\n%s", out) + } + if !strings.Contains(out, "expires in") { + t.Errorf("lease line must carry the countdown:\n%s", out) + } + + out = captureMetaStdout(t, func() { PrintItemMeta(leaseDisplayItem(false)) }) + if strings.Contains(out, "Lease:") { + t.Errorf("unleased item must not render a Lease line:\n%s", out) + } +} + +func TestRenderItemTable_MarksLeasedRows(t *testing.T) { + leased := *leaseDisplayItem(true) + plain := *leaseDisplayItem(false) + plain.Title = "Plain item" + + var sb strings.Builder + renderItemTable(&sb, []models.Item{leased, plain}, 120) + out := sb.String() + + leasedLine, plainLine := "", "" + for _, line := range strings.Split(out, "\n") { + if strings.Contains(line, "Leased item") { + leasedLine = line + } + if strings.Contains(line, "Plain item") { + plainLine = line + } + } + if leasedLine == "" || plainLine == "" { + t.Fatalf("both rows must render:\n%s", out) + } + if !strings.Contains(leasedLine, "»") { + t.Errorf("leased row missing the lease glyph:\n%s", leasedLine) + } + if strings.Contains(plainLine, "»") { + t.Errorf("unleased row must not carry the glyph:\n%s", plainLine) + } +} + +func TestLeaseCountdown(t *testing.T) { + if got := LeaseCountdown(time.Now().Add(-time.Second)); got != "expired" { + t.Errorf("past expiry = %q, want expired", got) + } + if got := LeaseCountdown(time.Now().Add(30 * time.Second)); !strings.HasPrefix(got, "expires in ") || !strings.HasSuffix(got, "s") { + t.Errorf("sub-minute countdown = %q, want seconds form", got) + } + if got := LeaseCountdown(time.Now().Add(12 * time.Minute)); got != "expires in 12m" { + t.Errorf("countdown = %q, want expires in 12m", got) + } +} diff --git a/internal/models/item.go b/internal/models/item.go index 2eeaec164..3bdf6232c 100644 --- a/internal/models/item.go +++ b/internal/models/item.go @@ -186,6 +186,16 @@ type Item struct { // or share-link path, and never emit a placeholder when the gate denies. MovedTo []ItemMovedTo `json:"moved_to,omitempty"` + // Lease is the item's LIVE execution lease (#1221) — who is actively + // executing this item right now, distinct from AssignedUserID's + // longer-term ownership. Populated on the single-item GET and on list + // responses; nil when unclaimed OR expired (an expired lease is absent + // on every read path — expiry is the reaper). Never persisted through + // item update paths: claiming/releasing goes through the dedicated + // claim/release endpoints and deliberately bumps neither updated_at + // nor the version history. + Lease *ItemLease `json:"lease,omitempty"` + DerivedClosure *ItemDerivedClosure `json:"derived_closure,omitempty"` CodeContext *ItemCodeContext `json:"code_context,omitempty"` Convention *ItemConventionMetadata `json:"convention,omitempty"` @@ -253,6 +263,16 @@ func (item *Item) ComputeRef() { } } +// ItemLease is a live execution lease on an item (#1221): holder is a +// freeform identity string (defaulting server-side to the authenticated +// user), and the lease is live iff ExpiresAt is in the future. Expired +// leases are never returned — expiry is absence. +type ItemLease struct { + Holder string `json:"holder"` + AcquiredAt time.Time `json:"acquired_at"` + ExpiresAt time.Time `json:"expires_at"` +} + // ItemMovedTo is one destination an archived item was moved to, rendered in // DISPLAYABLE terms. It deliberately carries no UUIDs: the consumer must be // able to render (and link to) the destination without a second call, and diff --git a/internal/server/decode_json_nul_test.go b/internal/server/decode_json_nul_test.go index 70d54d76c..7f96ffe63 100644 --- a/internal/server/decode_json_nul_test.go +++ b/internal/server/decode_json_nul_test.go @@ -268,6 +268,7 @@ func TestEveryRequestBodyReaderIsAccountedFor(t *testing.T) { "handlers_tokens.go": "guards on r.Body != nil && r.ContentLength != 0, then decodes THROUGH decodeJSON — so the body is read by the chokepoint, which applies the cap and the NUL rule. The earlier reason here said it never reads the body, which was simply false (codex round 29): a wrong reason in this list is the same defect as a missing entry, since both let a reader pass as reviewed", "handlers_oauth.go": "KNOWN GAP, tracked as BUG-2811: the OAuth handlers read FORM-encoded bodies (r.Form/FormValue), which no rule in this family covers — the transport rules see the query half of r.Form and not the body half. Listed so this test states the gap instead of being blind to it; measuring it needs a fosite-backed fixture.", "handlers_watches.go": "guards on r.Body != nil && r.ContentLength != 0, then decodes THROUGH decodeJSON — the closing-round-4 fix for the chunked-body drop; the one reader expression is the nil check itself, and the body bytes flow through the chokepoint", + "handlers_item_lease.go": "guards on r.Body != nil && r.ContentLength != 0, then decodes THROUGH decodeJSON — same shape as handlers_watches.go; the one reader expression is the nil check, and the body (optional holder/ttl_seconds) flows through the chokepoint, so the caller-text holder gets BUG-2803's NUL refusal before it can reach items.lease_holder", } // Reader counts as reviewed. A mismatch means this file gained or lost a @@ -287,6 +288,7 @@ func TestEveryRequestBodyReaderIsAccountedFor(t *testing.T) { "handlers_tokens.go": 1, "handlers_oauth.go": 13, "handlers_watches.go": 1, + "handlers_item_lease.go": 1, } accounted := map[string]entry{} for name, why := range accountedWhy { diff --git a/internal/server/handlers_item_lease.go b/internal/server/handlers_item_lease.go new file mode 100644 index 000000000..8d0a63a98 --- /dev/null +++ b/internal/server/handlers_item_lease.go @@ -0,0 +1,177 @@ +package server + +import ( + "errors" + "net/http" + "time" + + "github.com/go-chi/chi/v5" + + "github.com/PerpetualSoftware/pad/internal/store" +) + +// Item execution lease endpoints (#1221): +// +// POST /api/v1/workspaces/{slug}/items/{itemSlug}/claim +// POST /api/v1/workspaces/{slug}/items/{itemSlug}/release +// +// A lease answers "who is actively executing this right now" — distinct +// from assignment's longer-term ownership. The store's conditional UPDATE +// is the arbiter (see store/item_lease.go), so two callers racing on the +// same "unclaimed" snapshot produce one winner and one structured 409. +// +// Deliberately NOT wired: an SSE/activity event, an updated_at bump, or a +// version entry. A lease is runtime coordination state, not content — a +// claim must never 409 a concurrent editor's expected_updated_at token, +// and a heartbeat re-claim every few minutes must not spam the feed. + +const ( + defaultLeaseTTL = 15 * time.Minute + maxLeaseTTL = 24 * time.Hour +) + +// itemLeaseInput is the request body for claim and release. Everything is +// optional: holder defaults to the authenticated identity, ttl_seconds to +// defaultLeaseTTL (claim only). +type itemLeaseInput struct { + Holder string `json:"holder,omitempty"` + TTLSeconds int `json:"ttl_seconds,omitempty"` +} + +// handleClaimItem acquires (or, for the live holder, refreshes) the +// execution lease on an item. +func (s *Server) handleClaimItem(w http.ResponseWriter, r *http.Request) { + item, holder, input, ok := s.resolveLeaseRequest(w, r) + if !ok { + return + } + + ttl := defaultLeaseTTL + if input.TTLSeconds != 0 { + if input.TTLSeconds < 0 || time.Duration(input.TTLSeconds)*time.Second > maxLeaseTTL { + writeError(w, http.StatusBadRequest, "bad_request", + "ttl_seconds must be between 1 and 86400 (omit it for the 15-minute default)") + return + } + ttl = time.Duration(input.TTLSeconds) * time.Second + } + + lease, err := s.store.ClaimItemLease(item.ID, holder, ttl) + if err != nil { + var held *store.LeaseHeldError + if errors.As(err, &held) { + writeLeaseHeldError(w, item.Ref, held) + return + } + writeInternalError(w, err) + return + } + + // Structured success body (the BUG-1081 pattern): enough to be the + // caller's next source of truth without a follow-up read. + writeJSON(w, http.StatusOK, map[string]any{ + "ref": item.Ref, + "lease": lease, + }) +} + +// handleReleaseItem clears the caller's lease. Idempotent: releasing an +// absent or expired lease answers released=false, never an error. +func (s *Server) handleReleaseItem(w http.ResponseWriter, r *http.Request) { + item, holder, _, ok := s.resolveLeaseRequest(w, r) + if !ok { + return + } + + released, err := s.store.ReleaseItemLease(item.ID, holder) + if err != nil { + var held *store.LeaseHeldError + if errors.As(err, &held) { + writeLeaseHeldError(w, item.Ref, held) + return + } + writeInternalError(w, err) + return + } + + writeJSON(w, http.StatusOK, map[string]any{ + "ref": item.Ref, + "released": released, + }) +} + +// resolveLeaseRequest performs the shared claim/release preamble: resolve +// workspace + item, check visibility, require an authenticated user, and +// decode the optional body. holder falls back to the authenticated +// user's email (their durable, human-readable identity; the #879 named +// profiles become the natural source once layer 2 lands) and then to the +// user id when the account has no email. +func (s *Server) resolveLeaseRequest(w http.ResponseWriter, r *http.Request) (item *resolvedLeaseItem, holder string, input itemLeaseInput, ok bool) { + workspaceID, wok := s.getWorkspaceID(w, r) + if !wok { + return nil, "", input, false + } + + itemSlug := chi.URLParam(r, "itemSlug") + resolved, err := s.store.ResolveItem(workspaceID, itemSlug) + if err != nil { + writeInternalError(w, err) + return nil, "", input, false + } + if resolved == nil { + s.writeItemResolveError(w, r, workspaceID, itemSlug) + return nil, "", input, false + } + if !s.requireItemVisible(w, r, workspaceID, resolved) { + return nil, "", input, false + } + + user := currentUser(r) + if user == nil { + writeError(w, http.StatusUnauthorized, "unauthorized", "Authentication required") + return nil, "", input, false + } + + if r.Body != nil && r.ContentLength != 0 { + if err := decodeJSON(r, &input); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid JSON body") + return nil, "", input, false + } + } + + holder = input.Holder + if holder == "" { + holder = user.Email + } + if holder == "" { + holder = user.ID + } + + return &resolvedLeaseItem{ID: resolved.ID, Ref: resolved.Ref}, holder, input, true +} + +// resolvedLeaseItem is the slice of the resolved item the lease handlers +// need — id for the store call, ref for the response envelope. +type resolvedLeaseItem struct { + ID string + Ref string +} + +// writeLeaseHeldError emits the structured 409 for a live foreign lease — +// the same envelope discipline as update_conflict: a stable `code` plus +// details a client can act on (wait until expires_at, skip, or escalate) +// without parsing the message. +func writeLeaseHeldError(w http.ResponseWriter, ref string, held *store.LeaseHeldError) { + writeJSON(w, http.StatusConflict, map[string]any{ + "error": map[string]any{ + "code": "lease_held", + "message": held.Error(), + "details": map[string]any{ + "ref": ref, + "holder": held.Holder, + "acquired_at": held.AcquiredAt.UTC().Format(time.RFC3339), + "expires_at": held.ExpiresAt.UTC().Format(time.RFC3339), + }, + }, + }) +} diff --git a/internal/server/handlers_item_lease_test.go b/internal/server/handlers_item_lease_test.go new file mode 100644 index 000000000..da01886f4 --- /dev/null +++ b/internal/server/handlers_item_lease_test.go @@ -0,0 +1,322 @@ +package server + +import ( + "bytes" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/PerpetualSoftware/pad/internal/models" +) + +// Tests for the item claim/release endpoints (#1221): +// POST /api/v1/workspaces/{slug}/items/{itemSlug}/claim and /release. +// The claim is atomic (the store's conditional UPDATE is the arbiter); +// contention answers 409 code=lease_held naming the holder and expiry, +// the same structured-envelope discipline as update_conflict. + +type leaseFixtureServer struct { + srv *Server + wsSlug string + item models.Item + tokenA string // user-a@example.com + tokenB string // user-b@example.com +} + +func setupLeaseFixture(t *testing.T) *leaseFixtureServer { + t.Helper() + srv := testServer(t) + slug := createWSWithCollections(t, srv) + + rr := doRequest(srv, "POST", "/api/v1/workspaces/"+slug+"/collections/tasks/items", + map[string]interface{}{"title": "Contended item", "fields": `{"status":"open"}`}) + if rr.Code != http.StatusCreated { + t.Fatalf("seed item: %d %s", rr.Code, rr.Body.String()) + } + var item models.Item + parseJSON(t, rr, &item) + + ws, err := srv.store.GetWorkspaceBySlug(slug) + if err != nil { + t.Fatalf("GetWorkspaceBySlug: %v", err) + } + + mint := func(email string) string { + t.Helper() + user, err := srv.store.CreateUser(models.UserCreate{ + Email: email, Name: email, Password: "pw-test-12345", + }) + if err != nil { + t.Fatalf("CreateUser %s: %v", email, err) + } + if err := srv.store.AddWorkspaceMember(ws.ID, user.ID, "owner"); err != nil { + t.Fatalf("AddWorkspaceMember: %v", err) + } + tok, err := srv.store.CreateAPIToken(user.ID, models.APITokenCreate{ + Name: "lease-test", WorkspaceID: ws.ID, + }, 0, 0) + if err != nil { + t.Fatalf("CreateAPIToken: %v", err) + } + return tok.Token + } + + return &leaseFixtureServer{ + srv: srv, + wsSlug: slug, + item: item, + tokenA: mint("user-a@example.com"), + tokenB: mint("user-b@example.com"), + } +} + +func (f *leaseFixtureServer) call(t *testing.T, token, action string, body map[string]any) *httptest.ResponseRecorder { + t.Helper() + data, err := json.Marshal(body) + if err != nil { + t.Fatalf("marshal body: %v", err) + } + req := httptest.NewRequest("POST", + "/api/v1/workspaces/"+f.wsSlug+"/items/"+f.item.Slug+"/"+action, + bytes.NewReader(data)) + req.Header.Set("Content-Type", "application/json") + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + req.RemoteAddr = "127.0.0.1:0" + rec := httptest.NewRecorder() + f.srv.ServeHTTP(rec, req) + return rec +} + +// Claim → contended claim → release → re-claim: the full lifecycle, with +// the 409 envelope naming the live holder and expiry. +func TestItemLease_ClaimContendReleaseReclaim(t *testing.T) { + f := setupLeaseFixture(t) + + // A claims with an explicit holder string. + rr := f.call(t, f.tokenA, "claim", map[string]any{"holder": "sweep-runner", "ttl_seconds": 600}) + if rr.Code != http.StatusOK { + t.Fatalf("claim: expected 200, got %d (body: %s)", rr.Code, rr.Body.String()) + } + var claimResp struct { + Ref string `json:"ref"` + Lease models.ItemLease `json:"lease"` + } + if err := json.Unmarshal(rr.Body.Bytes(), &claimResp); err != nil { + t.Fatalf("parse claim response: %v", err) + } + if claimResp.Ref != f.item.Ref { + t.Errorf("ref = %q, want %q", claimResp.Ref, f.item.Ref) + } + if claimResp.Lease.Holder != "sweep-runner" { + t.Errorf("holder = %q, want sweep-runner", claimResp.Lease.Holder) + } + if !claimResp.Lease.ExpiresAt.After(time.Now().UTC().Add(9 * time.Minute)) { + t.Errorf("expiry %v not ~10m out", claimResp.Lease.ExpiresAt) + } + + // B's claim while A's is live: 409 lease_held naming holder + expiry. + rr = f.call(t, f.tokenB, "claim", map[string]any{"holder": "other-runner"}) + if rr.Code != http.StatusConflict { + t.Fatalf("contended claim: expected 409, got %d (body: %s)", rr.Code, rr.Body.String()) + } + var conflict struct { + Error struct { + Code string `json:"code"` + Details map[string]any `json:"details"` + } `json:"error"` + } + if err := json.Unmarshal(rr.Body.Bytes(), &conflict); err != nil { + t.Fatalf("parse conflict: %v", err) + } + if conflict.Error.Code != "lease_held" { + t.Errorf("code = %q, want lease_held", conflict.Error.Code) + } + if conflict.Error.Details["holder"] != "sweep-runner" { + t.Errorf("details.holder = %v, want sweep-runner", conflict.Error.Details["holder"]) + } + if conflict.Error.Details["expires_at"] == nil { + t.Error("details.expires_at missing — the caller can't decide to wait or skip without it") + } + + // The holder releases; released=true. + rr = f.call(t, f.tokenA, "release", map[string]any{"holder": "sweep-runner"}) + if rr.Code != http.StatusOK { + t.Fatalf("release: expected 200, got %d (body: %s)", rr.Code, rr.Body.String()) + } + var relResp map[string]any + if err := json.Unmarshal(rr.Body.Bytes(), &relResp); err != nil { + t.Fatalf("parse release response: %v", err) + } + if relResp["released"] != true { + t.Errorf("released = %v, want true", relResp["released"]) + } + + // B claims successfully now. + rr = f.call(t, f.tokenB, "claim", map[string]any{"holder": "other-runner"}) + if rr.Code != http.StatusOK { + t.Fatalf("post-release claim: expected 200, got %d (body: %s)", rr.Code, rr.Body.String()) + } +} + +// An empty holder defaults to the authenticated user's email — the #879 +// tie-in: the identity the token carries is the identity the lease +// records unless the caller says otherwise. +func TestItemLease_DefaultHolderIsAuthenticatedIdentity(t *testing.T) { + f := setupLeaseFixture(t) + + rr := f.call(t, f.tokenA, "claim", map[string]any{}) + if rr.Code != http.StatusOK { + t.Fatalf("claim: expected 200, got %d (body: %s)", rr.Code, rr.Body.String()) + } + var resp struct { + Lease models.ItemLease `json:"lease"` + } + if err := json.Unmarshal(rr.Body.Bytes(), &resp); err != nil { + t.Fatalf("parse: %v", err) + } + if resp.Lease.Holder != "user-a@example.com" { + t.Errorf("default holder = %q, want the authenticated user's email", resp.Lease.Holder) + } +} + +// ttl_seconds outside (0, 24h] is refused with 400 — a zero body means +// the 15-minute default, not an instant or eternal lease. +func TestItemLease_TTLBounds(t *testing.T) { + f := setupLeaseFixture(t) + + for _, ttl := range []int{-5, 86401} { + rr := f.call(t, f.tokenA, "claim", map[string]any{"ttl_seconds": ttl}) + if rr.Code != http.StatusBadRequest { + t.Errorf("ttl_seconds=%d: expected 400, got %d (body: %s)", ttl, rr.Code, rr.Body.String()) + } + } + + // Default: no ttl_seconds → ~15 minutes. + rr := f.call(t, f.tokenA, "claim", map[string]any{}) + if rr.Code != http.StatusOK { + t.Fatalf("default-ttl claim: expected 200, got %d", rr.Code) + } + var resp struct { + Lease models.ItemLease `json:"lease"` + } + if err := json.Unmarshal(rr.Body.Bytes(), &resp); err != nil { + t.Fatalf("parse: %v", err) + } + until := time.Until(resp.Lease.ExpiresAt) + if until < 14*time.Minute || until > 16*time.Minute { + t.Errorf("default TTL = %v, want ~15m", until) + } +} + +// Claiming requires authentication — a lease with no accountable identity +// behind it is exactly the out-of-band lock the feature replaces. An +// unauthenticated cookie-less POST is refused upstream of the handler +// (the CSRF middleware answers 403 before auth resolves), so accept +// either refusal shape; the invariant is that NO lease gets written. +func TestItemLease_ClaimRequiresAuth(t *testing.T) { + f := setupLeaseFixture(t) + + rr := f.call(t, "", "claim", map[string]any{}) + if rr.Code != http.StatusUnauthorized && rr.Code != http.StatusForbidden { + t.Errorf("unauthenticated claim: expected 401/403, got %d (body: %s)", rr.Code, rr.Body.String()) + } + lease, err := f.srv.store.GetItemLease(f.item.ID) + if err != nil { + t.Fatalf("GetItemLease: %v", err) + } + if lease != nil { + t.Errorf("unauthenticated claim wrote a lease: %+v", lease) + } +} + +// The single-item GET carries the live lease and omits it once expired — +// expiry is absence on the read path, with no reaper involved. +func TestItemLease_GetItemCarriesLiveLeaseOmitsExpired(t *testing.T) { + f := setupLeaseFixture(t) + + get := func() map[string]any { + t.Helper() + req := httptest.NewRequest("GET", "/api/v1/workspaces/"+f.wsSlug+"/items/"+f.item.Slug, nil) + req.Header.Set("Authorization", "Bearer "+f.tokenA) + req.RemoteAddr = "127.0.0.1:0" + rec := httptest.NewRecorder() + f.srv.ServeHTTP(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("get item: %d %s", rec.Code, rec.Body.String()) + } + var m map[string]any + if err := json.Unmarshal(rec.Body.Bytes(), &m); err != nil { + t.Fatalf("parse item: %v", err) + } + return m + } + + if _, present := get()["lease"]; present { + t.Error("unclaimed item must not carry a lease key") + } + + if rr := f.call(t, f.tokenA, "claim", map[string]any{"holder": "runner"}); rr.Code != http.StatusOK { + t.Fatalf("claim: %d", rr.Code) + } + leaseVal, present := get()["lease"] + if !present { + t.Fatal("live lease missing from item GET") + } + leaseMap, _ := leaseVal.(map[string]any) + if leaseMap["holder"] != "runner" { + t.Errorf("lease.holder = %v, want runner", leaseMap["holder"]) + } + + // Expire it directly at the store (a crashed holder), then re-read: + // the key must be gone with no reaper having run. + if _, err := f.srv.store.ClaimItemLease(f.item.ID, "runner", -2*time.Second); err != nil { + t.Fatalf("expire lease: %v", err) + } + if _, present := get()["lease"]; present { + t.Error("expired lease leaked into the item GET") + } +} + +// The workspace item list decorates leased items and only leased items. +func TestItemLease_ListDecoratesLeasedItems(t *testing.T) { + f := setupLeaseFixture(t) + + if rr := f.call(t, f.tokenA, "claim", map[string]any{"holder": "runner"}); rr.Code != http.StatusOK { + t.Fatalf("claim: %d", rr.Code) + } + + req := httptest.NewRequest("GET", "/api/v1/workspaces/"+f.wsSlug+"/items", nil) + req.Header.Set("Authorization", "Bearer "+f.tokenA) + req.RemoteAddr = "127.0.0.1:0" + rec := httptest.NewRecorder() + f.srv.ServeHTTP(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("list items: %d %s", rec.Code, rec.Body.String()) + } + var items []map[string]any + if err := json.Unmarshal(rec.Body.Bytes(), &items); err != nil { + t.Fatalf("parse list: %v", err) + } + found := false + for _, it := range items { + if it["id"] == f.item.ID { + lease, present := it["lease"].(map[string]any) + if !present { + t.Fatal("leased item's list row carries no lease") + } + if lease["holder"] != "runner" { + t.Errorf("list lease.holder = %v, want runner", lease["holder"]) + } + found = true + } else if _, present := it["lease"]; present { + t.Errorf("unleased item %v carries a lease key", it["id"]) + } + } + if !found { + t.Fatal("seeded item missing from the list") + } +} diff --git a/internal/server/handlers_items.go b/internal/server/handlers_items.go index f82da53ee..d96e14139 100644 --- a/internal/server/handlers_items.go +++ b/internal/server/handlers_items.go @@ -81,6 +81,22 @@ func (s *Server) handleListItems(w http.ResponseWriter, r *http.Request) { } s.enrichItemsWithParent(workspaceID, result, visibleIDs) + // Decorate live execution leases (#1221) with ONE workspace query + // rather than widening the shared item scan: leases are absent on + // most rows, and the map only carries live ones (expiry filtered in + // SQL), so unleased items keep their key omitted. + if leases, err := s.store.ListItemLeases(workspaceID); err != nil { + writeInternalError(w, err) + return + } else if len(leases) > 0 { + for i := range result { + if lease, ok := leases[result[i].ID]; ok { + leaseCopy := lease + result[i].Lease = &leaseCopy + } + } + } + writeJSON(w, http.StatusOK, result) } @@ -919,6 +935,17 @@ func (s *Server) handleGetItem(w http.ResponseWriter, r *http.Request) { // See item_moved_to.go for the disclosure rule. item.MovedTo = s.movedToDestinations(r, item) + // Live execution lease (#1221). Read here rather than in the shared + // scan path: a lease is display/coordination state, absent on most + // items, and one point read on the single-item GET keeps the many + // item SELECT column lists untouched. Expired leases read as nil. + lease, err := s.store.GetItemLease(item.ID) + if err != nil { + writeInternalError(w, err) + return + } + item.Lease = lease + writeJSON(w, http.StatusOK, item) } diff --git a/internal/server/server.go b/internal/server/server.go index 6df2db24d..db1496160 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -1808,6 +1808,10 @@ func (s *Server) setupRouter() { r.Get("/star", s.handleGetItemStarStatus) r.Post("/star", s.handleStarItem) r.Delete("/star", s.handleUnstarItem) + // Execution lease (#1221): atomic claim/checkout so + // concurrent pollers can't both "start" an item. + r.Post("/claim", s.handleClaimItem) + r.Post("/release", s.handleReleaseItem) // Watches (TASK-2533): durable per-item subscriptions // for the padd event-stream / plugin-monitor nudge // pipeline. `pad watch ` / `pad watch remove `. diff --git a/internal/store/item_lease.go b/internal/store/item_lease.go new file mode 100644 index 000000000..64d3c73f2 --- /dev/null +++ b/internal/store/item_lease.go @@ -0,0 +1,182 @@ +package store + +import ( + "database/sql" + "fmt" + "time" + + "github.com/PerpetualSoftware/pad/internal/models" +) + +// Item execution lease (#1221): an atomic claim/checkout so two pollers +// that both read "unclaimed" cannot both proceed. The conditional UPDATE +// is the arbiter — the same protocol the event-outbox claim (TASK-2714) +// and orphan GC (BUG-2415) established: the predicate on the UPDATE +// decides the winner, so a race produces one winner and one typed +// refusal, never two winners. +// +// Expiry is the reaper: an expired lease is treated as absent by the +// claim predicate and by every read, so a crashed holder strands nothing +// and no sweep job exists. Timestamps are RFC3339 whole-second UTC TEXT, +// the items-table convention, so string comparison in SQL and time +// comparison in Go agree. + +// LeaseHeldError reports a claim or release refused because another +// holder's lease is live. It names the holder and expiry so the caller +// can decide to wait, skip, or escalate rather than guessing. +type LeaseHeldError struct { + Holder string + AcquiredAt time.Time + ExpiresAt time.Time +} + +func (e *LeaseHeldError) Error() string { + return fmt.Sprintf("item is leased to %q until %s", e.Holder, e.ExpiresAt.Format(time.RFC3339)) +} + +// ClaimItemLease atomically acquires (or, for the live holder, refreshes) +// the execution lease on an item. It succeeds iff no live lease exists or +// the caller already holds it; a refresh moves the expiry forward but +// keeps the original acquired_at — extending is not re-acquiring. A +// negative ttl writes an already-expired lease (used by tests to model a +// crashed holder; the handler bounds ttl before it gets here). +// +// On contention it returns *LeaseHeldError naming the live holder. +func (s *Store) ClaimItemLease(itemID, holder string, ttl time.Duration) (*models.ItemLease, error) { + nowStr := now() + expiresStr := time.Now().UTC().Add(ttl).Format(time.RFC3339) + + // The predicate is the arbiter: unclaimed, expired, or already ours. + // All SET expressions read the PRE-update row (standard SQL, both + // dialects), so the CASE sees the old holder/expiry. + res, err := s.db.Exec(s.q(` + UPDATE items SET + lease_acquired_at = CASE + WHEN lease_holder = ? AND lease_expires_at > ? THEN lease_acquired_at + ELSE ? + END, + lease_holder = ?, + lease_expires_at = ? + WHERE id = ? + AND (lease_holder IS NULL OR lease_expires_at IS NULL OR lease_expires_at <= ? OR lease_holder = ?) + `), holder, nowStr, nowStr, holder, expiresStr, itemID, nowStr, holder) + if err != nil { + return nil, err + } + n, err := res.RowsAffected() + if err != nil { + return nil, err + } + if n == 1 { + return s.readItemLeaseRow(itemID) + } + + // Lost the race (or the item id doesn't exist). Read the live lease + // for the refusal; if it expired or was released between our UPDATE + // and this read, the caller's retry will win — say so. + lease, err := s.GetItemLease(itemID) + if err != nil { + return nil, err + } + if lease == nil { + return nil, fmt.Errorf("claim on item %s did not apply and no live lease exists — item missing or lease released mid-claim; retry", itemID) + } + return nil, &LeaseHeldError{Holder: lease.Holder, AcquiredAt: lease.AcquiredAt, ExpiresAt: lease.ExpiresAt} +} + +// ReleaseItemLease clears the caller's lease on an item. Idempotent: +// releasing an absent or expired lease is a no-op (released=false), never +// an error — cleanup code must not special-case "did I still hold this". +// Releasing another holder's LIVE lease is refused with *LeaseHeldError. +func (s *Store) ReleaseItemLease(itemID, holder string) (bool, error) { + res, err := s.db.Exec(s.q(` + UPDATE items SET lease_holder = NULL, lease_acquired_at = NULL, lease_expires_at = NULL + WHERE id = ? AND lease_holder = ? + `), itemID, holder) + if err != nil { + return false, err + } + n, err := res.RowsAffected() + if err != nil { + return false, err + } + if n == 1 { + return true, nil + } + + lease, err := s.GetItemLease(itemID) + if err != nil { + return false, err + } + if lease != nil && lease.Holder != holder { + return false, &LeaseHeldError{Holder: lease.Holder, AcquiredAt: lease.AcquiredAt, ExpiresAt: lease.ExpiresAt} + } + return false, nil +} + +// GetItemLease returns the item's live lease, or nil — absent and expired +// are indistinguishable on purpose (expiry is absence on every read path). +func (s *Store) GetItemLease(itemID string) (*models.ItemLease, error) { + lease, err := s.readItemLeaseRow(itemID) + if err != nil { + return nil, err + } + if lease == nil || !lease.ExpiresAt.After(time.Now().UTC()) { + return nil, nil + } + return lease, nil +} + +// ListItemLeases returns every LIVE lease in a workspace, keyed by item +// id. Expired leases are filtered in SQL so the map carries only what a +// list rendering should decorate. +func (s *Store) ListItemLeases(workspaceID string) (map[string]models.ItemLease, error) { + rows, err := s.db.Query(s.q(` + SELECT id, lease_holder, lease_acquired_at, lease_expires_at + FROM items + WHERE workspace_id = ? AND lease_holder IS NOT NULL AND lease_expires_at > ? + `), workspaceID, now()) + if err != nil { + return nil, err + } + defer rows.Close() + + leases := make(map[string]models.ItemLease) + for rows.Next() { + var id, holder, acquiredAt, expiresAt string + if err := rows.Scan(&id, &holder, &acquiredAt, &expiresAt); err != nil { + return nil, err + } + leases[id] = models.ItemLease{ + Holder: holder, + AcquiredAt: parseTime(acquiredAt), + ExpiresAt: parseTime(expiresAt), + } + } + return leases, rows.Err() +} + +// readItemLeaseRow reads the raw lease columns without the liveness +// filter — the claim path needs the row it just wrote even when a test +// wrote it pre-expired. Returns nil when the item has no lease columns +// set (or no such item exists). +func (s *Store) readItemLeaseRow(itemID string) (*models.ItemLease, error) { + var holder, acquiredAt, expiresAt *string + err := s.db.QueryRow(s.q(` + SELECT lease_holder, lease_acquired_at, lease_expires_at FROM items WHERE id = ? + `), itemID).Scan(&holder, &acquiredAt, &expiresAt) + if err == sql.ErrNoRows { + return nil, nil + } + if err != nil { + return nil, err + } + if holder == nil || expiresAt == nil { + return nil, nil + } + lease := &models.ItemLease{Holder: *holder, ExpiresAt: parseTime(*expiresAt)} + if acquiredAt != nil { + lease.AcquiredAt = parseTime(*acquiredAt) + } + return lease, nil +} diff --git a/internal/store/item_lease_test.go b/internal/store/item_lease_test.go new file mode 100644 index 000000000..0a701a206 --- /dev/null +++ b/internal/store/item_lease_test.go @@ -0,0 +1,249 @@ +package store + +import ( + "errors" + "sync" + "testing" + "time" +) + +// Tests for the item execution lease (#1221): an atomic claim/checkout so +// two pollers that both read "unclaimed" cannot both proceed. The +// conditional UPDATE is the arbiter — the same protocol the event-outbox +// claim (migration 083 / TASK-2714) and orphan GC (BUG-2415) established. + +func leaseFixture(t *testing.T) (*Store, string) { + t.Helper() + s := testStore(t) + ws := createTestWorkspace(t, s, "Lease WS") + col := createTestCollection(t, s, ws.ID, "Tasks") + item := createTestItem(t, s, ws.ID, col.ID, "Contended item", "") + return s, item.ID +} + +// A first claim on an unclaimed item succeeds and returns the lease. +func TestClaimItemLease_UnclaimedSucceeds(t *testing.T) { + s, itemID := leaseFixture(t) + + lease, err := s.ClaimItemLease(itemID, "sweep-runner", 15*time.Minute) + if err != nil { + t.Fatalf("claim: %v", err) + } + if lease.Holder != "sweep-runner" { + t.Errorf("holder = %q, want sweep-runner", lease.Holder) + } + if !lease.ExpiresAt.After(time.Now().UTC().Add(14 * time.Minute)) { + t.Errorf("expiry %v not ~15m out", lease.ExpiresAt) + } +} + +// A second claim by a different holder while the lease is live fails with +// a typed LeaseHeldError naming the holder and expiry — never a silent +// second winner. +func TestClaimItemLease_ContendedReturnsLeaseHeld(t *testing.T) { + s, itemID := leaseFixture(t) + + if _, err := s.ClaimItemLease(itemID, "winner", 15*time.Minute); err != nil { + t.Fatalf("first claim: %v", err) + } + _, err := s.ClaimItemLease(itemID, "loser", 15*time.Minute) + var held *LeaseHeldError + if !errors.As(err, &held) { + t.Fatalf("want LeaseHeldError, got %v", err) + } + if held.Holder != "winner" { + t.Errorf("error names holder %q, want winner", held.Holder) + } + if held.ExpiresAt.IsZero() { + t.Error("error must carry the expiry so the caller can decide to wait or skip") + } +} + +// N concurrent claimers produce exactly one winner; every loser gets +// LeaseHeldError. The predicate on the UPDATE is the arbiter. +func TestClaimItemLease_ConcurrentSingleWinner(t *testing.T) { + s, itemID := leaseFixture(t) + + const n = 8 + var wg sync.WaitGroup + errs := make([]error, n) + for i := 0; i < n; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + _, errs[i] = s.ClaimItemLease(itemID, string(rune('a'+i)), 15*time.Minute) + }(i) + } + wg.Wait() + + wins := 0 + for i, err := range errs { + switch { + case err == nil: + wins++ + default: + var held *LeaseHeldError + if !errors.As(err, &held) { + t.Errorf("claimer %d got a non-lease error: %v", i, err) + } + } + } + if wins != 1 { + t.Errorf("winners = %d, want exactly 1", wins) + } +} + +// A re-claim by the live holder refreshes the expiry (heartbeat) and +// keeps the original acquired_at — extending is not re-acquiring. +func TestClaimItemLease_HolderReclaimRefreshes(t *testing.T) { + s, itemID := leaseFixture(t) + + first, err := s.ClaimItemLease(itemID, "holder", 100*time.Second) + if err != nil { + t.Fatalf("first claim: %v", err) + } + second, err := s.ClaimItemLease(itemID, "holder", 30*time.Minute) + if err != nil { + t.Fatalf("re-claim by holder must succeed: %v", err) + } + if !second.ExpiresAt.After(first.ExpiresAt) { + t.Errorf("expiry did not move forward: %v -> %v", first.ExpiresAt, second.ExpiresAt) + } + if !second.AcquiredAt.Equal(first.AcquiredAt) { + t.Errorf("acquired_at changed on refresh: %v -> %v", first.AcquiredAt, second.AcquiredAt) + } +} + +// An expired lease is absent: a new holder claims through it with no +// reaper having run. +func TestClaimItemLease_ExpiredIsClaimable(t *testing.T) { + s, itemID := leaseFixture(t) + + if _, err := s.ClaimItemLease(itemID, "crashed", -2*time.Second); err != nil { + t.Fatalf("seed expired lease: %v", err) + } + lease, err := s.ClaimItemLease(itemID, "next", 15*time.Minute) + if err != nil { + t.Fatalf("claim over an expired lease must succeed: %v", err) + } + if lease.Holder != "next" { + t.Errorf("holder = %q, want next", lease.Holder) + } +} + +// GetItemLease returns nil for absent AND for expired leases — expiry is +// absence on every read path. +func TestGetItemLease_ExpiredReadsAsAbsent(t *testing.T) { + s, itemID := leaseFixture(t) + + if lease, err := s.GetItemLease(itemID); err != nil || lease != nil { + t.Fatalf("unclaimed item: lease=%v err=%v, want nil,nil", lease, err) + } + if _, err := s.ClaimItemLease(itemID, "crashed", -2*time.Second); err != nil { + t.Fatalf("seed expired lease: %v", err) + } + if lease, err := s.GetItemLease(itemID); err != nil || lease != nil { + t.Errorf("expired lease: lease=%v err=%v, want nil,nil", lease, err) + } + if _, err := s.ClaimItemLease(itemID, "live", 15*time.Minute); err != nil { + t.Fatalf("live claim: %v", err) + } + lease, err := s.GetItemLease(itemID) + if err != nil || lease == nil { + t.Fatalf("live lease: lease=%v err=%v, want non-nil", lease, err) + } + if lease.Holder != "live" { + t.Errorf("holder = %q, want live", lease.Holder) + } +} + +// Release by the holder clears the lease; the item is claimable again. +func TestReleaseItemLease_HolderReleases(t *testing.T) { + s, itemID := leaseFixture(t) + + if _, err := s.ClaimItemLease(itemID, "holder", 15*time.Minute); err != nil { + t.Fatalf("claim: %v", err) + } + released, err := s.ReleaseItemLease(itemID, "holder") + if err != nil { + t.Fatalf("release: %v", err) + } + if !released { + t.Error("release by the live holder should report released=true") + } + if _, err := s.ClaimItemLease(itemID, "someone-else", 15*time.Minute); err != nil { + t.Errorf("item must be claimable after release: %v", err) + } +} + +// Releasing an absent or expired lease is a no-op, never an error — +// cleanup code must not special-case "did I still hold this". +func TestReleaseItemLease_AbsentOrExpiredIsNoop(t *testing.T) { + s, itemID := leaseFixture(t) + + released, err := s.ReleaseItemLease(itemID, "nobody") + if err != nil { + t.Fatalf("release on unclaimed item: %v", err) + } + if released { + t.Error("nothing to release — released should be false") + } + if _, err := s.ClaimItemLease(itemID, "crashed", -2*time.Second); err != nil { + t.Fatalf("seed expired lease: %v", err) + } + if _, err := s.ReleaseItemLease(itemID, "someone-else"); err != nil { + t.Errorf("releasing over an expired foreign lease must be a no-op, got %v", err) + } +} + +// Releasing another holder's LIVE lease is refused with LeaseHeldError — +// refuse-on-ambiguity, not last-writer-wins. +func TestReleaseItemLease_ForeignLiveRefused(t *testing.T) { + s, itemID := leaseFixture(t) + + if _, err := s.ClaimItemLease(itemID, "holder", 15*time.Minute); err != nil { + t.Fatalf("claim: %v", err) + } + _, err := s.ReleaseItemLease(itemID, "intruder") + var held *LeaseHeldError + if !errors.As(err, &held) { + t.Fatalf("want LeaseHeldError, got %v", err) + } + if held.Holder != "holder" { + t.Errorf("error names %q, want holder", held.Holder) + } +} + +// ListItemLeases returns only live leases, keyed by item id. +func TestListItemLeases_LiveOnly(t *testing.T) { + s := testStore(t) + ws := createTestWorkspace(t, s, "Lease WS") + col := createTestCollection(t, s, ws.ID, "Tasks") + live := createTestItem(t, s, ws.ID, col.ID, "live item", "") + expired := createTestItem(t, s, ws.ID, col.ID, "expired item", "") + unclaimed := createTestItem(t, s, ws.ID, col.ID, "unclaimed item", "") + + if _, err := s.ClaimItemLease(live.ID, "runner", 15*time.Minute); err != nil { + t.Fatalf("live claim: %v", err) + } + if _, err := s.ClaimItemLease(expired.ID, "crashed", -2*time.Second); err != nil { + t.Fatalf("expired claim: %v", err) + } + + leases, err := s.ListItemLeases(ws.ID) + if err != nil { + t.Fatalf("list leases: %v", err) + } + if len(leases) != 1 { + t.Fatalf("live leases = %d, want 1 (map: %v)", len(leases), leases) + } + if lease, ok := leases[live.ID]; !ok || lease.Holder != "runner" { + t.Errorf("live item's lease missing or wrong: %v", leases) + } + if _, ok := leases[expired.ID]; ok { + t.Error("expired lease leaked into the live listing") + } + if _, ok := leases[unclaimed.ID]; ok { + t.Error("unclaimed item appeared in the listing") + } +} diff --git a/internal/store/migrations/084_nul_invariant_triggers.sql b/internal/store/migrations/084_nul_invariant_triggers.sql index 95c905ae3..fa102a628 100644 --- a/internal/store/migrations/084_nul_invariant_triggers.sql +++ b/internal/store/migrations/084_nul_invariant_triggers.sql @@ -1173,6 +1173,24 @@ BEGIN SELECT RAISE(ABORT, 'pad_nul_invariant: items.last_modified_by must not contain a NUL'); END; +CREATE TRIGGER IF NOT EXISTS pad_nul_items_lease_holder_ins +BEFORE INSERT ON items +FOR EACH ROW WHEN NEW.lease_holder IS NOT NULL AND ( + instr(NEW.lease_holder, char(0)) > 0 +) +BEGIN + SELECT RAISE(ABORT, 'pad_nul_invariant: items.lease_holder must not contain a NUL'); +END; + +CREATE TRIGGER IF NOT EXISTS pad_nul_items_lease_holder_upd +BEFORE UPDATE OF lease_holder ON items +FOR EACH ROW WHEN NEW.lease_holder IS NOT NULL AND ( + instr(NEW.lease_holder, char(0)) > 0 +) +BEGIN + SELECT RAISE(ABORT, 'pad_nul_invariant: items.lease_holder must not contain a NUL'); +END; + CREATE TRIGGER IF NOT EXISTS pad_nul_items_slug_ins BEFORE INSERT ON items FOR EACH ROW WHEN NEW.slug IS NOT NULL AND ( diff --git a/internal/store/migrations/085_item_lease.sql b/internal/store/migrations/085_item_lease.sql new file mode 100644 index 000000000..80e75b6c7 --- /dev/null +++ b/internal/store/migrations/085_item_lease.sql @@ -0,0 +1,28 @@ +-- Migration 085: item execution lease (#1221). +-- +-- Two processes that both read an item's status and conclude "no one else +-- has started this" have no way to make that decision atomically — whatever +-- they do next happens after the read, so both proceed. The lease is a +-- first-class, time-bounded "someone is executing this right now" state, +-- acquired by a conditional UPDATE whose predicate is the arbiter (the +-- protocol migration 083 / TASK-2714 established for the event outbox, and +-- BUG-2415 before that for orphan GC). +-- +-- lease_expires_at doubles as the reaper: an expired lease is treated as +-- absent by every read and by the claim predicate, so a crashed holder +-- strands nothing and no sweep job exists. The lease is DELIBERATELY not +-- part of the item's content state — claiming bumps neither updated_at nor +-- the version history, so a lease cannot 409 a concurrent editor's +-- expected_updated_at token. +-- +-- NOTE: no `IF NOT EXISTS` — SQLite's ALTER TABLE ADD COLUMN rejects it. +-- Nullable with no default; NULL = unclaimed, so no backfill is needed. +ALTER TABLE items ADD COLUMN lease_holder TEXT; +ALTER TABLE items ADD COLUMN lease_acquired_at TEXT; +ALTER TABLE items ADD COLUMN lease_expires_at TEXT; + +-- The workspace-wide live-lease listing (item list decoration) scans by +-- workspace and expiry; the partial index keeps unleased rows out of it. +CREATE INDEX IF NOT EXISTS idx_items_lease + ON items(workspace_id, lease_expires_at) + WHERE lease_holder IS NOT NULL; diff --git a/internal/store/nulcolumns.go b/internal/store/nulcolumns.go index 5f8d17706..9f7bc3c2d 100644 --- a/internal/store/nulcolumns.go +++ b/internal/store/nulcolumns.go @@ -77,6 +77,13 @@ var nulColumns = []nulColumn{ {"items", "last_modified_by", classText}, {"items", "source", classText}, + // Execution lease (#1221). The holder is a freeform caller string + // (`--holder` / body `holder`), so it carries caller text. decodeJSON's + // NUL refusal (BUG-2803) already blocks the API path; the trigger is the + // Layer B guarantee for writers that bypass this binary. The two lease + // timestamp columns are server-composed RFC3339 — see nulExcluded. + {"items", "lease_holder", classText}, + // collections {"collections", "schema", classJSON}, {"collections", "settings", classJSON}, @@ -318,6 +325,8 @@ var nulExcluded = map[string]string{ "mcp_audit_log.error_kind": "server enum, mcp_audit.go", "users.recovery_codes": "newline-joined bcrypt hashes of server-generated codes; looks like JSON, is not", "workspace_members.collection_access": "validated enum all/selected", + "items.lease_acquired_at": "server-composed RFC3339 (store now()/Add); the caller supplies only ttl_seconds, an int", + "items.lease_expires_at": "server-composed RFC3339 (store now()/Add); the caller supplies only ttl_seconds, an int", "item_yjs_updates.update_data": "BINARY (BLOB/BYTEA), the only such column in either schema. Raw Yjs updates legitimately contain NUL bytes; Layer A exempts it for the same reason and TestBinaryColumnCensus pins that. Surfaced here when the census's type filter was widened to include BLOB affinity, which is correct — the decision to exclude it is a judgement, not an oversight.", } diff --git a/internal/store/nulscan.go b/internal/store/nulscan.go index 559b42389..266894e22 100644 --- a/internal/store/nulscan.go +++ b/internal/store/nulscan.go @@ -199,8 +199,8 @@ func (r *NULScanReport) ByColumn() map[string]int { // the DESTINATION. Nothing about what any layer REFUSES changed. // // COST, stated because an operator should not be surprised by it: one -// unindexed scan per protected column — 131 of them today (24 JSON-classed, -// 107 text), measured from NULProtectedColumns rather than counted by hand. +// unindexed scan per protected column — 132 of them today (24 JSON-classed, +// 108 text), measured from NULProtectedColumns rather than counted by hand. // There is no index that would help, since the predicate is a substring search // over the value. That is cheap next to the migration it guards, which reads // every row of every table anyway. diff --git a/internal/store/nulscan_test.go b/internal/store/nulscan_test.go index d279da039..8d54201e1 100644 --- a/internal/store/nulscan_test.go +++ b/internal/store/nulscan_test.go @@ -331,7 +331,7 @@ func mustExec(t *testing.T, db *sql.DB, query string, args ...any) { // list is the source and the comment is the copy. func TestScanCostFiguresMatchTheList(t *testing.T) { const ( - wantTotal = 131 + wantTotal = 132 wantJSON = 24 ) diff --git a/internal/store/pgmigrations/062_item_lease.sql b/internal/store/pgmigrations/062_item_lease.sql new file mode 100644 index 000000000..b22002461 --- /dev/null +++ b/internal/store/pgmigrations/062_item_lease.sql @@ -0,0 +1,26 @@ +-- Migration 062: item execution lease (#1221). Mirrors SQLite migration 085. +-- +-- Two processes that both read an item's status and conclude "no one else +-- has started this" have no way to make that decision atomically — whatever +-- they do next happens after the read, so both proceed. The lease is a +-- first-class, time-bounded "someone is executing this right now" state, +-- acquired by a conditional UPDATE whose predicate is the arbiter (the +-- protocol the event-outbox claim / TASK-2714 established, and BUG-2415 +-- before that for orphan GC). Dialect-uniform on purpose — no FOR UPDATE +-- SKIP LOCKED special case; one implementation of one behaviour. +-- +-- lease_expires_at doubles as the reaper: an expired lease is treated as +-- absent by every read and by the claim predicate, so a crashed holder +-- strands nothing and no sweep job exists. The lease is DELIBERATELY not +-- part of the item's content state — claiming bumps neither updated_at nor +-- the version history, so a lease cannot 409 a concurrent editor's +-- expected_updated_at token. +ALTER TABLE items ADD COLUMN IF NOT EXISTS lease_holder TEXT; +ALTER TABLE items ADD COLUMN IF NOT EXISTS lease_acquired_at TEXT; +ALTER TABLE items ADD COLUMN IF NOT EXISTS lease_expires_at TEXT; + +-- The workspace-wide live-lease listing (item list decoration) scans by +-- workspace and expiry; the partial index keeps unleased rows out of it. +CREATE INDEX IF NOT EXISTS idx_items_lease + ON items(workspace_id, lease_expires_at) + WHERE lease_holder IS NOT NULL; diff --git a/web/src/lib/types/index.ts b/web/src/lib/types/index.ts index 10bb9bed3..2795c2f23 100644 --- a/web/src/lib/types/index.ts +++ b/web/src/lib/types/index.ts @@ -634,6 +634,16 @@ export interface ItemMovedTo { moved_at?: string; } +// ItemLease is a live execution lease on an item (#1221). holder is a +// freeform identity string (defaults server-side to the authenticated +// user); the lease is live iff expires_at is in the future, and the server +// never returns an expired one. +export interface ItemLease { + holder: string; + acquired_at: string; + expires_at: string; +} + export interface Item { id: string; workspace_id: string; @@ -691,6 +701,12 @@ export interface Item { // means "nothing to show", never "there is one you may not see". Do not // render a distinction between the two; there isn't one. moved_to?: ItemMovedTo[]; + // Live execution lease (#1221): who is actively executing this item right + // now, distinct from assigned_user_id's longer-term ownership. Present on + // the single-item GET and on list responses; the server omits the key for + // unclaimed items AND for expired leases (expiry is absence — there is no + // reaper), so `undefined` always means "no one holds this right now". + lease?: ItemLease; derived_closure?: ItemDerivedClosure; code_context?: ItemCodeContext; convention?: ItemConventionMetadata;