Repository navigation
Expand file tree
/
Copy pathsession.js
More file actions
274 lines (246 loc) · 14.4 KB
/
Copy pathsession.js
File metadata and controls
274 lines (246 loc) · 14.4 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
// The auto-okf FLAGSHIP example: FOUR AGENT WRITERS co-author a knowledge
// bundle on ONE shared multi-writer vault while PARTITIONED, then heal — and
// a human operator dumps the session as a readable OKF v0.1 markdown bundle
// (the committed out/ tree).
//
// node okf/examples/agent-session/session.js (from the repo root, after pnpm install)
//
// ThreadMode -> DocumentMode: the agents emit a THREAD (an append-only op
// log of findings, determinations, and an RFC); out/ is the refactored
// DOCUMENT (the bundle a human reads at the end). The op log stays the source
// of truth; out/ is a face regenerated from it, byte-for-byte (I9).
//
// The cast (each agent is its own vault peer with its own writer core):
// alpha — session coordinator (vault creator, genesis indexer)
// beta — investigator
// gamma — investigator
// delta — reviewer
//
// What it demonstrates, in order:
// 1. governance — alpha admits beta/gamma/delta as writers (logged §8 ops)
// 2. authorship — findings, determinations, and an RFC are ordinary
// concepts in the same log, cross-linked per OKF §5
// 3. partition — the four peers DISCONNECT and write concurrently:
// * beta & gamma both edit one shared concept — the tags
// OR-set UNIONS and the LWW `status` register picks one
// deterministic winner and FLAGS the overwrite (§5)
// * beta & gamma write the RFC body from the same base:
// the system keeps BOTH revisions and flags the
// conflict (§5); alpha resolves it with a merge op
// * one link points at an unwritten concept — it shows up
// in wanted/, the ranked demand index, never a flag (§9)
// 4. heal + dump — the peers reconnect, converge (I2), and alpha
// materializes the OKF bundle to out/ (concepts + per-dir
// index.md + root log.md)
//
// Determinism: writer key pairs are derived from fixed seeds and every
// concurrent phase starts from a fully converged barrier, so the final view —
// and therefore every byte under out/, including the `_id` frontmatter, the
// log.md dates (from each op's informational `at`), and the wanted ranking —
// is identical on every run (I9; asserted by okf/test/example-agent-session.test.js).
// Do not hand-edit out/: it is regenerated by this script.
import { createHash } from 'node:crypto'
import { mkdir } from 'node:fs/promises'
import { join, dirname, resolve, relative } from 'node:path'
import { fileURLToPath } from 'node:url'
import crypto from 'hypercore-crypto'
import b4a from 'b4a'
import { initVault, Vault } from '@auto-okf/store'
import { materialize, wanted } from '@auto-okf/faces'
import { keyspace as K, revHash } from '@auto-okf/core'
const HERE = dirname(fileURLToPath(import.meta.url))
const REPO = join(HERE, '..', '..', '..')
// Committed, human-readable output (override for tests).
const OUT = process.env.OKF_AGENT_SESSION_OUT ? resolve(process.env.OKF_AGENT_SESSION_OUT) : join(HERE, 'out')
// Per-run vault storage scratch (unique, never reused, never cleaned up).
const STORAGE = process.env.OKF_AGENT_SESSION_STORAGE
? resolve(process.env.OKF_AGENT_SESSION_STORAGE)
: join(REPO, 'sslop', '035', 'tmp', 'agent-session', `run-${process.pid}-${Date.now()}`)
const seed = (name) => createHash('sha256').update(`auto-okf/agent-session/${name}`).digest()
const say = (who, msg) => console.log(`[${who.padEnd(5)}] ${msg}`)
const sleep = (ms) => new Promise((r) => setTimeout(r, ms))
// §3: `at` is informational (never an input to merge) — it is the sole source
// for log.md date headings. Pinning it makes the dump byte-deterministic. The
// session spans two UTC days so log.md shows two `## YYYY-MM-DD` headings.
const DAY1 = Date.UTC(2026, 6, 2) // 2026-07-02: session setup + seed concepts
const DAY2 = Date.UTC(2026, 6, 3) // 2026-07-03: the partitioned work + heal
let ticks = 0
const at1 = () => ({ at: DAY1 + ticks++ * 1000 })
const at2 = () => ({ at: DAY2 + ticks++ * 1000 })
async function until(fn, ms = 30000, step = 50) {
const t0 = Date.now()
for (;;) {
if (await fn()) return true
if (Date.now() - t0 >= ms) return false
await sleep(step)
}
}
/** Direct replication stream pair; returns an unpair(). */
function pairVaults(a, b) {
const s1 = a.replicate(true)
const s2 = b.replicate(false)
s1.on('error', () => {})
s2.on('error', () => {})
s1.pipe(s2).pipe(s1)
return () => {
s1.destroy()
s2.destroy()
}
}
/** Full-mesh replication among all peers; returns a teardown(). */
function mesh(vaults) {
const unpairs = []
for (let i = 0; i < vaults.length; i++) {
for (let j = i + 1; j < vaults.length; j++) unpairs.push(pairVaults(vaults[i], vaults[j]))
}
return () => unpairs.forEach((u) => u())
}
/** Barrier: wait until every peer holds the identical view (I2). */
async function converge(vaults) {
const ok = await until(async () => {
const hashes = []
for (const v of vaults) {
await v.update()
hashes.push(await v.snapshotHash())
}
return hashes.every((h) => h === hashes[0])
})
if (!ok) throw new Error('peers failed to converge')
}
async function count(vault, prefix) {
let n = 0
for await (const _ of vault.scan(prefix)) n++
return n
}
async function main() {
await mkdir(STORAGE, { recursive: true })
console.log('=== agent-session: four agents, one shared vault, partitioned then healed ===\n')
// ---- session setup: one vault, four writers -----------------------------
const alpha = await initVault(join(STORAGE, 'alpha'), { keyPair: crypto.keyPair(seed('alpha')) })
const beta = await Vault.join(join(STORAGE, 'beta'), alpha.key, { keyPair: crypto.keyPair(seed('beta')) })
const gamma = await Vault.join(join(STORAGE, 'gamma'), alpha.key, { keyPair: crypto.keyPair(seed('gamma')) })
const delta = await Vault.join(join(STORAGE, 'delta'), alpha.key, { keyPair: crypto.keyPair(seed('delta')) })
const agents = [alpha, beta, gamma, delta]
say('alpha', `opened the vault (key ${alpha.keyHex.slice(0, 16)}…); alpha is the genesis indexer`)
let connected = mesh(agents)
for (const [name, v] of [['beta', beta], ['gamma', gamma], ['delta', delta]]) {
await alpha.addWriter(v.localKey, { indexer: false })
say('alpha', `admitted ${name} as a writer (AddWriter op, ${b4a.toString(v.localKey, 'hex').slice(0, 16)}…)`)
}
for (const v of [beta, gamma, delta]) {
if (!(await until(async () => { await v.update(); return v.writable }))) {
throw new Error('agent never became writable')
}
}
// ---- phase 1: alpha seeds the shared thread (all peers connected) --------
say('alpha', 'opening the thread: an overview note and an RFC stub, cross-linked (OKF §5)')
const { tag: overview } = await alpha.createConcept('overview', 'reference', at1())
await alpha.setField(overview, 'title', 'Incident: Cache-Layer Regression', at1())
await alpha.setField(overview, 'description', 'Shared root note for the session.', at1())
await alpha.setBody(overview, null, 'Coordinating the cache-layer incident. See the [RFC](/rfc/caching-layer.md).', at1())
const { tag: rfc } = await alpha.createConcept('rfc/caching-layer', 'rfc', at1())
await alpha.setField(rfc, 'title', 'RFC: Caching Layer', at1())
await alpha.setBody(rfc, null, 'Draft. Problem statement pending. See the [overview](/overview.md).', at1())
await converge(agents)
// The converged RFC head is the SHARED BASE the §5 body conflict needs: two
// revisions citing the same NON-null base are true concurrency (a null base
// is "no base cited" and never conflicts).
const rfcBase = await alpha.val(K.cBody(rfc))
say('all', 'barrier: the seed thread replicated; every peer holds a bit-identical view (I2)\n')
// ---- phase 2: the peers PARTITION and write concurrently -----------------
connected() // tear down replication: everything below is truly concurrent
console.log('--- the four agents are partitioned and write in parallel ---')
// beta investigates. The [sharding plan] link points at a concept nobody
// has written — it lands in wanted/, the ranked demand index (§9), not a flag.
say('beta', 'writing "Finding: Latency Spike" (links the RFC and an unwritten sharding plan)')
const { tag: latency } = await beta.createConcept('findings/latency-spike', 'finding', at2())
await beta.setField(latency, 'title', 'Finding: Latency Spike', at2())
await beta.setBody(
latency,
null,
'p99 latency tripled after the deploy. Bears on the [RFC](/rfc/caching-layer.md); the remedy is tracked in the [sharding plan](/determinations/sharding.md).',
at2()
)
say('gamma', 'writing "Finding: Cache Misses" and "Determination: Adopt LRU"')
const { tag: misses } = await gamma.createConcept('findings/cache-misses', 'finding', at2())
await gamma.setField(misses, 'title', 'Finding: Cache Misses', at2())
await gamma.setBody(misses, null, 'Cold-cache miss rate hit 40%. Bears on the [RFC](/rfc/caching-layer.md).', at2())
const { tag: lru } = await gamma.createConcept('determinations/adopt-lru', 'determination', at2())
await gamma.setField(lru, 'title', 'Determination: Adopt LRU', at2())
await gamma.setBody(lru, null, 'We will adopt LRU eviction. Grounded in [Cache Misses](/findings/cache-misses.md).', at2())
say('delta', 'writing "Determination: TTL Policy"')
const { tag: ttl } = await delta.createConcept('determinations/ttl-policy', 'determination', at2())
await delta.setField(ttl, 'title', 'Determination: TTL Policy', at2())
await delta.setBody(ttl, null, 'Cache entries expire after 300s. Grounded in the [RFC](/rfc/caching-layer.md).', at2())
// beta & gamma BOTH edit the shared overview note, concurrently:
// - tags: two AddTag ops in the same OR-set -> UNION (§5, I5)
// - status: two LWW writes to one register -> one deterministic winner,
// the loser recorded under flag/clobber (§5, I4 — no silent loss)
say('beta', "tagging overview `reviewed`; setting status `triage` (concurrent with gamma)")
await beta.addTag(overview, 'reviewed', at2())
await beta.setField(overview, 'status', ['triage'], at2())
say('gamma', "tagging overview `urgent`; setting status `under-review` (concurrent with beta)")
await gamma.addTag(overview, 'urgent', at2())
await gamma.setField(overview, 'status', ['under-review'], at2())
// beta & gamma BOTH rewrite the RFC body from the SAME converged base:
// the system keeps both revisions and raises exactly one flag/conflict (§5).
const betaRfc = 'RFC v2 (beta): adopt LRU eviction. Backed by the [Adopt LRU](/determinations/adopt-lru.md) determination.'
const gammaRfc = 'RFC v2 (gamma): adopt a 300s TTL. Backed by the [TTL Policy](/determinations/ttl-policy.md) determination.'
say('beta', 'rewriting the RFC body — adopt LRU (cites the shared base)')
await beta.setBody(rfc, rfcBase, betaRfc, at2())
say('gamma', 'rewriting the RFC body — adopt TTL (cites the SAME base; a true conflict)')
await gamma.setBody(rfc, rfcBase, gammaRfc, at2())
// ---- phase 3: the session heals — every disagreement is visible ---------
console.log('\n--- the agents reconnect: merge, audit, converge ---')
connected = mesh(agents)
await converge(agents)
const tagsRec = await alpha.val(K.cField(overview, 'tags'))
const tags = [...new Set(Object.values((tagsRec && tagsRec.els) || {}))].sort()
say('all', `overview tags UNIONED across writers: [${tags.join(', ')}] (OR-set, I5)`)
const statusRec = await alpha.val(K.cField(overview, 'status'))
say('all', `overview status: LWW register resolved to "${statusRec.values[0]}" (deterministic linearization, §5)`)
say('all', `cross-writer overwrites are NOT silent: ${await count(alpha, K.flagPrefix('clobber'))} flag/clobber record(s) (§5, I4)`)
say('all', `body conflict detected and BOTH revisions kept: ${await count(alpha, K.flagPrefix('conflict'))} flag/conflict record(s) (§5)`)
const demand = await wanted(alpha)
say('all', `the graph WANTS unwritten concepts: ${demand.map((w) => `${w.path} (${w.count} ref)`).join(', ')} (wanted/, §9 — never a flag)`)
// ---- phase 4: alpha resolves the conflict (editorial pass) ---------------
console.log('\n--- editorial pass: alpha resolves the RFC body conflict ---')
const keep = revHash(rfcBase, betaRfc)
const supersede = revHash(rfcBase, gammaRfc)
const merged =
'RFC v2 (merged): adopt LRU eviction with a 300s TTL backstop. Backed by the [Adopt LRU](/determinations/adopt-lru.md) and [TTL Policy](/determinations/ttl-policy.md) determinations.'
say('alpha', 'ResolveConflict on the RFC: keep beta\'s revision, supersede gamma\'s, append the merged text (§5)')
await alpha.resolveConflict(rfc, keep, { supersede: [supersede], text: merged, ...at2() })
await converge(agents)
say('all', `flag/conflict drained by the ResolveConflict op: ${await count(alpha, K.flagPrefix('conflict'))} left`)
// ---- phase 5: the human operator dumps the session -----------------------
console.log('\n--- session dump: the thread refactored into a readable document ---')
const report = await materialize(alpha, OUT)
if (report.pendingIngest.length > 0 || report.unsafe.length > 0) {
throw new Error(
`out/ has hand-edited or unsafe files (${[...report.pendingIngest, ...report.unsafe.map((u) => u.rel)].join(', ')}); ` +
'out/ is generated — revert it and re-run'
)
}
const files = [...report.written, ...report.unchanged].sort()
say('alpha', `materialized ${files.length} files to ${relative(process.cwd(), OUT)}/`)
for (const f of files) console.log(` ${f}`)
console.log(`
=== what just happened ===
Four agents wrote a thread — findings, determinations, and an RFC — as ops in
one shared, partitioned log. They tagged concurrently (the OR-set unioned),
disagreed on a register (LWW picked a winner, the loser is flagged), forked the
RFC body (both revisions kept, then merged by ResolveConflict), and left one
link pointing at an unwritten concept (it is queued in wanted/, ranked by
demand). Here is the graph as a knowledge bundle a human operator can read:
the Markdown files under out/ — one per concept, plus a derived index.md per
directory and a root log.md that narrates the op log, newest day first. The log
stays the source of truth; out/ is a face regenerated from it, byte-for-byte.
`)
connected() // destroy the replication streams before closing
for (const v of agents) await v.close()
}
main().catch((err) => {
console.error(err)
process.exit(1)
})