Remove experiments/ - concept proven

This commit is contained in:
Levi Neely 2026-08-04 16:39:15 +02:00
parent 593584a673
commit 6b3913a524
3 changed files with 0 additions and 487 deletions

View File

@ -1,5 +0,0 @@
module ollie/experiments/9p-stream
go 1.25
require 9fans.net/go v0.0.7

View File

@ -1,34 +0,0 @@
9fans.net/go v0.0.7 h1:H5CsYJTf99C8EYAQr+uSoEJnLP/iZU8RmDuhyk30iSM=
9fans.net/go v0.0.7/go.mod h1:Rxvbbc1e+1TyGMjAvLthGTyO97t+6JMQ6ly+Lcs9Uf0=
dmitri.shuralyov.com/gpu/mtl v0.0.0-20201218220906-28db891af037/go.mod h1:H6x//7gZCb22OMCxBHrMx7a5I7Hp++hsVxbQ4BYO7hU=
github.com/BurntSushi/xgb v0.0.0-20160522181843-27f122750802/go.mod h1:IVnqGOEym/WlBOVXweHU+Q+/VP0lqqI8lqeDx9IjBqo=
github.com/go-gl/glfw/v3.3/glfw v0.0.0-20200222043503-6f7a984d4dc4/go.mod h1:tQ2UAYgL5IevRw8kRxooKSPJfGvJ9fJQFa0TUsXzTg8=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/exp v0.0.0-20190731235908-ec7cb31e5a56/go.mod h1:JhuoJpWY28nO4Vef9tZUw9qufEGTyX1+7lmHxV5q5G4=
golang.org/x/exp v0.0.0-20210405174845-4513512abef3/go.mod h1:I6l2HNBLBZEcrOoCpyKLdY2lHoRZ8lI4x60KMCQDft4=
golang.org/x/image v0.0.0-20190227222117-0694c2d4d067/go.mod h1:kZ7UVZpmo3dzQBMxlp+ypCbDeSB+sBbTgSJuh5dn5js=
golang.org/x/image v0.0.0-20190802002840-cff245a6509b/go.mod h1:FeLwcggjj3mMvU+oOTbSwawSJRM1uh48EjtB4UJZlP0=
golang.org/x/mobile v0.0.0-20190312151609-d3739f865fa6/go.mod h1:z+o9i4GpDbdi3rU15maQ/Ox0txvL9dWGYEHz965HBQE=
golang.org/x/mobile v0.0.0-20201217150744-e6ae53a27f4f/go.mod h1:skQtrUTUwhdJvXM/2KKJzY8pDgNr9I/FOMqDVRPBUS4=
golang.org/x/mobile v0.0.0-20210220033013-bdb1ca9a1e08/go.mod h1:skQtrUTUwhdJvXM/2KKJzY8pDgNr9I/FOMqDVRPBUS4=
golang.org/x/mod v0.1.0/go.mod h1:0QHyrYULN0/3qlju5TqG8bIK38QM8yzMo5ekMj3DlcY=
golang.org/x/mod v0.1.1-0.20191105210325-c90efee705ee/go.mod h1:QqPTAvyqsEbceGzBzNggFXnrqF1CaUcvgkdR5Ot7KZg=
golang.org/x/mod v0.1.1-0.20191209134235-331c550502dd/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
golang.org/x/mod v0.3.1-0.20200828183125-ce943fd02449/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210415045647-66c3f260301c/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/tools v0.0.0-20190312151545-0bb0c0a6e846/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs=
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.0.0-20200117012304-6edc0a871e69/go.mod h1:TB2adYChydJhpapKDTa4BR/hXlZSLoq2Wpct/0txZ28=
golang.org/x/tools v0.0.0-20200207183749-b753a1ba74fa/go.mod h1:TB2adYChydJhpapKDTa4BR/hXlZSLoq2Wpct/0txZ28=
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=

View File

@ -1,448 +0,0 @@
// 9P streaming chat prototype.
//
// Demonstrates native 9P streaming: reads on `stream` block until
// new tokens arrive from the simulated LLM, then return them.
// No polling needed — the protocol handles the wait naturally.
//
// Files:
// ctl - write "prompt <text>" to start generation
// stream - blocking read: delivers tokens as they arrive, EOF when done
// chat - full conversation log (non-blocking)
//
// Usage:
// go run .
// echo "prompt hello" | 9p -a localhost:5640 write ctl
// 9p -a localhost:5640 read stream
// 9p -a localhost:5640 read chat
package main
import (
"context"
"fmt"
"log"
"math/rand"
"net"
"os"
"strings"
"sync"
"time"
"9fans.net/go/plan9"
)
const listenAddr = "localhost:5640"
// Simulated response — what the "LLM" produces
var cannedResponse = `I'll explain how 9P streaming works.
In Plan 9, the file server controls when it responds to read requests.
This gives us natural streaming semantics without polling:
1. Client opens the stream file
2. Client sends Tread
3. Server holds the request until tokens arrive
4. Server responds with new tokens
5. Repeat until response is complete (EOF)
The transport handles the wait — no timers, no busy loops.
` + "```go" + `
// The key: server delays Rread until data exists
func handleRead(req) {
wait until tokens available
respond with tokens
}
` + "```" + `
This is how /dev/cons, /dev/mouse, and other Plan 9 devices work.
Streaming is not a feature — it's the natural mode of operation.
`
func main() {
srv := &chatServer{}
srv.streamCond = sync.NewCond(&srv.mu)
ln, err := net.Listen("tcp", listenAddr)
if err != nil {
log.Fatal(err)
}
log.Printf("9P stream prototype on %s", listenAddr)
log.Printf(" echo 'prompt hello' | 9p -a %s write ctl", listenAddr)
log.Printf(" 9p -a %s read stream", listenAddr)
for {
conn, err := ln.Accept()
if err != nil {
log.Printf("accept: %v", err)
continue
}
go srv.serve(conn)
}
}
type chatServer struct {
mu sync.Mutex
streamCond *sync.Cond
// Current streaming response
streamBuf []byte // tokens produced so far
streamDone bool // generation complete
// Full conversation log
chat []byte
}
type fidState struct {
path string
offset int64
}
func (s *chatServer) serve(conn net.Conn) {
defer conn.Close()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
var fidsMu sync.Mutex
fids := make(map[uint32]*fidState)
var msize uint32 = 8192
responses := make(chan *plan9.Fcall, 16)
// Writer
go func() {
for resp := range responses {
if err := plan9.WriteFcall(conn, resp); err != nil {
cancel()
return
}
}
}()
var wg sync.WaitGroup
for {
fc, err := plan9.ReadFcall(conn)
if err != nil {
break
}
wg.Add(1)
go func(fc *plan9.Fcall) {
defer wg.Done()
resp := s.handle(ctx, fc, fids, &fidsMu, &msize)
select {
case responses <- resp:
case <-ctx.Done():
}
}(fc)
}
cancel()
// Wake any blocked stream readers
s.streamCond.Broadcast()
wg.Wait()
close(responses)
}
func (s *chatServer) handle(ctx context.Context, fc *plan9.Fcall, fids map[uint32]*fidState, fidsMu *sync.Mutex, msize *uint32) *plan9.Fcall {
switch fc.Type {
case plan9.Tversion:
if fc.Msize < *msize {
*msize = fc.Msize
}
return &plan9.Fcall{Type: plan9.Rversion, Tag: fc.Tag, Msize: *msize, Version: "9P2000"}
case plan9.Tauth:
return &plan9.Fcall{Type: plan9.Rerror, Tag: fc.Tag, Ename: "no auth required"}
case plan9.Tattach:
fidsMu.Lock()
fids[fc.Fid] = &fidState{path: "/"}
fidsMu.Unlock()
return &plan9.Fcall{Type: plan9.Rattach, Tag: fc.Tag, Qid: dirQid("/")}
case plan9.Twalk:
fidsMu.Lock()
src, ok := fids[fc.Fid]
fidsMu.Unlock()
if !ok {
return errResp(fc, "bad fid")
}
if len(fc.Wname) == 0 {
fidsMu.Lock()
fids[fc.Newfid] = &fidState{path: src.path}
fidsMu.Unlock()
return &plan9.Fcall{Type: plan9.Rwalk, Tag: fc.Tag}
}
path := src.path
var qids []plan9.Qid
for _, name := range fc.Wname {
switch name {
case "ctl", "stream", "chat":
path = "/" + name
qids = append(qids, fileQid(name))
default:
return errResp(fc, "not found: "+name)
}
}
fidsMu.Lock()
fids[fc.Newfid] = &fidState{path: path}
fidsMu.Unlock()
return &plan9.Fcall{Type: plan9.Rwalk, Tag: fc.Tag, Wqid: qids}
case plan9.Topen:
fidsMu.Lock()
f, ok := fids[fc.Fid]
fidsMu.Unlock()
if !ok {
return errResp(fc, "bad fid")
}
f.offset = 0
if f.path == "/stream" {
s.mu.Lock()
f.offset = int64(len(s.streamBuf))
s.mu.Unlock()
}
return &plan9.Fcall{Type: plan9.Ropen, Tag: fc.Tag, Qid: fileQid(strings.TrimPrefix(f.path, "/"))}
case plan9.Tread:
fidsMu.Lock()
f, ok := fids[fc.Fid]
fidsMu.Unlock()
if !ok {
return errResp(fc, "bad fid")
}
count := int(fc.Count)
if uint32(count) > *msize-plan9.IOHDRSZ {
count = int(*msize - plan9.IOHDRSZ)
}
switch f.path {
case "/":
return s.readDir(fc, f, count)
case "/chat":
return s.readChat(fc, f, count)
case "/stream":
return s.readStream(ctx, fc, f, count)
default:
return errResp(fc, "cannot read "+f.path)
}
case plan9.Twrite:
fidsMu.Lock()
f, ok := fids[fc.Fid]
fidsMu.Unlock()
if !ok {
return errResp(fc, "bad fid")
}
if f.path == "/ctl" {
cmd := strings.TrimSpace(string(fc.Data))
if strings.HasPrefix(cmd, "prompt ") {
prompt := strings.TrimPrefix(cmd, "prompt ")
go s.generateResponse(prompt)
return &plan9.Fcall{Type: plan9.Rwrite, Tag: fc.Tag, Count: uint32(len(fc.Data))}
}
return errResp(fc, "unknown command: "+cmd)
}
return errResp(fc, "cannot write "+f.path)
case plan9.Tstat:
fidsMu.Lock()
f, ok := fids[fc.Fid]
fidsMu.Unlock()
if !ok {
return errResp(fc, "bad fid")
}
dir := s.stat(f.path)
b, _ := dir.Bytes()
return &plan9.Fcall{Type: plan9.Rstat, Tag: fc.Tag, Stat: b}
case plan9.Tclunk:
fidsMu.Lock()
delete(fids, fc.Fid)
fidsMu.Unlock()
return &plan9.Fcall{Type: plan9.Rclunk, Tag: fc.Tag}
default:
return errResp(fc, fmt.Sprintf("unhandled %d", fc.Type))
}
}
// readStream blocks until new data is available or response is done.
func (s *chatServer) readStream(ctx context.Context, fc *plan9.Fcall, f *fidState, count int) *plan9.Fcall {
s.mu.Lock()
// Wait until there's data beyond our read position, or generation is done
for int64(len(s.streamBuf)) <= f.offset && !s.streamDone {
// Release lock and wait for signal
s.streamCond.Wait()
// Check if context cancelled (connection closed)
if ctx.Err() != nil {
s.mu.Unlock()
return errResp(fc, "interrupted")
}
}
// EOF: no more data and generation is done
if f.offset >= int64(len(s.streamBuf)) {
s.mu.Unlock()
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
}
// Return available data
end := int64(len(s.streamBuf))
if end-f.offset > int64(count) {
end = f.offset + int64(count)
}
data := make([]byte, end-f.offset)
copy(data, s.streamBuf[f.offset:end])
f.offset = end
s.mu.Unlock()
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(data)), Data: data}
}
// readChat returns the full log (non-blocking).
func (s *chatServer) readChat(fc *plan9.Fcall, f *fidState, count int) *plan9.Fcall {
s.mu.Lock()
data := s.chat
s.mu.Unlock()
offset := int64(fc.Offset)
if offset >= int64(len(data)) {
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
}
end := offset + int64(count)
if end > int64(len(data)) {
end = int64(len(data))
}
chunk := make([]byte, end-offset)
copy(chunk, data[offset:end])
f.offset = end
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(chunk)), Data: chunk}
}
func (s *chatServer) readDir(fc *plan9.Fcall, f *fidState, count int) *plan9.Fcall {
entries := []plan9.Dir{
{Name: "ctl", Qid: fileQid("ctl"), Mode: 0222, Atime: uint32(time.Now().Unix()), Mtime: uint32(time.Now().Unix())},
{Name: "stream", Qid: fileQid("stream"), Mode: 0444, Atime: uint32(time.Now().Unix()), Mtime: uint32(time.Now().Unix())},
{Name: "chat", Qid: fileQid("chat"), Mode: 0444, Atime: uint32(time.Now().Unix()), Mtime: uint32(time.Now().Unix())},
}
var buf []byte
for _, d := range entries {
b, _ := d.Bytes()
buf = append(buf, b...)
}
offset := int64(fc.Offset)
if offset >= int64(len(buf)) {
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
}
end := offset + int64(count)
if end > int64(len(buf)) {
end = int64(len(buf))
}
chunk := buf[offset:end]
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(chunk)), Data: chunk}
}
func (s *chatServer) stat(path string) plan9.Dir {
now := uint32(time.Now().Unix())
u := os.Getenv("USER")
switch path {
case "/":
return plan9.Dir{Name: "/", Qid: dirQid("/"), Mode: plan9.DMDIR | 0555, Atime: now, Mtime: now, Uid: u, Gid: u}
case "/ctl":
return plan9.Dir{Name: "ctl", Qid: fileQid("ctl"), Mode: 0222, Atime: now, Mtime: now, Uid: u, Gid: u}
case "/stream":
s.mu.Lock()
l := uint64(len(s.streamBuf))
s.mu.Unlock()
return plan9.Dir{Name: "stream", Qid: fileQid("stream"), Mode: 0444, Length: l, Atime: now, Mtime: now, Uid: u, Gid: u}
case "/chat":
s.mu.Lock()
l := uint64(len(s.chat))
s.mu.Unlock()
return plan9.Dir{Name: "chat", Qid: fileQid("chat"), Mode: 0444, Length: l, Atime: now, Mtime: now, Uid: u, Gid: u}
}
return plan9.Dir{}
}
// generateResponse simulates streaming LLM token generation with network jitter.
func (s *chatServer) generateResponse(prompt string) {
// Reset stream state
s.mu.Lock()
s.streamBuf = nil
s.streamDone = false
s.chat = append(s.chat, []byte("\n[user]\n"+prompt+"\n\n[assistant]\n")...)
s.mu.Unlock()
s.streamCond.Broadcast()
// Tokenize into small chunks (simulates token batching)
tokens := tokenize(cannedResponse)
for _, tok := range tokens {
// Simulate network latency: 30-150ms per batch
delay := time.Duration(30+rand.Intn(120)) * time.Millisecond
// 10% chance of longer pause (network jitter)
if rand.Intn(10) == 0 {
delay += time.Duration(200+rand.Intn(400)) * time.Millisecond
}
time.Sleep(delay)
s.mu.Lock()
s.streamBuf = append(s.streamBuf, []byte(tok)...)
s.chat = append(s.chat, []byte(tok)...)
s.mu.Unlock()
s.streamCond.Broadcast()
}
// Mark complete
s.mu.Lock()
s.streamDone = true
s.mu.Unlock()
s.streamCond.Broadcast()
log.Printf("response complete (%d bytes)", len(s.streamBuf))
}
func tokenize(text string) []string {
// Split into chunks of ~20-60 bytes each (simulates token batching from LLM)
var tokens []string
for len(text) > 0 {
n := 20 + rand.Intn(40)
if n > len(text) {
n = len(text)
}
// Try to break at whitespace
if n < len(text) {
for i := n; i > n-10 && i > 0; i-- {
if text[i] == ' ' || text[i] == '\n' {
n = i + 1
break
}
}
}
tokens = append(tokens, text[:n])
text = text[n:]
}
return tokens
}
func dirQid(name string) plan9.Qid {
return plan9.Qid{Type: plan9.QTDIR, Path: hash(name)}
}
func fileQid(name string) plan9.Qid {
return plan9.Qid{Type: plan9.QTFILE, Path: hash(name)}
}
func hash(s string) uint64 {
h := uint64(0xcbf29ce484222325)
for _, c := range s {
h ^= uint64(c)
h *= 0x100000001b3
}
return h
}
func errResp(fc *plan9.Fcall, msg string) *plan9.Fcall {
return &plan9.Fcall{Type: plan9.Rerror, Tag: fc.Tag, Ename: msg}
}