Remove ununsed ExecutionEvent from controller2 (#165)

This commit is contained in:
Jaana Dogan
2026-06-24 14:24:07 -07:00
committed by GitHub
parent 22a9daeddf
commit 094f057d9f
7 changed files with 6 additions and 223 deletions
+1 -4
View File
@@ -164,10 +164,7 @@ func TestController2_ExecHelloWorld(t *testing.T) {
t.Errorf("expected third event state to be COMPLETED, got %v", events[2].State)
}
// Verify that no Execution Events were logged (moving away from granular execution log)
if len(log.AllExecEvents) != 0 {
t.Errorf("expected 0 execution events in V2, got %d: %v", len(log.AllExecEvents), log.AllExecEvents)
}
}
func TestController2_ExecWithAgentID(t *testing.T) {
@@ -32,15 +32,9 @@ type EventLog interface {
// Append adds a conversation event to the end of the log.
Append(ctx context.Context, event *proto.ConversationEvent) (int32, error)
// AppendExec adds an execution event to the end of the log.
AppendExec(ctx context.Context, event *proto.ExecutionEvent) error
// Events returns all events for the conversation.
Events(ctx context.Context, conversationID string) ([]*proto.ConversationEvent, error)
// ExecEvents returns all events for a specific execution ID.
ExecEvents(ctx context.Context, execID string) ([]*proto.ExecutionEvent, error)
// DeleteAll deletes all events for a specific conversation ID.
DeleteAll(ctx context.Context, conversationID string) error
@@ -26,7 +26,6 @@ import (
type MemoryEventLog struct {
mu sync.Mutex
AllEvents []*proto.ConversationEvent
AllExecEvents []*proto.ExecutionEvent
}
func (m *MemoryEventLog) Append(_ context.Context, event *proto.ConversationEvent) (int32, error) {
@@ -48,14 +47,6 @@ func (m *MemoryEventLog) Append(_ context.Context, event *proto.ConversationEven
return seq, nil
}
func (m *MemoryEventLog) AppendExec(_ context.Context, event *proto.ExecutionEvent) error {
m.mu.Lock()
defer m.mu.Unlock()
m.AllExecEvents = append(m.AllExecEvents, event)
return nil
}
func (m *MemoryEventLog) Events(_ context.Context, conversationID string) ([]*proto.ConversationEvent, error) {
m.mu.Lock()
defer m.mu.Unlock()
@@ -69,19 +60,6 @@ func (m *MemoryEventLog) Events(_ context.Context, conversationID string) ([]*pr
return out, nil
}
func (m *MemoryEventLog) ExecEvents(_ context.Context, execID string) ([]*proto.ExecutionEvent, error) {
m.mu.Lock()
defer m.mu.Unlock()
out := make([]*proto.ExecutionEvent, 0)
for _, ev := range m.AllExecEvents {
if ev.ExecId == execID {
out = append(out, ev)
}
}
return out, nil
}
// Drop removes every event for which drop returns true.
// It is provided for testing and crash-simulation purposes.
func (m *MemoryEventLog) Drop(drop func(*proto.ConversationEvent) bool) {
@@ -101,34 +79,14 @@ func (m *MemoryEventLog) DeleteAll(_ context.Context, conversationID string) err
m.mu.Lock()
defer m.mu.Unlock()
var execIDs []string
var keptEvents []*proto.ConversationEvent
for _, ev := range m.AllEvents {
if ev.ConversationId == conversationID {
if ev.ExecId != "" {
execIDs = append(execIDs, ev.ExecId)
}
} else {
if ev.ConversationId != conversationID {
keptEvents = append(keptEvents, ev)
}
}
m.AllEvents = keptEvents
var keptExecEvents []*proto.ExecutionEvent
for _, ev := range m.AllExecEvents {
delete := false
for _, id := range execIDs {
if ev.ExecId == id {
delete = true
break
}
}
if !delete {
keptExecEvents = append(keptExecEvents, ev)
}
}
m.AllExecEvents = keptExecEvents
return nil
}
-15
View File
@@ -52,21 +52,6 @@ func OpenPostgresEventLog(dsn string) (EventLog, error) {
return nil, fmt.Errorf("postgres_eventlog: create conversation_log table: %w", err)
}
if _, err := db.ExecContext(ctx, `
CREATE TABLE IF NOT EXISTS execution_log (
exec_id TEXT NOT NULL,
payload TEXT NOT NULL,
timestamp TIMESTAMPTZ NOT NULL
)`); err != nil {
db.Close()
return nil, fmt.Errorf("postgres_eventlog: create execution_log table: %w", err)
}
// Create indexes if they don't exist.
if _, err := db.ExecContext(ctx, `CREATE INDEX IF NOT EXISTS idx_execution_log_exec_id ON execution_log(exec_id)`); err != nil {
db.Close()
return nil, fmt.Errorf("postgres_eventlog: create index exec_id: %w", err)
}
return &sqlEventLog{db: db}, nil
}
+2 -105
View File
@@ -18,7 +18,6 @@ import (
"context"
"database/sql"
"fmt"
"time"
"github.com/google/ax/proto"
)
@@ -63,30 +62,6 @@ func (l *sqlEventLog) Append(ctx context.Context, event *proto.ConversationEvent
return seq, nil
}
// AppendExec inserts an execution event into the database.
// TODO(anj): Remove execution_log table and AppendExec when legacy controller is removed.
func (l *sqlEventLog) AppendExec(ctx context.Context, event *proto.ExecutionEvent) error {
payload, err := marshalOpts.Marshal(event)
if err != nil {
return fmt.Errorf("eventlog: marshal exec: %w", err)
}
var timestamp time.Time
if event.Timestamp != nil {
timestamp = event.Timestamp.AsTime()
} else {
timestamp = time.Now()
}
if _, err := l.db.ExecContext(ctx,
"INSERT INTO execution_log (exec_id, payload, timestamp) VALUES ($1, $2, $3)",
event.ExecId, string(payload), timestamp); err != nil {
return fmt.Errorf("eventlog: insert exec: %w", err)
}
return nil
}
// Events retrieves all events from the database for a conversation, ordered by seq.
func (l *sqlEventLog) Events(ctx context.Context, conversationID string) ([]*proto.ConversationEvent, error) {
rows, err := l.db.QueryContext(ctx, "SELECT payload FROM conversation_log WHERE conversation_id = $1 ORDER BY seq", conversationID)
@@ -116,89 +91,11 @@ func (l *sqlEventLog) Events(ctx context.Context, conversationID string) ([]*pro
return events, nil
}
// ExecEvents retrieves all events from the database for a specific execution ID.
func (l *sqlEventLog) ExecEvents(ctx context.Context, execID string) ([]*proto.ExecutionEvent, error) {
rows, err := l.db.QueryContext(ctx, "SELECT payload FROM execution_log WHERE exec_id = $1 ORDER BY timestamp", execID)
if err != nil {
return nil, fmt.Errorf("eventlog: query exec: %w", err)
}
defer rows.Close()
var events []*proto.ExecutionEvent
for rows.Next() {
var payload string
if err := rows.Scan(&payload); err != nil {
return nil, fmt.Errorf("eventlog: scan exec: %w", err)
}
ev := &proto.ExecutionEvent{}
if err := unmarshalOpts.Unmarshal([]byte(payload), ev); err != nil {
continue
}
events = append(events, ev)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("eventlog: iterate exec: %w", err)
}
return events, nil
}
// DeleteAll deletes all events for a specific conversation ID and its child executions.
// DeleteAll deletes all events for a specific conversation ID.
func (l *sqlEventLog) DeleteAll(ctx context.Context, conversationID string) error {
tx, err := l.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("eventlog: begin tx: %w", err)
}
defer tx.Rollback()
// TODO(jbd): Update the schema to include conversation_id at every execution event.
// Get all exec_ids for this conversation.
rows, err := tx.QueryContext(ctx, "SELECT payload FROM conversation_log WHERE conversation_id = $1", conversationID)
if err != nil {
return fmt.Errorf("eventlog: query conversation: %w", err)
}
defer rows.Close()
var execIDs []string
for rows.Next() {
var payload string
if err := rows.Scan(&payload); err != nil {
return fmt.Errorf("eventlog: scan conversation: %w", err)
}
ev := &proto.ConversationEvent{}
if err := unmarshalOpts.Unmarshal([]byte(payload), ev); err != nil {
return fmt.Errorf("eventlog: unmarshal event: %w", err)
}
if ev.ExecId != "" {
execIDs = append(execIDs, ev.ExecId)
}
}
if err := rows.Err(); err != nil {
return fmt.Errorf("eventlog: iterate conversation: %w", err)
}
rows.Close()
// Delete from execution_log.
for _, execID := range execIDs {
if _, err := tx.ExecContext(ctx, "DELETE FROM execution_log WHERE exec_id = $1", execID); err != nil {
return fmt.Errorf("eventlog: delete exec %s: %w", execID, err)
}
}
// Delete from conversation_log.
if _, err := tx.ExecContext(ctx, "DELETE FROM conversation_log WHERE conversation_id = $1", conversationID); err != nil {
if _, err := l.db.ExecContext(ctx, "DELETE FROM conversation_log WHERE conversation_id = $1", conversationID); err != nil {
return fmt.Errorf("eventlog: delete conversation: %w", err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("eventlog: commit tx: %w", err)
}
return nil
}
+2 -35
View File
@@ -21,7 +21,6 @@ import (
"testing"
"github.com/google/ax/proto"
"google.golang.org/protobuf/types/known/timestamppb"
)
// testEventLog runs the EventLog contract against a backend. newLog returns a
@@ -67,29 +66,7 @@ func testEventLog(t *testing.T, newLog func(t *testing.T) EventLog) {
t.Errorf("conversation events mismatch: %q, %q", cEvents[0].ExecId, cEvents[1].ExecId)
}
// 2. Execution log.
ee1 := &proto.ExecutionEvent{
ExecId: task1,
State: proto.State_STATE_PENDING,
Timestamp: timestamppb.Now(),
Inputs: []*proto.Message{
{Role: "user", Content: &proto.Content{Type: &proto.Content_Text{Text: &proto.TextContent{Text: "hello"}}}},
},
}
if err := log.AppendExec(ctx, ee1); err != nil {
t.Fatalf("failed to append ee1: %v", err)
}
eEvents, err := log.ExecEvents(ctx, task1)
if err != nil {
t.Fatalf("failed to read execution events: %v", err)
}
if len(eEvents) != 1 {
t.Fatalf("expected 1 execution event, got %d", len(eEvents))
}
if eEvents[0].ExecId != task1 || eEvents[0].State != proto.State_STATE_PENDING {
t.Errorf("execution event mismatch: %+v", eEvents[0])
}
})
t.Run("Empty", func(t *testing.T) {
@@ -126,12 +103,7 @@ func testEventLog(t *testing.T, newLog func(t *testing.T) EventLog) {
if _, err := log.Append(ctx, &proto.ConversationEvent{ConversationId: conv2, Seq: 1, ExecId: task3}); err != nil {
t.Fatalf("append: %v", err)
}
if err := log.AppendExec(ctx, &proto.ExecutionEvent{ExecId: task1, State: proto.State_STATE_PENDING}); err != nil {
t.Fatalf("append exec: %v", err)
}
if err := log.AppendExec(ctx, &proto.ExecutionEvent{ExecId: task3, State: proto.State_STATE_PENDING}); err != nil {
t.Fatalf("append exec: %v", err)
}
if err := log.DeleteAll(ctx, conv1); err != nil {
t.Fatalf("failed to delete events: %v", err)
@@ -143,12 +115,7 @@ func testEventLog(t *testing.T, newLog func(t *testing.T) EventLog) {
if ev, _ := log.Events(ctx, conv2); len(ev) != 1 {
t.Errorf("expected 1 event for conv2, got %d", len(ev))
}
if ee, _ := log.ExecEvents(ctx, task1); len(ee) != 0 {
t.Errorf("expected 0 exec events for task1, got %d", len(ee))
}
if ee, _ := log.ExecEvents(ctx, task3); len(ee) != 1 {
t.Errorf("expected 1 exec event for task3, got %d", len(ee))
}
})
// AutoSeq exercises the seq==0 auto-assignment path: appends with Seq unset
-15
View File
@@ -51,21 +51,6 @@ func OpenSQLiteEventLog(path string) (EventLog, error) {
return nil, fmt.Errorf("sqlite_eventlog: create conversation_log table: %w", err)
}
if _, err := db.Exec(`
CREATE TABLE IF NOT EXISTS execution_log (
exec_id TEXT NOT NULL,
payload TEXT NOT NULL,
timestamp DATETIME NOT NULL
)`); err != nil {
db.Close()
return nil, fmt.Errorf("sqlite_eventlog: create execution_log table: %w", err)
}
// Create indexes if they don't exist.
if _, err := db.Exec(`CREATE INDEX IF NOT EXISTS idx_execution_log_exec_id ON execution_log(exec_id)`); err != nil {
db.Close()
return nil, fmt.Errorf("sqlite_eventlog: create index exec_id: %w", err)
}
return &sqlEventLog{db: db}, nil
}