Skip to content
Merged
Show file tree
Hide file tree
Changes from 20 commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
d0ba785
fix(sessions): reliably deliver startup user messages via event queue
jh0904 Jul 30, 2026
2151cf7
test(sessions): keep core startup delivery coverage only
jh0904 Jul 30, 2026
4029d81
refactor(sessions): simplify startup queue handoff after review
jh0904 Jul 30, 2026
112e515
fix(sessions): only queue startup messages for cloud environments
jh0904 Jul 30, 2026
6c573f5
fix(db): use stable tenant UUIDs for startup queue
jh0904 Jul 31, 2026
73b2d10
fix(sessions): reinject session history when activating code sessions
jh0904 Jul 31, 2026
7eef483
refactor(db): rename session startup window helper
arthur-zhang Jul 31, 2026
c79cbf5
refactor(db): simplify session event queue existence check
arthur-zhang Jul 31, 2026
2c06800
refactor(db): simplify listSessionEventQueueIdentityRows
arthur-zhang Jul 31, 2026
ee3eee6
refactor(db): simplify ListSessionEventQueueItems event lookup
arthur-zhang Jul 31, 2026
e05b3c7
refactor(db): simplify delete session event queue query
arthur-zhang Jul 31, 2026
f55b5ab
Merge origin/main into codex/fix-session-event-reliable-delivery
jh0904 Jul 31, 2026
88033f5
fix(db): renumber session event queue migration after main UUID series
jh0904 Jul 31, 2026
3609ca3
revert: drop unrelated merge fixes from startup-delivery branch
jh0904 Jul 31, 2026
c5a63a4
refactor(db): simplify startup queue SQL and drop cloud filter docs
jh0904 Jul 31, 2026
1ccc6fc
Merge remote-tracking branch 'origin/main' into codex/fix-session-eve…
jh0904 Jul 31, 2026
080310a
refactor(sessions): move activation tx orchestration out of DB
jh0904 Jul 31, 2026
56afac5
refactor(sessions): clarify startup queue delivery and activation han…
jh0904 Aug 1, 2026
3fa3cbe
fix(sessions): atomically replay activation history
jh0904 Aug 1, 2026
c8c0ed0
Migrate session activation SQL to generated yourbatis mappers
arthur-zhang Aug 3, 2026
5ea177c
refactor(db): migrate single code-session event append to yourbatis
jh0904 Aug 3, 2026
d3e211a
fix(ci): generate yourbatis mappers before Go typecheck
jh0904 Aug 3, 2026
f65c999
merge origin/main into session startup delivery branch
jh0904 Aug 3, 2026
c7c4c90
merge origin/main into session startup delivery branch
jh0904 Aug 4, 2026
8819d86
fix: reliably deliver startup session events
jh0904 Aug 4, 2026
4d76cdb
refactor(db): clarify code-session append naming and comments
jh0904 Aug 4, 2026
2b1bdee
fix(sessions): harden activation replay
jh0904 Aug 4, 2026
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
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -69,3 +69,4 @@ logs/
CLAUDE.local.md
AGENTS.override.md
.worktrees/
*.gen.go
429 changes: 429 additions & 0 deletions docs/design/be/session-startup-message-delivery.md

Large diffs are not rendered by default.

13 changes: 9 additions & 4 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,12 @@ require (
github.com/samber/lo v1.53.0
github.com/standard-webhooks/standard-webhooks/libraries v0.0.1
github.com/superduck-ai/e2b-go-sdk v0.0.1
github.com/superduck-ai/yourbatis v0.1.0
go.opentelemetry.io/proto/otlp v1.10.0
go.yaml.in/yaml/v3 v3.0.4
golang.org/x/net v0.53.0
golang.org/x/sync v0.20.0
golang.org/x/text v0.37.0
golang.org/x/net v0.56.0
golang.org/x/sync v0.21.0
golang.org/x/text v0.38.0
google.golang.org/protobuf v1.36.11
)

Expand Down Expand Up @@ -63,9 +64,13 @@ require (
go.uber.org/atomic v1.11.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
go.yaml.in/yaml/v4 v4.0.0-rc.2 // indirect
golang.org/x/mod v0.37.0 // indirect
golang.org/x/oauth2 v0.35.0 // indirect
golang.org/x/sys v0.44.0 // indirect
golang.org/x/sys v0.46.0 // indirect
golang.org/x/tools v0.47.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20260420184626-e10c466a9529 // indirect
google.golang.org/grpc v1.80.0 // indirect
)

tool github.com/superduck-ai/yourbatis/cmd/sqlmapgen
36 changes: 20 additions & 16 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -66,8 +66,8 @@ github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ4
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE=
github.com/go-sql-driver/mysql v1.8.1/go.mod h1:wEBSXgmK//2ZFJyE+qWnIsVGmvmEKlqwuVSjsCm7DZg=
github.com/go-sql-driver/mysql v1.9.3 h1:U/N249h2WzJ3Ukj8SowVFjdtZKfu9vlLZxjPXV1aweo=
github.com/go-sql-driver/mysql v1.9.3/go.mod h1:qn46aNg1333BRMNU69Lq93t8du/dwxI64Gl8i5p1WMU=
github.com/go-sql-driver/mysql v1.10.0 h1:Q+1LV8DkHJvSYAdR83XzuhDaTykuDx0l6fkXxoWCWfw=
github.com/go-sql-driver/mysql v1.10.0/go.mod h1:M+cqaI7+xxXGG9swrdeUIoPG3Y3KCkF0pZej+SK+nWk=
github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY=
github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
Expand Down Expand Up @@ -141,6 +141,8 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/superduck-ai/e2b-go-sdk v0.0.1 h1:yW23X/fEQvaPdJHhDnp8PgnsaIixn+UTVMt94X4Kr0s=
github.com/superduck-ai/e2b-go-sdk v0.0.1/go.mod h1:dWlhv18vJamYp5gQ3ktmjjrX7tUaMe2OyTJXUo8EeXY=
github.com/superduck-ai/yourbatis v0.1.0 h1:U49je+0PhoR6PVJKfR3mbNo31rr+2SzkWwBi28wndwE=
github.com/superduck-ai/yourbatis v0.1.0/go.mod h1:BlCyyT1yfU2Zxya89rDf89keXqsdcwP6Q3PFscIjvig=
github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
github.com/tidwall/gjson v1.18.0 h1:FIDeeyB800efLX89e5a8Y0BNH+LOngJyGrIWxG2FKQY=
github.com/tidwall/gjson v1.18.0/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
Expand Down Expand Up @@ -177,18 +179,20 @@ go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
go.yaml.in/yaml/v4 v4.0.0-rc.2 h1:/FrI8D64VSr4HtGIlUtlFMGsm7H7pWTbj6vOLVZcA6s=
go.yaml.in/yaml/v4 v4.0.0-rc.2/go.mod h1:aZqd9kCMsGL7AuUv/m/PvWLdg5sjJsZ4oHDEnfPPfY0=
golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA=
golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs=
golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ=
golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0=
golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o=
golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec=
golang.org/x/oauth2 v0.35.0 h1:Mv2mzuHuZuY2+bkyWXIHMfhNdJAdwW3FuWeCPYN5GVQ=
golang.org/x/oauth2 v0.35.0/go.mod h1:lzm5WQJQwKZ3nwavOZ3IS5Aulzxi68dUSgRHujetwEA=
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.44.0 h1:ildZl3J4uzeKP07r2F++Op7E9B29JRUy+a27EibtBTQ=
golang.org/x/sys v0.44.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc=
golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38=
golang.org/x/tools v0.44.0 h1:UP4ajHPIcuMjT1GqzDWRlalUEoY+uzoZKnhOjbIPD2c=
golang.org/x/tools v0.44.0/go.mod h1:KA0AfVErSdxRZIsOVipbv3rQhVXTnlU6UhKxHd1seDI=
golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw=
golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE=
golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4=
golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q=
golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA=
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=
google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 h1:JLQynH/LBHfCTSbDWl+py8C+Rg/k1OVH3xfcaiANuF0=
Expand All @@ -207,11 +211,11 @@ gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
modernc.org/libc v1.72.1 h1:db1xwJ6u1kE3KHTFTTbe2GCrczHPKzlURP0aDC4NGD0=
modernc.org/libc v1.72.1/go.mod h1:HRMiC/PhPGLIPM7GzAFCbI+oSgE3dhZ8FWftmRrHVlY=
modernc.org/libc v1.74.1 h1:bdR4VTKFMC4966QSNZ05XLGI/VwzVa2kTUX51Dm0riQ=
modernc.org/libc v1.74.1/go.mod h1:uH4t5bOx3G3g9Xcmj10YKlTcVISlRDwv8VoQJG9n8Os=
modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU=
modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg=
modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI=
modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw=
modernc.org/sqlite v1.49.1 h1:dYGHTKcX1sJ+EQDnUzvz4TJ5GbuvhNJa8Fg6ElGx73U=
modernc.org/sqlite v1.49.1/go.mod h1:m0w8xhwYUVY3H6pSDwc3gkJ/irZT/0YEXwBlhaxQEew=
modernc.org/sqlite v1.55.0 h1:hIFh0MCH0rGinQ/4KYb5/UbCkRkb+UP+OkLCVWa5MTM=
modernc.org/sqlite v1.55.0/go.mod h1:4ntCLuNmnH8+GNqjka1wNg7KJd5/Hi5FYp8K+XQ7GZw=
87 changes: 84 additions & 3 deletions internal/codesessions/managed_agent_code_session.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@ type ManagedAgentCreateInput struct {
PermissionMode string
DangerouslySkipPermissions bool
Config json.RawMessage
InitialEvents []json.RawMessage
}

// ManagedAgentCreateResult 只在创建链路内短暂携带两份明文凭证,调用方应立即交给
Expand Down Expand Up @@ -66,7 +65,7 @@ func (s *Service) CreateManagedAgentCodeSession(ctx context.Context, input Manag
WorkDir: strings.TrimSpace(input.WorkDir),
PermissionMode: strings.TrimSpace(input.PermissionMode),
Model: strings.TrimSpace(input.Model),
Status: "active",
Status: "initializing",
Comment thread
jh0904 marked this conversation as resolved.
Metadata: metadata,
// OAuth-compatible token 只落 SHA-256 hash;明文仅存在于当前返回值中。
OAuthAccessTokenHash: auth.HashAPIKey(oauthAccessToken),
Expand Down Expand Up @@ -99,7 +98,7 @@ func (s *Service) CreateManagedAgentCodeSession(ctx context.Context, input Manag
if err := s.queueInitialize(ctx, record, input.Config, now); err != nil {
return ManagedAgentCreateResult{}, err
}
if err := s.queueInitialPublicSessionEvents(ctx, record, input.InitialEvents, now); err != nil {
if err := s.CommitManagedAgentCodeSessionActivation(ctx, record); err != nil {
return ManagedAgentCreateResult{}, err
}
credentialContext, err := s.db.GetCodeSessionCredentialContextForIssue(
Expand All @@ -126,6 +125,88 @@ func (s *Service) CreateManagedAgentCodeSession(ctx context.Context, input Manag
}, nil
}

// CommitManagedAgentCodeSessionActivation locks the owning Session, loads the
// startup queue and complete public history, writes forwardable events in
// stable order, clears the queue, and activates the Code Session atomically.
func (s *Service) CommitManagedAgentCodeSessionActivation(
ctx context.Context,
codeSession db.CodeSession,
) error {
if s == nil || s.db == nil {
return db.ErrNotFound
}
return s.db.WithManagedAgentActivationTx(ctx, func(tx db.ManagedAgentActivationTx) error {
// lock session by session external id
lockedSession, err := tx.LockSessionForEvents(
ctx,
codeSession.WorkspaceUUID,
codeSession.SessionExternalID,
)
if err != nil {
return err
}
// lock code_session by code session id
lockedCodeSession, err := tx.LockInitializingCodeSession(ctx, codeSession.UUID)
if err != nil {
return err
}
// Queue rows are only a startup-delivery responsibility check. Inbound
// payloads and order come from locked public history below; after that
// succeeds the queue is cleared as the handoff completes.
queuedEvents, err := tx.ListSessionEventQueueItems(ctx, lockedSession)
if err != nil {
return err
}
for _, event := range queuedEvents {
if event.EventType != "user.message" {
return db.ErrInvalidState
}
}
sessionEvents, err := tx.ListSessionEventsForActivation(ctx, lockedSession)
if err != nil {
return err
}
inboundInputs := make([]db.AppendCodeSessionEventInput, 0, len(sessionEvents))
for _, event := range sessionEvents {
if !shouldForwardPublicEventToWorker(event.EventType) {
continue
}
inbound, err := s.convertSessionEventToInbound(lockedCodeSession.ExternalID, event)
if err != nil {
return err
}
inboundInputs = append(inboundInputs, inbound)
}
if err := tx.AppendCodeSessionInboundEvents(ctx, lockedCodeSession, inboundInputs); err != nil {
return err
}
if err := tx.DeleteSessionEventQueue(ctx, lockedSession.UUID); err != nil {
return err
}
statusUpdated, err := tx.ActivateCodeSession(ctx, lockedCodeSession.UUID, time.Now().UTC())
if err != nil {
return err
}
if !statusUpdated {
return db.ErrInvalidState
}
return nil
})
}

// convertSessionEventToInbound maps one public session event payload into a
// Code Session inbound append input.
func (s *Service) convertSessionEventToInbound(
codeSessionID string,
event db.SessionEvent,
) (db.AppendCodeSessionEventInput, error) {
payload, err := workerPayloadForPublicEvent(codeSessionID, event.Payload, event.ProcessedAt)
if err != nil {
return db.AppendCodeSessionEventInput{}, err
}
return newInboundEventInput(codeSessionID, payload, "public-session")
}

// TerminateManagedAgentCodeSession revokes a Code Session created for a
// sandbox launch that failed before the runtime became usable.
func (s *Service) TerminateManagedAgentCodeSession(
Expand Down
55 changes: 21 additions & 34 deletions internal/codesessions/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,30 +43,6 @@ func NewServiceWithCredentials(database *db.DB, credentials *SessionCredentials,
return &Service{db: database, credentials: credentials, logger: logger}
}

func (s *Service) queueInitialPublicSessionEvents(ctx context.Context, codeSession db.CodeSession, payloads []json.RawMessage, now time.Time) error {
if len(payloads) == 0 {
return nil
}
workerPayloads := make([]json.RawMessage, 0, len(payloads))
for _, raw := range payloads {
object, err := decodeJSONObject(raw)
if err != nil {
s.logger.WarnContext(ctx, "skip initial code session event", "code_session_id", codeSession.ExternalID, "error", err)
continue
}
if !forwardPublicEventToWorker(stringField(object, "type")) {
continue
}
payload, err := workerPayloadForPublicEvent(codeSession.ExternalID, raw, now)
if err != nil {
s.logger.ErrorContext(ctx, "convert initial code session event", "code_session_id", codeSession.ExternalID, "error", err)
continue
}
workerPayloads = append(workerPayloads, payload)
}
return s.QueueRawPublicSessionEvents(ctx, codeSession, workerPayloads)
}

func (s *Service) QueuePublicSessionEvents(ctx context.Context, session db.Session, events []db.SessionEvent) error {
if s == nil || len(events) == 0 {
return nil
Expand All @@ -78,9 +54,12 @@ func (s *Service) QueuePublicSessionEvents(ctx context.Context, session db.Sessi
}
return err
}
if codeSession.Status != "active" {
return nil
Comment thread
jh0904 marked this conversation as resolved.
}
payloads := make([]json.RawMessage, 0, len(events))
for _, event := range events {
if !forwardPublicEventToWorker(event.EventType) {
if !shouldForwardPublicEventToWorker(event.EventType) {
continue
}
if event.EventType == "user.tool_confirmation" {
Expand Down Expand Up @@ -216,7 +195,7 @@ func (s *Service) AppendWorkerOutputEventsForEpoch(ctx context.Context, codeSess
return err
}
}
if event.Ephemeral || !publicWorkerOutputEvent(event.EventType) {
if event.Ephemeral || !isPublicWorkerOutputEvent(event.EventType) {
continue
}
publicPayloads, ok, err := publicPayloadsFromWorkerEvent(codeSessionID, event, publicObject)
Expand Down Expand Up @@ -282,7 +261,7 @@ func (s *Service) appendWorkerEvent(ctx context.Context, codeSessionID string, r
if meta.EventType == "control_request" && meta.EventSubtype == "can_use_tool" {
return s.handleToolPermissionRequest(ctx, codeSessionID, object, meta)
}
if hiddenWorkerEvent(meta.EventType) {
if isHiddenWorkerEvent(meta.EventType) {
return nil
}
publicPayloads, ok, err := publicPayloadsFromWorkerEvent(codeSessionID, event, object)
Expand Down Expand Up @@ -324,15 +303,23 @@ func (s *Service) queueInitialize(ctx context.Context, codeSession db.CodeSessio
}

func (s *Service) appendInboundPayload(ctx context.Context, codeSessionID string, payload json.RawMessage, source string) (db.CodeSessionEvent, bool, error) {
meta, err := BuildEventMetadata(codeSessionID, "inbound", payload)
input, err := newInboundEventInput(codeSessionID, payload, source)
if err != nil {
return db.CodeSessionEvent{}, false, err
}
return s.db.AppendCodeSessionInboundEvent(ctx, codeSessionID, input)
}

func newInboundEventInput(codeSessionID string, payload json.RawMessage, source string) (db.AppendCodeSessionEventInput, error) {
meta, err := BuildEventMetadata(codeSessionID, "inbound", payload)
if err != nil {
return db.AppendCodeSessionEventInput{}, err
}
eventID, err := ids.New("csev_")
if err != nil {
return db.CodeSessionEvent{}, false, err
return db.AppendCodeSessionEventInput{}, err
}
return s.db.AppendCodeSessionInboundEvent(ctx, codeSessionID, db.AppendCodeSessionEventInput{
return db.AppendCodeSessionEventInput{
ExternalID: eventID,
EventType: meta.EventType,
EventSubtype: meta.EventSubtype,
Expand All @@ -344,7 +331,7 @@ func (s *Service) appendInboundPayload(ctx context.Context, codeSessionID string
DeliveryStatus: "queued",
Source: strings.TrimSpace(source),
CreatedAt: time.Now().UTC(),
})
}, nil
}

func (s *Service) publishPublicPayloads(ctx context.Context, codeSessionID string, payloads []json.RawMessage) error {
Expand Down Expand Up @@ -460,7 +447,7 @@ func (s *Service) subagentThreadMappings(ctx context.Context, codeSession db.Cod
return threadByAgent, nil
}

func forwardPublicEventToWorker(eventType string) bool {
func shouldForwardPublicEventToWorker(eventType string) bool {
switch eventType {
case "user.message", "user.interrupt", "user.tool_confirmation", "user.tool_result", "user.custom_tool_result":
return true
Expand All @@ -469,7 +456,7 @@ func forwardPublicEventToWorker(eventType string) bool {
}
}

func hiddenWorkerEvent(eventType string) bool {
func isHiddenWorkerEvent(eventType string) bool {
switch eventType {
case "control_request", "control_response", "control_cancel_request":
return true
Expand All @@ -478,7 +465,7 @@ func hiddenWorkerEvent(eventType string) bool {
}
}

func publicWorkerOutputEvent(eventType string) bool {
func isPublicWorkerOutputEvent(eventType string) bool {
return maevents.IsWorkerOutputEvent(eventType) || maevents.IsStreamDelta(eventType)
}

Expand Down
25 changes: 25 additions & 0 deletions internal/db/code_session_inbound_event_mapper.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
package db

import (
"context"

"github.com/google/uuid"
)

//go:generate go tool sqlmapgen -mapper CodeSessionInboundEventMapper -sql ./code_session_inbound_event_mapper.xml -dialect postgres

// CodeSessionInboundEventMapper contains queries whose primary table is
// code_session_inbound_events.
type CodeSessionInboundEventMapper interface {
ListExistingActivationInboundEvents(
ctx context.Context,
organizationUUID uuid.UUID,
workspaceUUID uuid.UUID,
idempotencyKeys []string,
) ([]codeSessionInboundEventIdentityRow, error)

InsertCodeSessionInboundEvents(
ctx context.Context,
rows []codeSessionInboundEventInsertRow,
) (int64, error)
}
Loading
Loading