Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions docs/en/changes/changes.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,12 @@
before it, and it now says so. It used to say the request before was not among the loaded bodies,
here and for the first call of any chain, where nothing is missing. The renderer is Horizon's, now
pinned at `88e0110`, which also draws a light theme's scrollbars and other native controls light.
- A landed LangChain file is named for the collected time in its header. The collector read the
clock once per pass for the name and once per file for the header, so a server that stores files by
session and sequence and names them from the header, as the OAP does, named them apart from asz. On
a real root, all 50 LangChain files were named apart from their header, by up to 99 milliseconds,
and none of the 7,185 Claude Code files was. Files landed before this keep their names, because
landed files are never rewritten.

## Claude Code

Expand Down Expand Up @@ -98,6 +104,11 @@
nested agents and Claude Code's synchronous returns gain it; nothing is re-landed, and the next
round of an existing conversation adds it.

- A stream's `opened_by` lists its origins in relation id order. A call the assembler could not tie
to one stream is an origin of each stream it could have started, and the origins were listed in the
order a Go map gave them, so the document changed from one read to the next. On a real
conversation whose stream had three, five reads gave three orders.

## Release

- The `pypi` skill says what a release manager sets up before an upload. PyPI requires two-factor
Expand Down
6 changes: 3 additions & 3 deletions internal/adapters/langsmith/bodies.go
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,7 @@ type landedBodies struct {
// landBodies writes one session's bodies out of one request, in files of the
// session's own sequence, under the session lock the caller holds.
func (c *Collector) landBodies(session string, state *storage.SessionState,
bodies []body, stamp string, request Waiting) (landedBodies, error) {
bodies []body, now time.Time, request Waiting) (landedBodies, error) {
var out landedBodies
if len(bodies) == 0 {
return out, nil
Expand All @@ -179,15 +179,15 @@ func (c *Collector) landBodies(session string, state *storage.SessionState,
}
seq := state.Take()
header := &sessiondata.Header{
Seq: seq, At: c.Now().UTC().Format(time.RFC3339Nano),
Seq: seq, At: now.Format(time.RFC3339Nano),
Kind: sessiondata.KindProviderBody, Adapter: Name + "/" + Version,
Dialect: Dialect, Src: ".", Session: session,
}
dir := c.Zone.ProviderDir(session)
if err := os.MkdirAll(dir, 0o755); err != nil {
return err
}
path := filepath.Join(dir, storage.LandedName(string(sessiondata.KindProviderBody), stamp, seq))
path := filepath.Join(dir, storage.LandedName(string(sessiondata.KindProviderBody), storage.Stamp(now), seq))
err := storage.WriteExclusive(path, storage.PermLanded, func(w io.Writer) error {
writer, err := sessiondata.NewWriter(w, header)
if err != nil {
Expand Down
20 changes: 12 additions & 8 deletions internal/adapters/langsmith/collect.go
Original file line number Diff line number Diff line change
Expand Up @@ -400,7 +400,11 @@ func (c *Collector) convertWith(request Waiting, open *pending, wait bool) (Land
}
g.items = append(g.items, item)
}
stamp := c.Now().UTC().Format("20060102T150405.000000000Z")
// One time for the whole pass: it names every file the pass lands and is
// the collected time in each file's header, as in every other adapter, so
// a reader that has only the header, such as a server that stores files by
// session and sequence, derives the file's name from it.
now := c.Now().UTC()
// Every session of the request is checked before any is landed. Checking
// one at a time landed the first and then set the request aside on the
// second, and the first was indexed but never reported for parsing.
Expand All @@ -415,7 +419,7 @@ func (c *Collector) convertWith(request Waiting, open *pending, wait bool) (Land
if len(g.items) == 0 {
continue
}
if err := c.land(g, stamp, request, open); err != nil {
if err := c.land(g, now, request, open); err != nil {
return out, err
}
out.Files += g.files
Expand Down Expand Up @@ -620,7 +624,7 @@ func (c *Collector) check(g *grouped, open *pending, request Waiting) error {
// session, and the sequence is what lets assembly track its progress with one
// watermark. Two requests carrying the same session are therefore serialised
// here, however they arrived.
func (c *Collector) land(g *grouped, stamp string, request Waiting, open *pending) error {
func (c *Collector) land(g *grouped, now time.Time, request Waiting, open *pending) error {
dir := c.Zone.SessionDir(g.session)
p := c.place(g, open)
sh, shapePath, streams, byStream := p.shape, p.shapePath, p.streams, p.byStream
Expand All @@ -642,12 +646,12 @@ func (c *Collector) land(g *grouped, stamp string, request Waiting, open *pendin
return err
}
for _, stream := range streams {
if err := c.landStream(g, state, stream, byStream[stream], stamp, request); err != nil {
if err := c.landStream(g, state, stream, byStream[stream], now, request); err != nil {
return err
}
}
if c.ProviderBodies {
landed, err := c.landBodies(g.session, state, bodiesOf(g.items, p.placed), stamp, request)
landed, err := c.landBodies(g.session, state, bodiesOf(g.items, p.placed), now, request)
if err != nil {
return err
}
Expand Down Expand Up @@ -747,7 +751,7 @@ func (c *Collector) encodable(session, stream string, records []sessiondata.Reco

// landStream writes one stream's records out of one request.
func (c *Collector) landStream(g *grouped, state *storage.SessionState, stream string,
records []sessiondata.Record, stamp string, request Waiting) error {
records []sessiondata.Record, now time.Time, request Waiting) error {
if len(records) == 0 {
return nil
}
Expand Down Expand Up @@ -780,12 +784,12 @@ func (c *Collector) landStream(g *grouped, state *storage.SessionState, stream s
}
seq := state.Take()
header := &sessiondata.Header{
Seq: seq, At: c.Now().UTC().Format(time.RFC3339Nano),
Seq: seq, At: now.Format(time.RFC3339Nano),
Kind: sessiondata.KindTranscript, Adapter: Name + "/" + Version,
Dialect: Dialect, Src: filepath.Base(request.Path),
Session: g.session, Stream: stream,
}
path := filepath.Join(streamDir, storage.LandedName("transcript", stamp, seq))
path := filepath.Join(streamDir, storage.LandedName("transcript", storage.Stamp(now), seq))
batch := records[start:end]
// A failure here is almost always the disk, and a disk is retried:
// treating every write failure as a bad request took a valid batch
Expand Down
12 changes: 10 additions & 2 deletions internal/view/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -487,10 +487,18 @@ func streamRows(c *Conversation, talks []talkRow) []sessionview.Stream {
}
parent := map[string]string{}
from := map[string][]origin{}
// In relation id order, as the relations list is. The fold holds them in a
// map, and a stream with more than one candidate listed its origins in a
// different order on each read: measured on a real conversation whose
// stream had three, five reads gave three orders.
var starts []*sessionflow.Relation
for _, r := range c.View.Relations {
if r.Type != model.RelStarts {
continue
if r.Type == model.RelStarts {
starts = append(starts, r)
}
}
sort.Slice(starts, func(i, j int) bool { return starts[i].ID < starts[j].ID })
for _, r := range starts {
if n := c.View.Nodes[r.From]; n != nil && n.Stream != "" {
parent[r.To] = n.Stream
from[r.To] = append(from[r.To], origin{r.From, n.Stream, r.Quality, c.talkOf(r.From)})
Expand Down
49 changes: 49 additions & 0 deletions internal/view/order_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,12 @@
package view

import (
"encoding/json"
"sort"
"strings"
"testing"

"github.com/apache/skywalking-ai-sessionizer/internal/storage"
"github.com/apache/skywalking-ai-sessionizer/pkg/model"
"github.com/apache/skywalking-ai-sessionizer/pkg/sessionflow"
)
Expand Down Expand Up @@ -71,3 +73,50 @@ func TestRecordTimesSortByInstant(t *testing.T) {
t.Fatalf("%s, want %s", got, want)
}
}

// TestAStreamsOriginsKeepOneOrder. A call the assembler could not tie to one
// stream leaves several candidates, and each is listed as an origin of the
// stream. The fold holds its relations in a map, so listing them in the order
// the map gave them changed the document from one read to the next: on a real
// conversation whose stream had three, five reads gave three orders. They are
// listed by relation id, as the relations list is.
func TestAStreamsOriginsKeepOneOrder(t *testing.T) {
node := func(id, kind, stream, attrs string, seq uint64) *sessionflow.Node {
n := &sessionflow.Node{Entity: sessionflow.Entity{ID: id}, Kind: kind, Stream: stream,
Ref: &sessionflow.Ref{Seq: seq, Row: 1}}
if attrs != "" {
n.Attrs = json.RawMessage(attrs)
}
return n
}
nodes := map[string]*sessionflow.Node{}
for _, n := range []*sessionflow.Node{
node("stream/main", model.KindStream, "main", `{"role":"main"}`, 1),
node("stream/c1", model.KindStream, "c1", `{"role":"child"}`, 2),
node("tool/a", model.KindTool, "main", "", 1),
node("tool/b", model.KindTool, "main", "", 1),
node("tool/c", model.KindTool, "main", "", 1),
} {
nodes[n.ID] = n
}
rels := map[string]*sessionflow.Relation{}
for _, r := range []struct{ id, from string }{{"rel/3", "tool/a"}, {"rel/1", "tool/c"}, {"rel/2", "tool/b"}} {
rels[r.id] = &sessionflow.Relation{Entity: sessionflow.Entity{ID: r.id}, Type: model.RelStarts, From: r.from, To: "stream/c1"}
}
c := &Conversation{View: &sessionflow.View{Nodes: nodes, Relations: rels}, Session: "s", zone: storage.NewZone(t.TempDir())}
want := "tool/c tool/b tool/a"
for i := 0; i < 50; i++ {
var got []string
for _, st := range streamRows(c, nil) {
if st.ID != "stream/c1" {
continue
}
for _, o := range st.OpenedBy {
got = append(got, o.Step)
}
}
if strings.Join(got, " ") != want {
t.Fatalf("read %d: the origins are %v, want %s, by relation id", i, got, want)
}
}
}
41 changes: 41 additions & 0 deletions tests/langchain_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,47 @@ func landIntoBudget(t *testing.T, zone *storage.Zone, kase string, budget int64)
return landed.Sessions
}

// TestAFileIsNamedForTheTimeItsHeaderSays. Every file one pass lands carries
// the pass's time twice: in its name and as the collected time in its header.
// A server that stores files by session and sequence keeps only the header,
// and names the file from it, so the two must agree. The collector took the
// name's time once per pass and the header's once per file. On a real root,
// all 50 LangChain files were named apart from their header, by up to 99
// milliseconds, and none of the 7,185 Claude Code files was.
func TestAFileIsNamedForTheTimeItsHeaderSays(t *testing.T) {
zone, sessions := land(t, "subagent")
checked := 0
for _, session := range sessions {
files, err := storage.LandedFiles(zone, session)
if err != nil {
t.Fatal(err)
}
for _, lf := range files {
f, err := os.Open(lf.Path)
if err != nil {
t.Fatal(err)
}
rd, err := sessiondata.NewReader(f)
if err != nil {
f.Close()
t.Fatal(err)
}
at, err := time.Parse(time.RFC3339Nano, rd.Header().At)
f.Close()
if err != nil {
t.Fatalf("%s: header time %q: %v", lf.Path, rd.Header().At, err)
}
if !strings.Contains(filepath.Base(lf.Path), "-"+storage.Stamp(at)+"-") {
t.Fatalf("%s is not named for its header's time %s", filepath.Base(lf.Path), rd.Header().At)
}
checked++
}
}
if checked < 3 {
t.Fatalf("checked %d files; the capture lands a transcript per stream and its provider bodies", checked)
}
}

// fold reads the conversation back out of its rounds.
func fold(t *testing.T, zone *storage.Zone, session string) *sessionflow.View {
t.Helper()
Expand Down
Loading