diff --git a/design/runner-protocol.md b/design/runner-protocol.md index 54eb6cae..71b9bfcb 100644 --- a/design/runner-protocol.md +++ b/design/runner-protocol.md @@ -1,12 +1,13 @@ # Runner protocol — sequential jobs on warm workers -**Status:** design agreed, not yet implemented. Replaces the existing +**Status:** implemented (`rust/crates/server/src/runner.rs`: `/runner/poll`, +`/runner/result`). Replaces the existing backends outright: `dispatch_docker`/`dispatch_serve`/`dispatch_fly`, the `Backend` enum, the worker-slot semaphore, `caos entrypoint`, and `caos serve` are all deleted, and the dev stack gains `caos runnerd` as a required daemon. Builds on the runner-pool decomposition (`runner-pool-and-cloud-builds.md`): that doc removes the per-worker *image*; this one removes the per-job -*container start*. +*container start*. See also [actors](../std/actor/README.md) for state that outlives a job. --- diff --git a/std/README.md b/std/README.md index a51c1c84..5785bc77 100644 --- a/std/README.md +++ b/std/README.md @@ -48,6 +48,7 @@ test suite exercises them. | `llm-call` | A single model call as an entry, for an expression that wants one without a conversation. | | `rgrep` | The search worker behind the step's `grep` tool. | | `run-and-update-ref` | The async worker: one binary, two stages, behind `run_async` and the subagent tools. | +| `actor` | The actor wrapper: runs an inner `(state, message) -> (state', reply)` request against state on a Git branch, publishing by moving the branch with a compare-and-swap ([`actor/README.md`](actor/README.md)); a Go program on `std/go`. | | `hello` | The smallest possible entry, used by `tests/hello` and by hand when something is deeply broken. | | `llm-stub` | A scripted stand-in for the model, so `tests/llm-*` run with no API key and no network. | | `llm-test` / `llm-test-tool` | Fixtures the llm tests drive: a test harness entry and a tool for it to call. | diff --git a/std/actor/.caos-expr b/std/actor/.caos-expr new file mode 100644 index 00000000..d1c54c5d --- /dev/null +++ b/std/actor/.caos-expr @@ -0,0 +1,5 @@ +# The actor wrapper (README.md), a std/go worker. One program, two +# positions like run-and-update-ref: start reads the branch head and tail-calls +# the inner request; finish publishes the new state by moving the branch with a +# compare-and-swap. std/go has git, which start uses to read the head. +curry --base:@=DEEP-DEPS/go --worker1:@=worker.go diff --git a/std/actor/DEPS b/std/actor/DEPS new file mode 100644 index 00000000..b70714ab --- /dev/null +++ b/std/actor/DEPS @@ -0,0 +1,2 @@ +# The image this worker runs on (format ` `). +../go go diff --git a/std/actor/README.md b/std/actor/README.md new file mode 100644 index 00000000..dc0d4beb --- /dev/null +++ b/std/actor/README.md @@ -0,0 +1,368 @@ +# Actors — persistent state in Git, single writer by compare-and-swap + +**Status:** wrapper (`std/actor`) and a reference key-value inner +(`tests/actor`) implemented; the spike is resolved (see [Spike results](#spike-results)). +`tests/actor` covers the build plan's cases: concurrent writers converging, a +forced lost race that fails uncached and succeeds on retry, an idempotent +re-apply (which is also the crash-after-push retry), a read making no commit, +the inner's lazy view of the state, and the inner's cache hit. Two caveats: a +real crash between push and reply is not injected (a re-applied message is the +same observable), and laziness is checked from inside the inner, not by +counting server object reads. Both are Go programs on `std/go`. Open question 6 +(history fetched on every write) is resolved for the write path: the wrapper +moves the branch with a direct receive-pack command and an empty pack instead +of `git push`, so it fetches no history (see the question). Daemons are deliberately set aside; see +[Deferred: daemons](#deferred-daemons). + +Builds on [client-owned conversation refs](../../design/client-owned-conversation-refs.md) +(the Git protocol workers already use for conversation heads) and follows the +start/finish shape of `std/run-and-update-ref`. + +--- + +## Problem + +caos runs pure, cached, hermetic jobs to completion. We also want named things +whose state outlives any one job: small servers, and eventually a test stack +that conversations drive. Today the only way to keep state is to hand it back +through a caller. + +## Model + +An **actor** is a pure function from `(state, message)` to `(state', reply)`, +with its `state` kept on a Git branch. Cloudflare Durable Objects are the +closest analog. Four rules define it: + +1. **No Start message.** Nothing is created or started. The first request for a + name finds an empty branch and runs against empty state. +2. **State lives on a branch.** Each request names the branch. The head's tree + is the state. +3. **A lost race fails the request.** Publishing the new state is a + compare-and-swap on the ref (`--force-with-lease=:`). If it + loses, the request fails and the caller retries. caos already treats failures + as retryable: they are never cached. +4. **Messages are idempotent.** Applying a message twice has the same effect as + applying it once. That is the whole duplicate-delivery story: caos does not + dedupe for the actor and the actor does not dedupe for itself. + +Nothing in the server changes. An actor is a std tool plus a convention for the +inner worker. + +### Principle: never manifest the whole tree + +caos never materializes a whole Git tree, and actors follow that. The state +reaches the inner worker as a **tree oid** that it reads lazily +(`caos get /cas/args/state/foo`), and the new state leaves as an oid it staged +with `caos put`. Neither the wrapper nor the inner ever checks the state out. + +## Design + +### Request contract + +An actor request is a call to the `actor` wrapper tool with these args: + +| arg | meaning | +|---|---| +| `state-ref` | the branch holding the actor's state: `refs/heads/actors/` | +| `inner` | the inner actor: any caos worker request template, in any image | +| `nonce` | any value that makes this ArgTree unique, so the outer request is never answered from the cache | +| `message` | the message: a blob or a tree, opaque to the wrapper and passed to the inner unchanged | + +`message` is one entry so that a message field can never collide with `state` +or with the wrapper's own args, and so the inner's input and output are +symmetric (`{state, message}` in, `{state, reply}` out). + +The nonce is only a cache-buster. Reusing it across retries or minting a new +one per attempt are both correct, because messages are idempotent. + +### Branch layout + +``` +state/ the actor's state: an ordinary tree, opaque to caos +``` + +The branch tree is `{state}`. The wrapper hands the inner only the `state/` +subtree, so everything else in the tree is the wrapper's. Nothing else is used +today and nothing is reserved; a later feature (claims, counters) can add a +sibling of `state/` without changing the inner's view. Every update is **one +commit whose parent is the observed head**, a linear chain. + +### The inner actor + +The inner is an ordinary caos worker with **any image**, including a +`docker://` image or a flake. It is a pure function of its args: + +| | | +|---|---| +| in | `/cas/args` with two entries, read lazily with `caos get`: `state`, the tree oid of the `state/` subtree (empty for a new actor), and `message` | +| out | `/cas/out`: a tree `{state, reply}`. `state` is the new state tree (staged with `caos put`); `reply` is a blob or tree | + +Two choices here differ from a first instinct: + +- **New state goes inside `/cas/out`, not beside it.** The result is what caos + caches, so a separate `/cas/state-out` outside it would be lost on a cache + hit. (`/cas/out-trace` is the precedent for data that deliberately stays out + of the cache key; state must not.) A small helper in `worker-common` can + expose `state-out` as a path to authors. +- **The state is passed as a tree oid, not a commit.** A commit hash includes + its parent, so every update would change the inner's cache key and nothing + would ever be shared. A tree oid is content-only: the same state and message + hit the cache. + +Because the inner is a pure function of `(state tree, message)`, **caching it is +correct**. That is why the nonce goes only on the outer request. A retry after +a lost race reaches a different head and so a different inner request, but a +retry after an unrelated failure with the head unchanged reuses the cached +inner result. + +An inner that finishes with `state` equal to its input makes **no commit and no +push**. Reads therefore never race with writers, and observe the head at some +moment during the request. + +### The wrapper: a start/finish pair + +The wrapper has two positions, like `run-and-update-ref`. It holds no container +while the inner runs, and it needs `git` (selected through `git-runner` in its +`.caos-expr`), which the inner does not. + +**Start** + +1. Read the branch head's oid with `git ls-remote`, then fetch **only that + commit** (`--depth=1 --filter=tree:0`; the server allows filters and + fetch-by-oid). The commit names its root tree, and the `state/` entry in that + tree gives the state oid. Absent branch means empty state. +2. Take the `state/` subtree oid from the head. +3. Build the inner request R from `inner`, `state` (the oid) and `message`. +4. Emit `run-request-then R`, carrying the observed head (and the branch name) + into the callback. + +**Finish**, given R's result `{state, reply}` and the observed head: + +1. If `state` equals the input state, return `reply`. Nothing to publish. +2. Otherwise build the root tree `{state: }` by oid, with no + checkout, and a commit on the observed head. +3. Make the commit available to a throwaway scratch repository (origin + `CAOS_SERVER_URL`) **without downloading the new state tree**. The new state + exists only on the server, so this is the main technical risk; see the spike. +4. Push with `git push --force-with-lease=: :`. + - ok: return `reply`; + - lease rejected: **fail the request**; the caller retries; + - ambiguous: fetch the ref again. Treat "my commit is the head" as success, + "head changed" as a lost race, and "head unchanged" as an infrastructure + failure. + +The window between start and finish is the race window. It is small compared +with the inner's run time, and the compare-and-swap makes it safe. + +### Failure, retry and caching + +| event | what happens | +|---|---| +| lost race | finish fails; not cached; caller retries; retry reads the new head | +| inner fails | the request fails; not cached | +| crash after the push, before the reply is posted | the job fails; a retry re-applies the message, which is idempotent | +| inner succeeds, finish fails | the inner's result stays cached; a retry on an unchanged head reuses it | +| duplicate concurrent requests | single-flight coalesces identical outer requests; distinct nonces both run and one loses the race | + +## Research: what the existing code gives us + +### 1. Compare-and-swap on the server's Git transport + +- The server delegates smart-HTTP to `git http-backend` + (`rust/crates/server/src/git.rs`). It does not set + `receive.denyNonFastForwards`, and does not need to: the expected old value is + part of the push command, which receive-pack checks under the ref lock, so + `--force-with-lease=:` is a server-side atomic compare-and-swap. + Fast-forward-ness is **not** enforced by the server; the wrapper keeps the + chain linear by always committing on the observed head. This matches the + "transport, not policy" stance of the conversation-refs doc. +- That doc already defines the client protocol used above: a scratch repo whose + origin is `CAOS_SERVER_URL`, an exact-ref fetch, a lease push, and the + re-fetch rule after an ambiguous failure. +- Git is opt-in: only workers bound to `git-runner` have `git`. + +### 2. Failures are not cached; successes are + +In `compute.rs` (`run_dispatch_inner`), only an `Ok` result reaches +`cache_set`; an `Err` is returned uncached. Successes are cached under the +ArgTree hash with no expiry (Redis is best-effort, and the key is namespaced by +`cache_namespace`). I did not find the server retrying a failed job by itself, +so the retry in rule 3 is the **caller's**. + +### 3. No credentials needed + +The Git paths are unauthenticated: `handle()` in `main.rs` routes them before +anything else, and only `/runner/*` checks a token. A wrapper needs only +`CAOS_SERVER_URL`, which every worker has. + +### 4. Existing machinery to reuse + +**`std/run-and-update-ref`** is the async worker behind `llm-step`'s +`run_async` and `spawn_agent` tools (bound in `std/llm-step/.caos-expr`, tested +in `tests/run-and-update-ref`). For `run_async` it runs an already-built +request and appends the task's terminal status, with the result oid, to the +conversation ref that started it. For `spawn_agent` it checkpoints the child +conversation's head onto the parent. It has the start/finish structure the +actor wrapper wants (start emits `run-request-then`, finish updates a ref), and +`refs.rs` has the exact-ref fetch and lease-push logic for conversation refs. +The conversation semantics live in `refs.rs`, but the Git plumbing it uses is +generic and sits in the `conversation-protocol` crate (`git-cli` feature), which +the actor wrapper can depend on directly: + +- `GitStore::scratch(name, remote)` makes a bare scratch repo whose `origin` is + the server. `read_ref` is a cheap `ls-remote`. `push(&[RefUpdate])` pushes + with `--force-with-lease=:` (and `--atomic` for several refs). + `GitStore` also implements `ObjectStore` (`read_tree`, `write_tree`, + `write_commit`, ...), so it can write the commit. +- `cas_append` in `refs.rs` is the right ambiguous-push rule, already tested + with a fake store: after a failed push, re-read the ref; if the candidate is an + ancestor of the observed head the push succeeded; if the head is unchanged the + failure is real; otherwise it was a lost race. + +**One thing not to reuse as is:** `GitStore::fetch_ref` fetches with no depth +and no filter into a scratch repo that is cleared for every job. For a +conversation ref that is cheap by design (its trees hold gitlinks). For an actor +it would download the whole history and the whole state closure on every +request. The wrapper must use `read_ref` plus a depth-1, `tree:0` fetch of the +head commit instead. + +**`TreeBuilder`** (`conversation_protocol::v3::tree`) builds trees by oid with +no checkout: `put_oid(path, mode, oid)`, `delete(path)`, `build(store)`. +`llm-step` uses it to seed a child conversation from the parent's tree. The +wrapper can use it to build `{state: }` the same way. + +### Caveats + +- Anyone who can reach the server can rewrite an actor branch. Conversation refs + already accept this; actors inherit it. +- Git advertises every ref on every push and fetch, so the number of actors is a + soft scaling limit. +- GC is deliberately off, so actor history is never reclaimed. + +## Build plan + +The wrapper is a new std tool, `std/actor`: a single Go program run by `std/go` +(Go is the language for new workers), with start and finish as two positions of +one program like `std/run-and-update-ref`. It shells out to the `caos` CLI and, +for the head read, to `git`. The sections above describe the first design, a +Rust wrapper that pushed from a scratch repository; the shipped wrapper replaced +that with the direct receive-pack command described under open question 6. + +0. **Spike (verify before building).** An integration test against the test + stack, using real `git`: + - **Push a commit whose new state tree exists only on the server.** Stage a + state tree with `caos put`, build a commit on the observed head that points + at it (`TreeBuilder` plus `write_commit`), and push it with a lease, without + ever fetching that tree into the scratch repo. My best guess is a partial + (promisor) scratch repo, where the missing objects are "promised" and the + push sends nothing the server lacks; I have not tried it. Fallbacks, both + worse: fetch the state closure (fine for small state, but it breaks the + no-manifest rule), or add a small push-by-oid endpoint to the server (which + breaks "no server change"). + - Confirm the depth-1, `tree:0` head fetch reads the state oid cheaply. + - Confirm a start/finish pair can carry the observed head through the + callback. +1. **Wrapper and a reference inner.** The inner is a small key-value actor + (`put`, `get`; `put` is idempotent) in a non-runner image, which exercises + the "any image" claim. +2. **Tests.** + - concurrent `put`s to one actor, retried on failure, converge with no lost + update; + - a forced lost race fails the request and is not cached; + - a crash after the push followed by a retry reaches the same state; + - a read makes no commit; + - **laziness:** a state with many entries and a message touching one of them + fetches only that entry's objects (assert on the server's object reads); + - the inner's result is a cache hit when state and message repeat. +3. **Docs.** Refresh the `runner-protocol.md` status line (it still says "not + yet implemented", but `runner.rs` implements it), and link this doc. + +### Spike results + +Measured against the test stack with real `git`: + +- **The promisor guess works, with one addition.** A scratch repo configured as + a partial clone of the server (`extensions.partialClone=origin`, + `remote.origin.promisor=true`, filter `tree:0`) can push a commit whose new + state tree was never downloaded, but only if that tree's *root object* was + fetched through the filter first. A commit pointing at an object the repo has + never seen fails in pack-objects (`Could not read `); one fetched with + `--filter=tree:0 origin ` makes its children "promised" and the + push goes through. Cost: one tree object per update. +- **Shallow fetches cannot push.** The server answers `shallow pushes are not + accepted`. Start may read the head with `--depth=1`; finish, which pushes, + fetches the parent commit without depth (commits only, no trees). That is + linear in history length, so it sharpens open question 2 (history growth). +- The start/finish pair carries the observed head and input state through the + callback by currying them onto the wrapper's own ArgTree. + +## Non-goals + +- Server-side ordering, leases or ref policy. The server stays transport. +- Built-in dedupe, reply caching or exactly-once delivery. +- External effects. The inner must be pure; effectful actors are an optional + later pattern (a write-ahead claim commit in a sibling of `state/`). +- Strong isolation between actors. + +## Deferred: daemons + +Daemons are not designed here. The direction I would take when we return to +them: keep **state and liveness** in a pure actor (messages like `claim`, +`renew` and `release` with a lease), run the live process as a **detached +long-running job in the author's own image** whose supervisor only sends caos +requests, and let clients reach it directly at an address recorded in the actor +state. That keeps git and the runner token out of the author's image. The +unverified parts are how a custom image also gets the supervisor, how a detached +job behaves across a server restart, runner-network reachability, and idle +detection. + +## Open questions + +1. **Caller retry.** The design assumes callers retry a failed request. Which + layer owns that for conversations and `map-then`, and does it back off under + contention on a hot actor? +2. **History growth.** Every state change is a commit and GC is off. Options: + periodic squash into a new root, or a per-actor compaction worker. +6. **Finish fetches every commit on the branch.** The spike showed that a push + cannot come from a shallow repository, so finish fetches the parent commit + without `--depth` (`--filter=tree:0`, so commits only, no trees). That + downloads the actor's whole commit history on every write: cost and latency + grow linearly with the number of updates, and GC is off, so the history + never shrinks. Start is unaffected (it reads the head at depth 1). This is + the same problem as question 2 seen from the write path, and it is the + reason that question matters now rather than later. Options to investigate: + - make the parent promised rather than present, by dropping the `shallow` + file after a depth-1 `tree:0` fetch so the head is the only commit held + and its parent is absent but listed by a `.promisor` pack. **Tried; it + does not work with `git push`** (40-commit branch, empty scratch repo): + the push fails with `Could not read ` / `could not parse commit + `, with or without `--no-thin`. The depth-1 pack is marked + promisor, but the pack-objects that `send-pack` starts walks the head's + parents to mark them uninteresting and does not tolerate a missing one. + With the `shallow` file kept, the server refuses the push instead + (`shallow pushes are not accepted`). Only the full-history fetch + pushes; + - **bypass `git push` and speak receive-pack directly. Tried; it works** + (`tests/actor-ref`). A push is a command ` ` plus a pack, + and the pack may be empty when the server already has the new object. A + commit made with `caos put-commit` is already on the server, so finish + POSTs one pkt-line command and an empty pack to + `$CAOS_SERVER_URL/git-receive-pack`: no scratch repository, no promisor + setup, no fetch of the parent or of any history. The server does the + compare-and-swap: a stale `` is answered `ng ` and the ref does + not move, and the right `` is accepted. The result is an ordinary + branch that git can fetch. This resolves the history cost of this + question for the write path, and needs no server change. `std/actor` + should use it; + - squash periodically (question 2), which bounds the chain; + - add a small server-side push-by-oid endpoint, which breaks "no server + change". +3. **Read consistency.** A read observes some head during its run and is not + linearized against concurrent writes. Is that acceptable, or should a read + optionally confirm the head at the end? +4. **Large states.** A state change that touches one path rewrites one path's + spine in the tree. Does the inner have an easy way to build a new state from + the old by oid, with no checkout? This is the main usability question for + authors, and the `state-out` helper in `worker-common` is meant to answer it. +5. **Authorization.** Do we want per-namespace write control on Git pushes + before actors hold anything sensitive? diff --git a/std/actor/worker.go b/std/actor/worker.go new file mode 100644 index 00000000..e0855a5a --- /dev/null +++ b/std/actor/worker.go @@ -0,0 +1,283 @@ +// The actor wrapper (README.md): run an inner `(state, message) -> +// (state', reply)` request against state kept on a Git branch, and publish the +// new state with a compare-and-swap. +// +// `Q = actor { state-ref, inner, nonce, message }` has two positions: +// +// - start reads the branch head (`git ls-remote`, then a depth-1 `tree:0` +// fetch of that one commit), takes the `state/` subtree oid from the head, +// builds the inner request and tail-calls it with Q (plus the observed head +// and the input state) as the callback; +// - finish receives the inner's `{state, reply}`. An unchanged state returns +// the reply without touching Git; otherwise it mints `{state: }` +// as a commit on the observed head with `caos put-commit` and moves the +// branch with a compare-and-swap. A lost race fails the request, which is +// never cached, so the caller retries. +// +// Neither position checks the state out: it travels as a tree oid. +// +// THE BRANCH IS MOVED WITHOUT `git push`. A push is a command line +// " " plus a pack, and the pack may be empty when the server +// already has the new object, which it does: `caos put-commit` put it there. +// So finish POSTs that one command and an empty pack to git-receive-pack, and +// the server does the compare-and-swap (a stale is answered `ng`). No +// scratch repository, no fetch of the parent, no history. `git push` cannot +// do this: it resolves the new commit in a local repository and walks its +// ancestry to build a pack, which needs every ancestor commit ("a deep +// checkout"), and a partial clone with a promisor remote does not avoid that +// (README.md, open question 6; tests/actor-ref proves the direct route). +package main + +import ( + "bytes" + "crypto/sha1" + "fmt" + "io" + "net/http" + "os" + "os/exec" + "path/filepath" + "strconv" + "strings" + + "caos/w" +) + +const ( + stateEntry = "state" + noHead = "none" + zeros = "0000000000000000000000000000000000000000" + gitDir = "/tmp/actor-git" +) + +// run runs a command and returns its trimmed stdout, failing with its stderr. +func run(name string, args ...string) string { + cmd := exec.Command(name, args...) + cmd.Env = append(os.Environ(), "GIT_TERMINAL_PROMPT=0") + var stderr bytes.Buffer + cmd.Stderr = &stderr + out, err := cmd.Output() + w.True(err == nil, "%s %s: %v: %s", name, strings.Join(args, " "), err, strings.TrimSpace(stderr.String())) + return strings.TrimSpace(string(out)) +} + +func caos(args ...string) string { return run("caos", args...) } + +func exists(path string) bool { + _, err := os.Lstat(path) + return err == nil +} + +// readArg fetches a blob argument and returns it trimmed. +func readArg(name string) string { + path := "/cas/args/" + name + caos("get", path) + return strings.TrimSpace(string(w.Check(os.ReadFile(path)))) +} + +func serverURL() string { + url := strings.TrimRight(os.Getenv("CAOS_SERVER_URL"), "/") + w.True(url != "", "CAOS_SERVER_URL not set") + return url +} + +// readRef is the branch's head on the server, or "" if the branch is absent. +func readRef(ref string) string { + out := run("git", "ls-remote", "--refs", serverURL(), ref) + for _, line := range strings.Split(out, "\n") { + fields := strings.Fields(line) + if len(fields) == 2 && fields[1] == ref { + return fields[0] + } + } + return "" +} + +// stateOf is the `state/` subtree oid of head, reading only the commit and its +// root tree: a depth-1 `tree:0` fetch into a throwaway partial-clone repository. +// Start only reads, so cutting the history off is fine here. +func stateOf(head string) string { + w.Must(os.RemoveAll(gitDir)) + run("git", "init", "-q", "--bare", gitDir) + git := func(args ...string) string { return run("git", append([]string{"-C", gitDir}, args...)...) } + git("config", "core.repositoryformatversion", "1") + git("config", "extensions.partialClone", "origin") + git("config", "remote.origin.url", serverURL()) + git("config", "remote.origin.promisor", "true") + git("config", "remote.origin.partialclonefilter", "tree:0") + git("fetch", "--quiet", "--no-tags", "--no-write-fetch-head", "--depth=1", "--filter=tree:0", "origin", head) + for _, line := range strings.Split(git("ls-tree", head), "\n") { + // "040000 tree \t" + meta, name, ok := strings.Cut(line, "\t") + if ok && name == stateEntry { + return strings.Fields(meta)[2] + } + } + return "" +} + +// emptyState is the empty tree, as a CAS path (so it can be bound by path). +func emptyState() string { + dir := "/tmp/actor-empty-state" + w.Must(os.RemoveAll(dir)) + w.Must(os.MkdirAll(dir, 0o755)) + caos("put", dir, "/cas/empty-state") + return "/cas/empty-state" +} + +func main() { + w.Main(func() { + stateRef := readArg("state-ref") + w.True(strings.HasPrefix(stateRef, "refs/heads/actors/") && !strings.Contains(stateRef, ".."), + "state-ref %q must be under refs/heads/actors/", stateRef) + if exists("/cas/args/result") { + finish(stateRef) + } else { + start(stateRef) + } + }) +} + +func start(stateRef string) { + head := readRef(stateRef) + statePath, stateOid := "", "" + if head != "" { + stateOid = stateOf(head) + } + if stateOid != "" { + caos("get-hash", stateOid, "/cas/state") + statePath = "/cas/state" + } else { + statePath = emptyState() + stateOid = caos("hash", statePath) + } + + request := caos("prepare-request", "--base:@=/cas/args/inner", + "--state:@="+statePath, "--message:@=/cas/args/message") + + // The callback is this same Q, carrying what finish needs to publish. + q := caos("hash", "/cas/args") + headText := head + if headText == "" { + headText = noHead + } + callback := caos("curry", "--base:hash="+q, "--head="+headText, "--old-state="+stateOid) + caos("run-request-then", request, "--then:hash="+callback) +} + +func finish(stateRef string) { + result := "/cas/args/result" + // List the result's children as hash-tagged entries; nothing is downloaded. + caos("get", result) + newState := caos("hash", filepath.Join(result, stateEntry)) + if newState != readArg("old-state") { + head := readArg("head") + if head == noHead { + head = "" + } + publish(stateRef, head, newState) + } + caos("forward", filepath.Join(result, "reply"), "/cas/out") +} + +// publish mints the commit {state: newState} on head and moves the branch. +func publish(stateRef, head, newState string) { + // The root tree {state: }: a symlink to the already-fetched + // result entry, which `caos put` resolves to its recorded hash. + root := "/tmp/actor-root" + w.Must(os.RemoveAll(root)) + w.Must(os.MkdirAll(root, 0o755)) + w.Must(os.Symlink("/cas/args/result/"+stateEntry, filepath.Join(root, stateEntry))) + caos("put", root, "/cas/new-root") + tree := caos("hash", "/cas/new-root") + + text := "tree " + tree + "\n" + if head != "" { + text += "parent " + head + "\n" + } + text += "author actor 0 +0000\ncommitter actor 0 +0000\n\nactor state\n" + w.Must(os.WriteFile("/tmp/actor-commit", []byte(text), 0o644)) + candidate := caos("put-commit", "/tmp/actor-commit", "/cas/new-commit") + + old := head + if old == "" { + old = zeros + } + status, err := setRef(serverURL(), stateRef, old, candidate) + if err == nil && status == "" { + return + } + // Ambiguous or refused: re-read the ref to learn what actually happened. + switch observed := readRef(stateRef); { + case observed == candidate: + return + case observed == head: + w.True(false, "moving %s: %s %v", stateRef, status, err) + default: + w.True(false, "lost the race for %s: %s %v", stateRef, status, err) + } +} + +func pkt(s string) string { return fmt.Sprintf("%04x%s", len(s)+4, s) } + +// emptyPack is a valid pack holding no objects: "PACK", version 2, count 0, +// and the SHA-1 of those twelve bytes. +func emptyPack() []byte { + header := []byte("PACK\x00\x00\x00\x02\x00\x00\x00\x00") + sum := sha1.Sum(header) + return append(header, sum[:]...) +} + +// setRef asks git-receive-pack to move ref from old to new, sending no objects. +// It returns "" when the server accepted the update, else the server's +// complaint (a stale arrives as `ng `). +func setRef(url, ref, old, new string) (string, error) { + var body bytes.Buffer + body.WriteString(pkt(fmt.Sprintf("%s %s %s\x00 report-status agent=caos-actor\n", old, new, ref))) + body.WriteString("0000") + body.Write(emptyPack()) + req, err := http.NewRequest("POST", url+"/git-receive-pack", &body) + if err != nil { + return "", err + } + req.Header.Set("Content-Type", "application/x-git-receive-pack-request") + req.Header.Set("Accept", "application/x-git-receive-pack-result") + resp, err := http.DefaultClient.Do(req) + if err != nil { + return "", err + } + defer resp.Body.Close() + data, err := io.ReadAll(resp.Body) + if err != nil { + return "", err + } + if resp.StatusCode != 200 { + return "", fmt.Errorf("git-receive-pack answered %s: %s", resp.Status, data) + } + unpacked, accepted := false, false + var complaint []string + for rest := string(data); len(rest) >= 4; { + n, err := strconv.ParseUint(rest[:4], 16, 16) + if err != nil || n < 4 && n != 0 || int(n) > len(rest) { + return "", fmt.Errorf("malformed report-status: %q", data) + } + if n == 0 { + rest = rest[4:] + continue + } + line := strings.TrimSpace(rest[4:n]) + rest = rest[n:] + switch { + case line == "unpack ok": + unpacked = true + case line == "ok "+ref: + accepted = true + default: + complaint = append(complaint, line) + } + } + if unpacked && accepted && len(complaint) == 0 { + return "", nil + } + return strings.Join(complaint, "; "), nil +} diff --git a/tests/actor-ref/.caos-expr b/tests/actor-ref/.caos-expr new file mode 100644 index 00000000..f0cd6032 --- /dev/null +++ b/tests/actor-ref/.caos-expr @@ -0,0 +1,3 @@ +# SPIKE for std/actor/README.md open question 6 (see worker.go): a std/go worker +# that moves a branch by speaking git-receive-pack with an empty pack. +curry --base:@=DEEP-DEPS/go --worker1:@=worker.go diff --git a/tests/actor-ref/DEPS b/tests/actor-ref/DEPS new file mode 100644 index 00000000..dbcdefbb --- /dev/null +++ b/tests/actor-ref/DEPS @@ -0,0 +1,2 @@ +# What this test reaches for (format ` `). +../../std/go go diff --git a/tests/actor-ref/worker.go b/tests/actor-ref/worker.go new file mode 100644 index 00000000..b6d3e6f2 --- /dev/null +++ b/tests/actor-ref/worker.go @@ -0,0 +1,136 @@ +// tests/actor-ref — SPIKE for std/actor/README.md, open question 6. +// +// Can a worker move a branch with a compare-and-swap WITHOUT any scratch +// repository and without fetching any history? A git push is a command line +// " " plus a pack, and the pack may be empty if the server +// already has the new object. `caos put-commit` puts the commit on the server +// without a push, so this speaks git-receive-pack directly with an empty pack. +// +// Proves, against the test stack: +// 1. creating a ref (old = zeros) at a commit made by `caos put-commit`; +// 2. updating it with the right to a child commit; +// 3. a stale is REJECTED and the ref does not move; +// 4. the result is a normal branch: git can fetch it and see both commits. +package main + +import ( + "bytes" + "crypto/sha1" + "fmt" + "io" + "net/http" + "os" + "os/exec" + "strings" + + "caos/w" +) + +const zeros = "0000000000000000000000000000000000000000" + +func run(name string, args ...string) string { + cmd := exec.Command(name, args...) + cmd.Env = append(os.Environ(), "GIT_TERMINAL_PROMPT=0") + var stderr bytes.Buffer + cmd.Stderr = &stderr + out, err := cmd.Output() + w.True(err == nil, "%s %s: %v: %s", name, strings.Join(args, " "), err, strings.TrimSpace(stderr.String())) + return strings.TrimSpace(string(out)) +} + +func pkt(s string) string { return fmt.Sprintf("%04x%s", len(s)+4, s) } + +// emptyPack is a valid pack holding no objects: "PACK", version 2, count 0, +// and the SHA-1 of those twelve bytes. +func emptyPack() []byte { + header := []byte("PACK\x00\x00\x00\x02\x00\x00\x00\x00") + sum := sha1.Sum(header) + return append(header, sum[:]...) +} + +// setRef asks git-receive-pack to move ref from old to new, sending no +// objects. It returns the server's status lines. +func setRef(url, ref, old, new string) string { + var body bytes.Buffer + body.WriteString(pkt(fmt.Sprintf("%s %s %s\x00 report-status agent=caos-spike\n", old, new, ref))) + body.WriteString("0000") + body.Write(emptyPack()) + req := w.Check(http.NewRequest("POST", url+"/git-receive-pack", &body)) + req.Header.Set("Content-Type", "application/x-git-receive-pack-request") + req.Header.Set("Accept", "application/x-git-receive-pack-result") + resp := w.Check(http.DefaultClient.Do(req)) + defer resp.Body.Close() + text := string(w.Check(io.ReadAll(resp.Body))) + w.True(resp.StatusCode == 200, "git-receive-pack answered %s: %s", resp.Status, text) + return text +} + +// commit mints a commit over a state subtree holding v=, via +// `caos put` and `caos put-commit`, and returns its hash. Nothing is pushed. +func commit(tag, value, parent string) string { + dir := "/tmp/root-" + tag + w.Must(os.RemoveAll(dir)) + w.Must(os.MkdirAll(dir+"/state", 0o755)) + w.Must(os.WriteFile(dir+"/state/v", []byte(value+"\n"), 0o644)) + run("caos", "put", dir, "/cas/root-"+tag) + tree := run("caos", "hash", "/cas/root-"+tag) + text := "tree " + tree + "\n" + if parent != "" { + text += "parent " + parent + "\n" + } + text += "author actor 0 +0000\ncommitter actor 0 +0000\n\nactor state " + tag + "\n" + file := "/tmp/commit-" + tag + w.Must(os.WriteFile(file, []byte(text), 0o644)) + return run("caos", "put-commit", file, "/cas/commit-"+tag) +} + +func remoteHead(url, ref string) string { + out := run("git", "ls-remote", "--refs", url, ref) + if out == "" { + return "" + } + return strings.Fields(out)[0] +} + +func main() { + w.Main(func() { + url := strings.TrimRight(os.Getenv("CAOS_SERVER_URL"), "/") + w.True(url != "", "this test needs CAOS_SERVER_URL from the runner") + run("caos", "get", "/cas/args/test-salt") + salt := strings.TrimSpace(string(w.Check(os.ReadFile("/cas/args/test-salt")))) + ref := fmt.Sprintf("refs/heads/actors/ref-%s-%d", salt, os.Getpid()) + + w.Step("mint two commits with caos put-commit (no push)") + c1 := commit("one", "1", "") + c2 := commit("two", "2", c1) + w.True(remoteHead(url, ref) == "", "fresh ref already exists") + + w.Step("create the ref with an empty pack") + resp := setRef(url, ref, zeros, c1) + fmt.Fprintf(os.Stderr, "create response: %q\n", resp) + w.True(strings.Contains(resp, "ok "+ref), "create was not accepted: %q", resp) + w.True(remoteHead(url, ref) == c1, "ref is %q, want %s", remoteHead(url, ref), c1) + + w.Step("a STALE is rejected and the ref does not move") + resp = setRef(url, ref, zeros, c2) + fmt.Fprintf(os.Stderr, "stale response: %q\n", resp) + w.True(strings.Contains(resp, "ng "+ref), "a stale was not rejected: %q", resp) + w.True(remoteHead(url, ref) == c1, "the ref moved on a stale update") + + w.Step("update with the right ") + resp = setRef(url, ref, c1, c2) + fmt.Fprintf(os.Stderr, "update response: %q\n", resp) + w.True(strings.Contains(resp, "ok "+ref), "update was not accepted: %q", resp) + w.True(remoteHead(url, ref) == c2, "ref is %q, want %s", remoteHead(url, ref), c2) + + w.Step("it is an ordinary branch") + w.Must(os.RemoveAll("/tmp/check")) + run("git", "init", "-q", "--bare", "/tmp/check") + run("git", "-C", "/tmp/check", "fetch", "-q", url, ref) + w.True(run("git", "-C", "/tmp/check", "rev-list", "--count", "FETCH_HEAD") == "2", "history is not two commits") + w.True(run("git", "-C", "/tmp/check", "show", "FETCH_HEAD:state/v") == "2", "state/v is not 2") + w.True(run("git", "-C", "/tmp/check", "show", "FETCH_HEAD~1:state/v") == "1", "parent state/v is not 1") + + w.Report("actor-ref: ALL PASS\n") + }) +} diff --git a/tests/actor/.caos-expr b/tests/actor/.caos-expr new file mode 100644 index 00000000..fe5246ab --- /dev/null +++ b/tests/actor/.caos-expr @@ -0,0 +1,5 @@ +# tests/actor, as an ENTRY: a std/go worker test (std/go has git, which is how it +# reads the actor's branch on the server) driving std/actor with a reference +# key-value inner. `probe` is an impure inner and `mapper` a concurrent writer, +# both for the race and concurrency cases. Every program runs on std/go. +curry --base:@=DEEP-DEPS/go --worker1:@=worker.go --actor:@=DEEP-DEPS/actor --kv:@=kv.go --probe:@=probe.go --mapper:@=mapper.go diff --git a/tests/actor/DEPS b/tests/actor/DEPS new file mode 100644 index 00000000..c1ec3469 --- /dev/null +++ b/tests/actor/DEPS @@ -0,0 +1,3 @@ +# What this test reaches for (format ` `). +../../std/go go +../../std/actor actor diff --git a/tests/actor/kv.go b/tests/actor/kv.go new file mode 100644 index 00000000..693d9004 --- /dev/null +++ b/tests/actor/kv.go @@ -0,0 +1,79 @@ +// Reference inner actor (std/actor/README.md): a key-value store. A pure function +// of (state tree, message) -> {state, reply}. Messages are idempotent: +// +// put set key (applying twice is the same as once) +// get reply with the value, state unchanged +// getcheck like get, and fail if any OTHER entry's content was +// materialized: the inner sees the state lazily +// +// It runs as an actor's inner on std/go, and probe.go runs it too: probe copies +// this file into the prelude module as its own command, so keep it a single +// self-contained `package main`. +package main + +import ( + "os" + "os/exec" + "path/filepath" + "strings" + + "caos/w" +) + +func caos(args ...string) { + cmd := exec.Command("caos", args...) + cmd.Stderr = os.Stderr + w.True(cmd.Run() == nil, "caos %s failed", strings.Join(args, " ")) +} + +func main() { + w.Main(func() { + caos("get", "/cas/args/message") + line := strings.TrimRight(string(w.Check(os.ReadFile("/cas/args/message"))), "\n") + fields := strings.SplitN(line, " ", 3) + w.True(len(fields) >= 2, "kv: bad message: %q", line) + op, key, value := fields[0], fields[1], "" + if len(fields) == 3 { + value = strings.TrimSpace(fields[2]) + } + w.True(key != "" && !strings.Contains(key, "/") && !strings.HasPrefix(key, "."), "kv: bad key: %s", key) + + // List the state's entries without reading their content. + caos("get", "/cas/args/state") + entries := w.Check(os.ReadDir("/cas/args/state")) + + w.Must(os.RemoveAll("/tmp/out")) + w.Must(os.MkdirAll("/tmp/out", 0o755)) + switch op { + case "put": + w.Must(os.Mkdir("/tmp/out/state", 0o755)) + for _, e := range entries { + if e.Name() != key { + w.Must(os.Symlink(filepath.Join("/cas/args/state", e.Name()), filepath.Join("/tmp/out/state", e.Name()))) + } + } + w.Must(os.WriteFile(filepath.Join("/tmp/out/state", key), []byte(value+"\n"), 0o644)) + w.Must(os.WriteFile("/tmp/out/reply", []byte("ok\n"), 0o644)) + case "get", "getcheck": + w.Must(os.Symlink("/cas/args/state", "/tmp/out/state")) + reply := []byte{} + if _, err := os.Lstat(filepath.Join("/cas/args/state", key)); err == nil { + caos("get", filepath.Join("/cas/args/state", key)) + reply = w.Check(os.ReadFile(filepath.Join("/cas/args/state", key))) + } + w.Must(os.WriteFile("/tmp/out/reply", reply, 0o644)) + if op == "getcheck" { + for _, e := range entries { + if e.Name() == key { + continue + } + info := w.Check(os.Stat(filepath.Join("/cas/args/state", e.Name()))) + w.True(info.Size() == 0, "kv: %s was materialized by a read of %s", e.Name(), key) + } + } + default: + w.True(false, "kv: unknown op: %s", op) + } + caos("put", "/tmp/out", "/cas/out") + }) +} diff --git a/tests/actor/mapper.go b/tests/actor/mapper.go new file mode 100644 index 00000000..fffeee0c --- /dev/null +++ b/tests/actor/mapper.go @@ -0,0 +1,72 @@ +// One concurrent writer for tests/actor. Used as a map-then `map`, it is called +// with --in=; it sends that message to the actor and, when the +// request loses the race for the branch (the wrapper fails it, uncached), sends +// it again with a new nonce, up to maxAttempts. Its callback is this same +// program with --attempt and --msg curried on and --result or --error supplied. +package main + +import ( + "fmt" + "os" + "os/exec" + "strconv" + "strings" + + "caos/w" +) + +const maxAttempts = 24 + +func caos(args ...string) string { + cmd := exec.Command("caos", args...) + cmd.Stderr = os.Stderr + out, err := cmd.Output() + w.True(err == nil, "caos %s: %v", strings.Join(args, " "), err) + return strings.TrimSpace(string(out)) +} + +func exists(path string) bool { + _, err := os.Lstat(path) + return err == nil +} + +func readArg(name string) string { + path := "/cas/args/" + name + caos("get", path) + return strings.TrimSpace(string(w.Check(os.ReadFile(path)))) +} + +func main() { + w.Main(func() { + if exists("/cas/args/result") { + caos("forward", "/cas/args/result", "/cas/out") + return + } + + attempt := 0 + if exists("/cas/args/attempt") { + attempt = w.Check(strconv.Atoi(readArg("attempt"))) + } + if exists("/cas/args/error") { + attempt++ + w.True(attempt < maxAttempts, "still losing the race after %d attempts: %s", maxAttempts, readArg("error")) + } + + msg := "/cas/args/msg" + if !exists(msg) { + caos("get", "/cas/args/in") + msg = "/cas/args/in" + } + stateRef, salt := readArg("state-ref"), readArg("test-salt") + nonce := fmt.Sprintf("%s-%d-%s", caos("hash", msg), attempt, salt) + + inner := caos("curry", "--base:@=/cas/args/base", "--worker1:@=/cas/args/kv") + request := caos("prepare-request", "--base:@=/cas/args/actor", "--state-ref="+stateRef, + "--inner:hash="+inner, "--nonce="+nonce, "--message:@="+msg) + callback := caos("curry", "--base:@=/cas/args/base", "--worker1:@=/cas/args/worker1", + "--actor:@=/cas/args/actor", "--kv:@=/cas/args/kv", + "--state-ref="+stateRef, "--test-salt:@=/cas/args/test-salt", + "--attempt="+strconv.Itoa(attempt), "--msg:@="+msg) + caos("run-request-then", request, "--then:hash="+callback, "--catch") + }) +} diff --git a/tests/actor/probe.go b/tests/actor/probe.go new file mode 100644 index 00000000..bda6505a --- /dev/null +++ b/tests/actor/probe.go @@ -0,0 +1,111 @@ +// An IMPURE inner, for tests only: it does what kv.go does, after a side effect +// on the server that the test then observes. Real inners must be pure. +// +// --race-ref=R if R does not exist yet, push a competing commit to it, so the +// wrapper's leased push (which observed no head) loses the race +// --count-ref=C push one new commit to C per execution, so the number of +// commits on C is the number of times this inner actually ran +package main + +import ( + "bytes" + "fmt" + "math/rand" + "os" + "os/exec" + "path/filepath" + "strings" + "time" + + "caos/w" +) + +func exists(path string) bool { + _, err := os.Lstat(path) + return err == nil +} + +func optArg(name string) string { + path := "/cas/args/" + name + if !exists(path) { + return "" + } + cmd := exec.Command("caos", "get", path) + w.True(cmd.Run() == nil, "reading --%s", name) + return strings.TrimSpace(string(w.Check(os.ReadFile(path)))) +} + +// git runs git in the scratch repository with stdin and returns trimmed stdout. +func git(stdin string, args ...string) string { + cmd := exec.Command("git", append([]string{"-C", "/tmp/probe"}, args...)...) + cmd.Env = append(os.Environ(), "GIT_TERMINAL_PROMPT=0") + cmd.Stdin = strings.NewReader(stdin) + var stderr bytes.Buffer + cmd.Stderr = &stderr + out, err := cmd.Output() + w.True(err == nil, "git %s: %v: %s", strings.Join(args, " "), err, strings.TrimSpace(stderr.String())) + return strings.TrimSpace(string(out)) +} + +func headOf(ref string) string { + out := git("", "ls-remote", "--refs", "caos", ref) + if out == "" { + return "" + } + return strings.Fields(out)[0] +} + +func main() { + w.Main(func() { + url := strings.TrimRight(os.Getenv("CAOS_SERVER_URL"), "/") + w.True(url != "", "needs CAOS_SERVER_URL from the runner") + raceRef, countRef := optArg("race-ref"), optArg("count-ref") + + w.Must(os.RemoveAll("/tmp/probe")) + w.Must(os.MkdirAll("/tmp/probe", 0o755)) + git("", "init", "-q", ".") + git("", "config", "user.email", "probe@caos") + git("", "config", "user.name", "probe") + git("", "config", "gc.auto", "0") + git("", "remote", "add", "caos", url) + + if raceRef != "" && headOf(raceRef) == "" { + blob := git("0\n", "hash-object", "-w", "--stdin") + sub := git(fmt.Sprintf("100644 blob %s\tx\n", blob), "mktree") + root := git(fmt.Sprintf("040000 tree %s\tstate\n", sub), "mktree") + winner := git("", "commit-tree", root, "-m", "competing writer") + git("", "push", "-q", "--force-with-lease="+raceRef+":", "caos", winner+":"+raceRef) + } + + if countRef != "" { + empty := git("", "mktree") + prior := headOf(countRef) + msg := fmt.Sprintf("ran %d-%d", time.Now().UnixNano(), rand.Intn(32768)) + var run string + if prior != "" { + git("", "fetch", "-q", "caos", prior) + run = git("", "commit-tree", empty, "-p", prior, "-m", msg) + } else { + run = git("", "commit-tree", empty, "-m", msg) + } + git("", "push", "-q", "--force-with-lease="+countRef+":"+prior, "caos", run+":"+countRef) + } + + // Then behave as kv: build the --kv program as a command inside the + // prelude module (this worker runs there, in /tmp/run) and run it. + w.Must(execCmd("", "caos", "get", "/cas/args/kv")) + dir := "/tmp/run/kvcmd" + w.Must(os.RemoveAll(dir)) + w.Must(os.MkdirAll(dir, 0o755)) + w.Must(os.WriteFile(filepath.Join(dir, "main.go"), w.Check(os.ReadFile("/cas/args/kv")), 0o644)) + w.Must(execCmd("/tmp/run", "go", "run", "./kvcmd")) + }) +} + +// execCmd runs a command with this worker's stdio, in dir ("" for the current one). +func execCmd(dir, name string, args ...string) error { + cmd := exec.Command(name, args...) + cmd.Dir = dir + cmd.Stdout, cmd.Stderr = os.Stdout, os.Stderr + return cmd.Run() +} diff --git a/tests/actor/worker.go b/tests/actor/worker.go new file mode 100644 index 00000000..47a5c5f4 --- /dev/null +++ b/tests/actor/worker.go @@ -0,0 +1,273 @@ +// tests/actor: the actor wrapper + a reference key-value inner, in stages (a +// worker cannot block on a run, so each stage tail-calls the next with +// run-request-then). Every program here runs on std/go, which has git. +// +// start put a=1 -> one commit, state/a == 1 +// after-put get a -> reply 1, head unchanged (a read commits nothing) +// after-get put a=1 again -> head unchanged (same state, no commit, no push). +// This is also the crash-after-push case: a retry +// re-applies the message and reaches the same head. +// after-idem put b=2 -> a second commit whose parent is the first +// after-b fresh branch, impure inner pushes a competing commit mid-request +// raced the request FAILED (lost race, not cached); the competing head stands; +// send the identical request again +// retried the retry succeeded on top of the winner: state has x (winner) and a +// after-conc 12 concurrent puts (map-then, retried on a lost race) all landed +// after-lazy a read touched one entry and the inner saw the others unmaterialized +// after-hit1/after-hit2 +// the same read twice with different nonces ran the inner once +package main + +import ( + "bytes" + "fmt" + "math/rand" + "os" + "os/exec" + "strings" + "time" + + "caos/w" +) + +var ( + salt string + stateRef string + url string +) + +// try runs a command, returning its trimmed stdout and whether it succeeded. +func try(dir, name string, args ...string) (string, bool) { + cmd := exec.Command(name, args...) + cmd.Dir = dir + cmd.Env = append(os.Environ(), "GIT_TERMINAL_PROMPT=0") + var stderr bytes.Buffer + cmd.Stderr = &stderr + out, err := cmd.Output() + return strings.TrimSpace(string(out)), err == nil +} + +func run(name string, args ...string) string { + cmd := exec.Command(name, args...) + cmd.Env = append(os.Environ(), "GIT_TERMINAL_PROMPT=0") + var stderr bytes.Buffer + cmd.Stderr = &stderr + out, err := cmd.Output() + w.True(err == nil, "%s %s: %v: %s", name, strings.Join(args, " "), err, strings.TrimSpace(stderr.String())) + return strings.TrimSpace(string(out)) +} + +func caos(args ...string) string { return run("caos", args...) } + +func git(args ...string) string { return run("git", append([]string{"-C", "/tmp/repo"}, args...)...) } + +func exists(path string) bool { + _, err := os.Lstat(path) + return err == nil +} + +func readArg(name string) string { + path := "/cas/args/" + name + w.True(exists(path), "reading --%s", name) + caos("get", path) + return strings.TrimSpace(string(w.Check(os.ReadFile(path)))) +} + +// remoteHead is the branch's head on the server, or "". +func remoteHead(ref string) string { + out, ok := try("", "git", "ls-remote", "--refs", url, ref) + w.True(ok, "ls-remote %s", ref) + if out == "" { + return "" + } + return strings.Fields(out)[0] +} + +func fetch(oid string) { + git("fetch", "-q", "caos", oid) +} + +func stateFile(commit, name string) string { + fetch(commit) + return git("show", commit+":state/"+name) +} + +// next is the ArgTree of the following stage; extra are more --name=value args. +func next(stage string, extra ...string) string { + args := []string{"curry", "--base:@=/cas/args/base", "--worker1:@=/cas/args/worker1", + "--stage=" + stage, "--test-salt:@=/cas/args/test-salt", + "--actor:@=/cas/args/actor", "--kv:@=/cas/args/kv", + "--probe:@=/cas/args/probe", "--mapper:@=/cas/args/mapper", + "--state-ref=" + stateRef} + return caos(append(args, extra...)...) +} + +func kvInner() string { + return caos("curry", "--base:@=/cas/args/base", "--worker1:@=/cas/args/kv") +} + +// probeInner is the impure inner, in this image: --race-ref=R, --count-ref=C. +func probeInner(opts ...string) string { + return caos(append([]string{"curry", "--base:@=/cas/args/base", + "--worker1:@=/cas/args/probe", "--kv:@=/cas/args/kv"}, opts...)...) +} + +// actorRequest is the complete request for one message. +func actorRequest(message, nonce, inner string) string { + w.Must(os.WriteFile("/tmp/msg", []byte(message+"\n"), 0o644)) + _ = os.Remove("/cas/msg") + caos("put", "/tmp/msg", "/cas/msg") + return caos("prepare-request", "--base:@=/cas/args/actor", "--state-ref="+stateRef, + "--inner:hash="+inner, "--nonce="+nonce+"-"+salt, "--message:@=/cas/msg") +} + +// call sends message and continues at nextStage. With catch, a failed request +// reaches the next stage as --error instead of failing the test. +func call(message, nonce, inner, nextStage string, catch bool, extra ...string) { + request := actorRequest(message, nonce, inner) + args := []string{"run-request-then", request, "--then:hash=" + next(nextStage, extra...)} + if catch { + args = append(args, "--catch") + } + caos(args...) +} + +func freshRef(tag string) string { + return fmt.Sprintf("refs/heads/actors/test-%s-%d-%d-%d", tag, time.Now().UnixNano(), os.Getpid(), rand.Intn(32768)) +} + +func main() { + w.Main(func() { + stage := "start" + if exists("/cas/args/stage") { + stage = readArg("stage") + } + salt = readArg("test-salt") + url = strings.TrimRight(os.Getenv("CAOS_SERVER_URL"), "/") + w.True(url != "", "this test needs CAOS_SERVER_URL from the runner") + + w.Must(os.RemoveAll("/tmp/repo")) + run("git", "init", "-q", "/tmp/repo") + git("config", "user.email", "test@caos") + git("config", "user.name", "caos") + git("config", "gc.auto", "0") + git("remote", "add", "caos", url) + + if stage == "start" { + stateRef = freshRef("main") + } else { + stateRef = readArg("state-ref") + } + + switch stage { + case "start": + w.True(remoteHead(stateRef) == "", "fresh ref already exists") + call("put a 1", "n1", kvInner(), "after-put", false) + + case "after-put": + h1 := remoteHead(stateRef) + w.True(h1 != "", "put created no branch") + w.True(stateFile(h1, "a") == "1", "state/a is not 1") + w.True(git("rev-list", "--count", h1) == "1", "first update is not a root commit") + call("get a", "n2", kvInner(), "after-get", false, "--h1="+h1) + + case "after-get": + h1 := readArg("h1") + w.True(remoteHead(stateRef) == h1, "a read changed the head") + reply := readArg("result") + w.True(reply == "1", "get a replied '%s'", reply) + call("put a 1", "n3", kvInner(), "after-idem", false, "--h1="+h1) + + case "after-idem": + h1 := readArg("h1") + w.True(remoteHead(stateRef) == h1, "an unchanged put made a commit") + call("put b 2", "n4", kvInner(), "after-b", false, "--h1="+h1) + + case "after-b": + h1 := readArg("h1") + h2 := remoteHead(stateRef) + w.True(h2 != "", "branch vanished") + w.True(h2 != h1, "put b made no commit") + fetch(h2) + w.True(git("rev-parse", h2+"^1") == h1, "second update is not on the first") + w.True(git("rev-list", "--count", h2) == "2", "history is not a linear chain of two") + w.True(stateFile(h2, "a") == "1", "state/a lost") + w.True(stateFile(h2, "b") == "2", "state/b missing") + // A forced lost race: on a fresh branch the impure inner pushes a + // competing commit while the request is in flight, so the wrapper's + // lease (no head) fails. + stateRef = freshRef("race") + call("put a 1", "n5", probeInner("--race-ref="+stateRef), "raced", true) + + case "raced": + w.True(exists("/cas/args/error"), "the raced request did not fail (--error missing)") + winner := remoteHead(stateRef) + w.True(winner != "", "the competing writer left no branch") + w.True(stateFile(winner, "x") == "0", "the head is not the competing commit") + _, has := try("/tmp/repo", "git", "cat-file", "-e", winner+":state/a") + w.True(!has, "the lost request published anyway") + // The identical request again (same nonce): a cached failure would + // replay the failure; instead it re-runs against the new head. + call("put a 1", "n5", probeInner("--race-ref="+stateRef), "retried", false, "--winner="+winner) + + case "retried": + winner := readArg("winner") + head := remoteHead(stateRef) + w.True(head != "", "branch vanished") + w.True(head != winner, "the retry published nothing") + fetch(head) + w.True(git("rev-parse", head+"^1") == winner, "the retry is not on top of the winner") + w.True(stateFile(head, "x") == "0", "the winner's entry was lost") + w.True(stateFile(head, "a") == "1", "state/a missing after the retry") + // Concurrent writers, each retrying a lost race, must converge with + // no lost update. + stateRef = freshRef("conc") + w.Must(os.RemoveAll("/tmp/msgs")) + w.Must(os.MkdirAll("/tmp/msgs", 0o755)) + for n := 1; n <= 12; n++ { + w.Must(os.WriteFile(fmt.Sprintf("/tmp/msgs/m%d", n), []byte(fmt.Sprintf("put c%d v%d\n", n, n)), 0o644)) + } + caos("put", "/tmp/msgs", "/cas/msgs") + mapper := caos("curry", "--base:@=/cas/args/base", "--worker1:@=/cas/args/mapper", + "--actor:@=/cas/args/actor", "--kv:@=/cas/args/kv", + "--state-ref="+stateRef, "--test-salt:@=/cas/args/test-salt") + caos("map-then", "/cas/msgs", "--map:hash="+mapper, "--then:hash="+next("after-conc")) + + case "after-conc": + head := remoteHead(stateRef) + w.True(head != "", "no branch after the concurrent puts") + for n := 1; n <= 12; n++ { + w.True(stateFile(head, fmt.Sprintf("c%d", n)) == fmt.Sprintf("v%d", n), "update c%d was lost", n) + } + w.True(git("rev-list", "--count", head) == "12", "expected a linear chain of 12 commits") + call("getcheck c3", "n6", kvInner(), "after-lazy", false, "--h="+head) + + case "after-lazy": + head := readArg("h") + w.True(remoteHead(stateRef) == head, "a read changed the head") + reply := readArg("result") + w.True(reply == "v3", "getcheck replied '%s'", reply) + // The same read twice, different nonces: the inner (pure, so + // cached) runs once. + countRef := fmt.Sprintf("refs/heads/actors-count/%d-%d-%d", time.Now().UnixNano(), os.Getpid(), rand.Intn(32768)) + call("get c4", "n7", probeInner("--count-ref="+countRef), "after-hit1", false, "--count-ref="+countRef) + + case "after-hit1": + countRef := readArg("count-ref") + w.True(remoteHead(countRef) != "", "the inner did not run") + call("get c4", "n8", probeInner("--count-ref="+countRef), "after-hit2", false, "--count-ref="+countRef) + + case "after-hit2": + countRef := readArg("count-ref") + last := remoteHead(countRef) + w.True(last != "", "count ref vanished") + fetch(last) + runs := git("rev-list", "--count", last) + w.True(runs == "1", "the inner ran %s times; the repeat should hit the cache", runs) + w.Report("actor: ALL PASS\n") + + default: + w.True(false, "unknown --stage: %s", stage) + } + }) +}