Class: SimpleAcp::Server::Base

Inherits:
Object
  • Object
show all
Defined in:
lib/simple_acp/server/base.rb

Overview

Main ACP Server class for hosting agents and handling requests.

The server manages agent registration, run execution (sync, async, stream), session state, and exposes an HTTP API via Roda/Falcon.

Examples:

Creating and running a server

server = SimpleAcp::Server::Base.new
server.agent("echo", description: "Echoes input") do |context|
  SimpleAcp::Models::Message.agent(context.input.first.text_content)
end
server.run(port: 8000)

Using custom storage

storage = SimpleAcp::Storage::Redis.new(url: "redis://localhost:6379")
server = SimpleAcp::Server::Base.new(storage: storage)

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(storage: nil, **options) ⇒ Base

Initialize a new ACP server.

Parameters:

  • storage (Storage::Base, nil) (defaults to: nil) —

    storage backend (defaults to Memory)

  • options (Hash) —

    additional configuration options



34
35
36
37
38
39
# File 'lib/simple_acp/server/base.rb', line 34

def initialize(storage: nil, **options)
  @agents = {}
  @storage = storage || SimpleAcp::Storage::Memory.new
  @options = options
  @running_contexts = Concurrent::Map.new
end

Instance Attribute Details

#agents ⇒ Hash<String, Agent> (readonly)

Returns registered agents indexed by name.

Returns:

  • (Hash<String, Agent>) —

    registered agents indexed by name



22
23
24
# File 'lib/simple_acp/server/base.rb', line 22

def agents
  @agents
end

#options ⇒ Hash (readonly)

Returns additional configuration options.

Returns:

  • (Hash) —

    additional configuration options



28
29
30
# File 'lib/simple_acp/server/base.rb', line 28

def options
  @options
end

#storage ⇒ Storage::Base (readonly)

Returns storage backend for runs, sessions, and events.

Returns:

  • (Storage::Base) —

    storage backend for runs, sessions, and events



25
26
27
# File 'lib/simple_acp/server/base.rb', line 25

def storage
  @storage
end

Instance Method Details

#agent(name = nil, description: nil, **options) {|Context| ... } ⇒ Agent, Proc

Register an agent using block syntax or decorator-style.

Examples:

Block syntax

server.agent("greeter", description: "Greets users") do |context|
  name = context.input.first&.text_content || "World"
  SimpleAcp::Models::Message.agent("Hello, #{name}!")
end

Streaming agent

server.agent("counter") do |context|
  Enumerator.new do |yielder|
    3.times { |i| yielder << SimpleAcp::Models::Message.agent("Count: #{i}") }
  end
end

Parameters:

  • name (String, nil) (defaults to: nil) —

    agent name (must follow RFC 1123 DNS label format)

  • description (String, nil) (defaults to: nil) —

    human-readable description

  • options (Hash) —

    additional options

Options Hash (**options):

  • :input_content_types (Array<String>) —

    accepted MIME types

  • :output_content_types (Array<String>) —

    produced MIME types

  • :metadata (Hash) —

    agent metadata

Yields:

  • (Context) —

    block that handles agent requests

Returns:

  • (Agent, Proc) —

    the registered agent or a decorator lambda



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
# File 'lib/simple_acp/server/base.rb', line 64

def agent(name = nil, description: nil, **options, &block)
  if block_given?
    # Direct registration with block
    agent_obj = AgentDSL.define(
      name: name,
      description: description,
      **options,
      &block
    )
    register(agent_obj)
    agent_obj
  else
    # Return a lambda for decorator-style usage
    ->(handler) do
      agent_name = name || handler_name(handler)
      agent_obj = Agent.new(
        manifest: Models::AgentManifest.new(
          name: agent_name,
          description: description,
          **options
        ),
        handler: handler
      )
      register(agent_obj)
      handler
    end
  end
end

#cancel_run(run_id) ⇒ Models::Run

Cancel a running agent execution.

Parameters:

  • run_id (String) —

    the run ID to cancel

Returns:

Raises:



369
370
371
372
373
374
375
376
377
378
379
# File 'lib/simple_acp/server/base.rb', line 369

def cancel_run(run_id)
  run = @storage.get_run(run_id)
  raise SimpleAcp::NotFoundError, "Run '#{run_id}' not found" unless run

  context = @running_contexts[run_id]
  context&.cancel!

  run.cancelled!
  @storage.save_run(run)
  run
end

#register(agent) ⇒ Agent

Register an agent instance directly.

Parameters:

  • agent (Agent) —

    the agent to register

Returns:

  • (Agent) —

    the registered agent

Raises:



99
100
101
102
103
104
105
# File 'lib/simple_acp/server/base.rb', line 99

def register(agent)
  raise SimpleAcp::ValidationError, "Invalid agent" unless agent.valid?
  raise SimpleAcp::ConfigurationError, "Agent '#{agent.name}' already registered" if @agents.key?(agent.name)

  @agents[agent.name] = agent
  agent
end

#resume_stream(run_id:, await_resume:) {|Models::Event| ... } ⇒ void

This method returns an undefined value.

Resume an awaited run with streaming output.

Parameters:

  • run_id (String) —

    the run ID to resume

  • await_resume (Models::AwaitResume) —

    the resume payload with client response

Yields:

Raises:



322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
# File 'lib/simple_acp/server/base.rb', line 322

def resume_stream(run_id:, await_resume:)
  run, context = prepare_resume(run_id, await_resume)

  begin
    run.start!
    @storage.save_run(run)
    yield Models::RunInProgressEvent.new(run_id: run.run_id)
    @storage.add_event(run.run_id, Models::RunInProgressEvent.new(run_id: run.run_id))

    output_messages = run.output.dup

    execute_agent(context) do |yielded|
      case yielded
      when RunYield
        message = yielded.message
        output_messages << message

        yield Models::MessageCreatedEvent.new(message: message)
        @storage.add_event(run.run_id, Models::MessageCreatedEvent.new(message: message))

        yield Models::MessageCompletedEvent.new(message: message)
        @storage.add_event(run.run_id, Models::MessageCompletedEvent.new(message: message))
      when RunYieldAwait
        yield Models::RunAwaitingEvent.new(run_id: run.run_id, await_request: yielded.request)
        @storage.add_event(run.run_id, Models::RunAwaitingEvent.new(run_id: run.run_id, await_request: yielded.request))
        return
      end
    end

    run.complete!(output_messages)
    yield Models::RunCompletedEvent.new(run: run)
    @storage.add_event(run.run_id, Models::RunCompletedEvent.new(run: run))
  rescue StandardError => e
    run.fail!(e.message)
    yield Models::RunFailedEvent.new(run_id: run.run_id, error: run.error)
    @storage.add_event(run.run_id, Models::RunFailedEvent.new(run_id: run.run_id, error: run.error))
  ensure
    @storage.save_run(run)
    @running_contexts.delete(run.run_id)
  end
end

#resume_sync(run_id:, await_resume:) ⇒ Models::Run

Resume an awaited run synchronously.

When an agent yields a RunYieldAwait, the run enters an "awaiting" state. Use this method to provide the requested input and continue execution.

Parameters:

  • run_id (String) —

    the run ID to resume

  • await_resume (Models::AwaitResume) —

    the resume payload with client response

Returns:

Raises:



286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
# File 'lib/simple_acp/server/base.rb', line 286

def resume_sync(run_id:, await_resume:)
  run, context = prepare_resume(run_id, await_resume)

  begin
    run.start!
    @storage.save_run(run)

    output_messages = run.output.dup

    execute_agent(context) do |yielded|
      case yielded
      when RunYield
        output_messages << yielded.message
      when RunYieldAwait
        return run
      end
    end

    run.complete!(output_messages)
    update_session_history(context.session, [await_resume.message].compact, output_messages)
  rescue StandardError => e
    run.fail!(e.message)
  end

  @storage.save_run(run)
  run
end

#run(port: 8000, host: "0.0.0.0", **options) ⇒ void

This method returns an undefined value.

Start the HTTP server using Falcon.

Falcon provides fiber-based concurrency for efficient handling of SSE streams and long-lived connections.

Parameters:

  • port (Integer) (defaults to: 8000) —

    port to listen on (default: 8000)

  • host (String) (defaults to: "0.0.0.0") —

    host to bind to (default: "0.0.0.0")

  • options (Hash) —

    additional Falcon configuration options



400
401
402
403
404
405
406
407
408
# File 'lib/simple_acp/server/base.rb', line 400

def run(port: 8000, host: "0.0.0.0", **options)
  require_relative "falcon_runner"

  app = to_app

  puts "Registered agents: #{@agents.keys.join(', ')}"

  FalconRunner.run(app, port: port, host: host, **options)
end

#run_async(agent_name:, input:, session_id: nil, session: nil) ⇒ Models::Run

Run an agent asynchronously, returning immediately with a run ID.

The agent executes in a background thread. Use #cancel_run to stop or poll the storage to check status.

Parameters:

  • agent_name (String) —

    name of the agent to run

  • input (Array<Models::Message>, Models::Message, String) —

    input messages

  • session_id (String, nil) (defaults to: nil) —

    optional session ID

  • session (Models::Session, Hash, nil) (defaults to: nil) —

    optional session data

Returns:

  • (Models::Run) —

    the run (status will be :created or :in_progress)

Raises:



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
# File 'lib/simple_acp/server/base.rb', line 162

def run_async(agent_name:, input:, session_id: nil, session: nil)
  run, context = prepare_run(agent_name, input, session_id, session)

  Thread.new do
    begin
      run.start!
      @storage.save_run(run)

      output_messages = []
      awaiting = false

      execute_agent(context) do |yielded|
        case yielded
        when RunYield
          output_messages << yielded.message
        when RunYieldAwait
          awaiting = true
          break
        end
      end

      # Check if cancelled or awaiting before completing
      if context.cancelled?
        run.cancelled!
      elsif !awaiting
        run.complete!(output_messages)
        update_session_history(context.session, input, output_messages)
      end
      # If awaiting, run is already in awaiting state from await_message
    rescue StandardError => e
      run.fail!(e.message)
    ensure
      @storage.save_run(run)
      @running_contexts.delete(run.run_id)
    end
  end

  run
end

#run_stream(agent_name:, input:, session_id: nil, session: nil) {|Models::Event| ... } ⇒ void

This method returns an undefined value.

Run an agent with streaming output via Server-Sent Events.

Yields events as the agent executes, enabling real-time response streaming.

Examples:

server.run_stream(agent_name: "echo", input: "Hello") do |event|
  case event
  when Models::MessagePartEvent
    print event.part.content
  when Models::RunCompletedEvent
    puts "\nDone!"
  end
end

Parameters:

  • agent_name (String) —

    name of the agent to run

  • input (Array<Models::Message>, Models::Message, String) —

    input messages

  • session_id (String, nil) (defaults to: nil) —

    optional session ID

  • session (Models::Session, Hash, nil) (defaults to: nil) —

    optional session data

Yields:

Yield Parameters:

Raises:



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
# File 'lib/simple_acp/server/base.rb', line 224

def run_stream(agent_name:, input:, session_id: nil, session: nil)
  run, context = prepare_run(agent_name, input, session_id, session)

  begin
    yield Models::RunCreatedEvent.new(run: run)
    @storage.add_event(run.run_id, Models::RunCreatedEvent.new(run: run))

    run.start!
    @storage.save_run(run)
    yield Models::RunInProgressEvent.new(run_id: run.run_id)
    @storage.add_event(run.run_id, Models::RunInProgressEvent.new(run_id: run.run_id))

    output_messages = []

    execute_agent(context) do |yielded|
      case yielded
      when RunYield
        message = yielded.message
        output_messages << message

        yield Models::MessageCreatedEvent.new(message: message)
        @storage.add_event(run.run_id, Models::MessageCreatedEvent.new(message: message))

        message.parts.each do |part|
          yield Models::MessagePartEvent.new(part: part)
          @storage.add_event(run.run_id, Models::MessagePartEvent.new(part: part))
        end

        yield Models::MessageCompletedEvent.new(message: message)
        @storage.add_event(run.run_id, Models::MessageCompletedEvent.new(message: message))
      when RunYieldAwait
        yield Models::RunAwaitingEvent.new(run_id: run.run_id, await_request: yielded.request)
        @storage.add_event(run.run_id, Models::RunAwaitingEvent.new(run_id: run.run_id, await_request: yielded.request))
        return
      end
    end

    run.complete!(output_messages)
    update_session_history(context.session, input, output_messages)

    yield Models::RunCompletedEvent.new(run: run)
    @storage.add_event(run.run_id, Models::RunCompletedEvent.new(run: run))
  rescue StandardError => e
    run.fail!(e.message)
    yield Models::RunFailedEvent.new(run_id: run.run_id, error: run.error)
    @storage.add_event(run.run_id, Models::RunFailedEvent.new(run_id: run.run_id, error: run.error))
  ensure
    @storage.save_run(run)
    @running_contexts.delete(run.run_id)
  end
end

#run_sync(agent_name:, input:, session_id: nil, session: nil) ⇒ Models::Run

Run an agent synchronously, blocking until completion.

Parameters:

  • agent_name (String) —

    name of the agent to run

  • input (Array<Models::Message>, Models::Message, String) —

    input messages

  • session_id (String, nil) (defaults to: nil) —

    optional session ID for stateful interactions

  • session (Models::Session, Hash, nil) (defaults to: nil) —

    optional session data

Returns:

Raises:



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
# File 'lib/simple_acp/server/base.rb', line 123

def run_sync(agent_name:, input:, session_id: nil, session: nil)
  run, context = prepare_run(agent_name, input, session_id, session)

  begin
    run.start!
    @storage.save_run(run)

    output_messages = []
    execute_agent(context) do |yielded|
      case yielded
      when RunYield
        output_messages << yielded.message
      when RunYieldAwait
        # Agent is awaiting - save state and return
        return run
      end
    end

    run.complete!(output_messages)
    update_session_history(context.session, input, output_messages)
  rescue StandardError => e
    run.fail!(e.message)
  end

  @storage.save_run(run)
  run
end

#to_app ⇒ Roda

Create a Rack-compatible application.

Returns:

  • (Roda) —

    the Rack application



384
385
386
387
388
389
# File 'lib/simple_acp/server/base.rb', line 384

def to_app
  # Create a subclass to avoid freezing the base App class
  app_class = Class.new(App)
  app_class.configure(self)
  app_class.freeze.app
end

#unregister(name) ⇒ Agent?

Unregister an agent by name.

Parameters:

  • name (String) —

    the agent name to remove

Returns:

  • (Agent, nil) —

    the removed agent or nil if not found



111
112
113
# File 'lib/simple_acp/server/base.rb', line 111

def unregister(name)
  @agents.delete(name)
end