Skip to content
Open
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
2 changes: 1 addition & 1 deletion AGENT_VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
v0.6.53
v0.6.54
2 changes: 1 addition & 1 deletion lib/console/deployments/agents.ex
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ defmodule Console.Deployments.Agents do


@spec upsert_agent_runtime(map, Cluster.t | nil, User.t | Cluster.t | nil) :: agent_runtime_resp
def upsert_agent_runtime(%{name: name} = attrs, %Cluster{id: cid} = cluster, actor) do
def upsert_agent_runtime(%{name: name} = attrs, %Cluster{id: cid}, actor) do
runtime = get_agent_runtime(cid, name) |> Repo.preload([:create_bindings])

with {:ok, _} <- allow(runtime || %AgentRuntime{cluster_id: cid}, actor, :write),
Expand Down
5 changes: 5 additions & 0 deletions lib/console/deployments/policies/rbac.ex
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ defmodule Console.Deployments.Policies.Rbac do
WorkbenchTool,
WorkbenchCron,
WorkbenchPrompt,
QueuedPrompt,
WorkbenchSkill,
WorkbenchEval,
WorkbenchEvalResult,
Expand Down Expand Up @@ -160,6 +161,8 @@ defmodule Console.Deployments.Policies.Rbac do
do: recurse(cron, user, action, & &1.workbench)
def evaluate(%WorkbenchPrompt{} = prompt, user, action),
do: recurse(prompt, user, action, & &1.workbench)
def evaluate(%QueuedPrompt{} = prompt, user, action),
do: recurse(prompt, user, action, & &1.workbench_job)
def evaluate(%WorkbenchSkill{} = skill, user, action),
do: recurse(skill, user, action, & &1.workbench)
def evaluate(%WorkbenchEval{} = eval, user, action),
Expand Down Expand Up @@ -309,6 +312,8 @@ defmodule Console.Deployments.Policies.Rbac do
do: Repo.preload(cron, [workbench: [:read_bindings, :write_bindings, project: @bindings]])
def preload(%WorkbenchPrompt{} = prompt),
do: Repo.preload(prompt, [workbench: [:read_bindings, :write_bindings, project: @bindings]])
def preload(%QueuedPrompt{} = prompt),
do: Repo.preload(prompt, [workbench_job: [workbench: [:read_bindings, :write_bindings, project: @bindings]]])
def preload(%WorkbenchSkill{} = skill),
do: Repo.preload(skill, [workbench: [:read_bindings, :write_bindings, project: @bindings]])
def preload(%WorkbenchEval{} = eval),
Expand Down
8 changes: 6 additions & 2 deletions lib/console/deployments/pubsub/recurse.ex
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,10 @@ defimpl Console.PubSub.Recurse, for: [Console.PubSub.PullRequestCreated, Console
with %PullRequest{stack: %Stack{} = stack} <- Repo.preload(pr, [stack: :repository]),
_ <- sleep(stack.repository),
_ <- Discovery.kick(stack.repository),
do: Stacks.poll(stack)
{:ok, run} <- Stacks.poll(stack) do
PullRequest.changeset(pr, %{stack_run_id: run.id})
|> Repo.update()
end
end

def process(%@for{item: %PullRequest{status: :merged, service_id: id} = pr}) when is_binary(id) do
Expand Down Expand Up @@ -237,7 +240,7 @@ end

defimpl Console.PubSub.Recurse, for: [Console.PubSub.StackRunCompleted] do
alias Console.Schema.{Stack, StackRun, PullRequest}
alias Console.Deployments.Stacks
alias Console.Deployments.{Stacks, Workbenches}

def process(%{item: %StackRun{id: id} = run}) do
run = Console.Repo.preload(run, [:pull_request, :stack])
Expand All @@ -250,6 +253,7 @@ defimpl Console.PubSub.Recurse, for: [Console.PubSub.StackRunCompleted] do
Stacks.post_comment(run)
Stacks.dequeue(pr)
%StackRun{stack: %Stack{} = stack} ->
Comment thread
michaeljguarino marked this conversation as resolved.
Workbenches.kick_workbench(run)
Stacks.dequeue(stack)
end
end
Expand Down
93 changes: 92 additions & 1 deletion lib/console/deployments/workbenches.ex
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,13 @@ defmodule Console.Deployments.Workbenches do
WorkbenchJobActivityAgentRun,
WorkbenchJobThought,
PullRequest,
FlowWorkbench
FlowWorkbench,
StackRun,
QueuedPrompt
}
alias Console.AI.{Provider, VectorStore}
alias Console.AI.Tools.Workbench.{FunctionCall, KubeRequest, SavedPrompt}
alias Console.Services.Users
alias Console.Deployments.Settings
alias Console.PubSub

Expand All @@ -38,6 +41,7 @@ defmodule Console.Deployments.Workbenches do
@type activity_resp :: {:ok, WorkbenchJobActivity.t()} | error
@type cron_resp :: {:ok, WorkbenchCron.t()} | error
@type prompt_resp :: {:ok, WorkbenchPrompt.t()} | error
@type queued_prompt_resp :: {:ok, QueuedPrompt.t()} | error
@type skill_resp :: {:ok, WorkbenchSkill.t()} | error
@type eval_resp :: {:ok, WorkbenchEval.t()} | error
@type webhook_resp :: {:ok, WorkbenchWebhook.t()} | error
Expand Down Expand Up @@ -67,6 +71,8 @@ defmodule Console.Deployments.Workbenches do
def get_workbench_cron(id), do: Repo.get(WorkbenchCron, id)
def get_workbench_prompt!(id), do: Repo.get!(WorkbenchPrompt, id)
def get_workbench_prompt(id), do: Repo.get(WorkbenchPrompt, id)
def get_queued_prompt!(id), do: Repo.get!(QueuedPrompt, id)
def get_queued_prompt(id), do: Repo.get(QueuedPrompt, id)
def get_workbench_skill!(id), do: Repo.get!(WorkbenchSkill, id)
def get_workbench_skill(id), do: Repo.get(WorkbenchSkill, id)
def get_workbench_webhook!(id), do: Repo.get!(WorkbenchWebhook, id)
Expand Down Expand Up @@ -786,6 +792,76 @@ defmodule Console.Deployments.Workbenches do
|> then(&create_message(attrs, &1, user))
end

@doc """
Queues a prompt to be sent to a workbench job after `dequeable_at`.
Requires read/prompt access to the target job.
"""
@spec create_queued_prompt(map, binary | WorkbenchJob.t(), User.t()) :: queued_prompt_resp
def create_queued_prompt(attrs, %WorkbenchJob{id: job_id}, %User{id: user_id} = user) do
%QueuedPrompt{workbench_job_id: job_id, user_id: user_id}
|> QueuedPrompt.changeset(attrs)
|> allow(user, :read)
|> when_ok(:insert)
end

def create_queued_prompt(attrs, job_id, %User{} = user) when is_binary(job_id) do
get_workbench_job!(job_id)
|> then(&create_queued_prompt(attrs, &1, user))
end

@doc """
Deletes a queued prompt before or after it has been consumed.
Requires read/prompt access to the target job.
"""
@spec delete_queued_prompt(binary, User.t()) :: queued_prompt_resp
def delete_queued_prompt(id, %User{} = user) do
get_queued_prompt!(id)
|> allow(user, :read)
|> when_ok(:delete)
end

@spec dequeue_prompt(QueuedPrompt.t()) :: activity_resp
def dequeue_prompt(%QueuedPrompt{} = prompt) do
%{user: user, workbench_job: job} = Repo.preload(prompt, [:workbench_job, user: [:groups]])

start_transaction()
|> add_operation(:consume, fn _ ->
prompt
|> QueuedPrompt.changeset(%{consumed_at: DateTime.utc_now()})
|> Repo.update()
end)
|> add_operation(:job, fn %{consume: prompt} ->
create_message(%{
prompt: prompt.prompt
}, job, user)
end)
|> execute(extract: :job)
end

def kick_workbench(%StackRun{status: :successful, id: id}) do
WorkbenchJob.for_stack_run(id)
|> WorkbenchJob.with_limit(1)
|> Repo.one()
|> case do
%WorkbenchJob{} = job ->
%PullRequest{} = pr = PullRequest.for_stack_run(id)
|> Repo.one!()
create_message(%{
prompt: String.trim(stack_run_verification_prompt(pr: pr))
}, job, Users.get_bot!("console"))
Comment thread
michaeljguarino marked this conversation as resolved.
nil ->
{:error, "no workbench job found for stack run #{id}"}
end
end
def kick_workbench(_), do: :ok

EEx.function_from_file(
:defp,
:stack_run_verification_prompt,
"priv/prompts/workbench/stack_run_verification.md.eex",
[:assigns]
)

@doc """
Creates a new message for the job associated with a pull request.
"""
Expand All @@ -801,6 +877,21 @@ defmodule Console.Deployments.Workbenches do
end
end

@doc """
Queues a prompt for the job associated with a pull request.
"""
@spec pr_queued_prompt(map, binary, User.t()) :: queued_prompt_resp
def pr_queued_prompt(attrs, url, %User{} = user) do
Repo.get_by(PullRequest, url: url)
|> Repo.preload([:workbench_job])
|> case do
%PullRequest{workbench_job: %WorkbenchJob{} = job} ->
create_queued_prompt(attrs, job, user)
_ ->
{:error, "pull request not found"}
end
end

@doc """
Creates a new activity for a job, and bookkeeps job status and timestamp.
"""
Expand Down
57 changes: 57 additions & 0 deletions lib/console/graphql/deployments/workbench.ex
Original file line number Diff line number Diff line change
Expand Up @@ -418,6 +418,11 @@ defmodule Console.GraphQl.Deployments.Workbench do
field :prompt, non_null(:string), description: "the prompt for the message"
end

input_object :queued_prompt_attributes do
field :prompt, non_null(:string), description: "the prompt to send when dequeued"
field :dequeable_at, non_null(:datetime), description: "when this prompt becomes eligible to dequeue"
end

object :workbench do
field :id, non_null(:string), description: "the id of the workbench"
field :name, non_null(:string), description: "the name of the workbench"
Expand Down Expand Up @@ -536,6 +541,10 @@ defmodule Console.GraphQl.Deployments.Workbench do
resolve &Deployments.list_workbench_job_activities/3
end

connection field :queued_prompts, node_type: :queued_prompt do
resolve &Deployments.list_queued_prompts/3
end

field :metrics_tool, list_of(:workbench_job_activity_metric) do
arg :name, :string, description: "the name of the metrics tool"
arg :arguments, :json, description: "the arguments for the metrics tool"
Expand Down Expand Up @@ -817,6 +826,19 @@ defmodule Console.GraphQl.Deployments.Workbench do
timestamps()
end

object :queued_prompt do
field :id, non_null(:string), description: "the id of the queued prompt"
field :prompt, :string, description: "the prompt text"
field :dequeable_at, :datetime, description: "when this prompt becomes eligible to dequeue"
field :consumed_at, :datetime, description: "when this prompt was consumed"
field :user_id, :id, description: "user this prompt will run as"

field :workbench_job, :workbench_job, resolve: dataloader(Deployments), description: "the job this prompt will be sent to"
field :user, :user, resolve: dataloader(User), description: "the user who queued this prompt"

timestamps()
end

object :workbench_prompt do
field :id, non_null(:string), description: "the id of the saved prompt"
field :title, non_null(:string), description: "display title for the saved prompt", resolve: fn prompt, _, _ ->
Expand Down Expand Up @@ -1229,6 +1251,7 @@ defmodule Console.GraphQl.Deployments.Workbench do
connection node_type: :workbench_job
connection node_type: :workbench_job_activity
connection node_type: :workbench_job_thought
connection node_type: :queued_prompt
connection node_type: :workbench_cron
connection node_type: :workbench_prompt
connection node_type: :workbench_skill
Expand Down Expand Up @@ -1749,6 +1772,29 @@ defmodule Console.GraphQl.Deployments.Workbench do
resolve &Deployments.create_workbench_job/2
end

@desc "Queues a prompt to be sent to a workbench job later. Requires prompt access to the job."
field :create_queued_prompt, :queued_prompt do
middleware Authenticated
middleware Scope,
resource: :workbench,
action: :write
arg :job_id, non_null(:id), description: "the workbench job to queue a prompt for"
arg :attributes, non_null(:queued_prompt_attributes)

resolve &Deployments.create_queued_prompt/2
end

@desc "Deletes a queued prompt. Requires prompt access to the queued prompt's job."
field :delete_queued_prompt, :queued_prompt do
middleware Authenticated
middleware Scope,
resource: :workbench,
action: :write
arg :id, non_null(:id)

resolve &Deployments.delete_queued_prompt/2
end

field :create_workbench_message, :workbench_job_activity do
middleware Authenticated
middleware Scope,
Expand All @@ -1771,6 +1817,17 @@ defmodule Console.GraphQl.Deployments.Workbench do
resolve &Deployments.workbench_pr_followup/2
end

field :enqueue_workbench_pr_followup, :queued_prompt do
middleware Authenticated
middleware Scope,
resource: :workbench,
action: :write
arg :url, non_null(:string), description: "the pull request url to queue a follow-up prompt for"
arg :attributes, non_null(:queued_prompt_attributes), description: "queued prompt attributes"

resolve &Deployments.enqueue_workbench_pr_followup/2
end

@desc "Approves and invokes a pending workbench function activity. Requires read access to the job's workbench."
field :approve_workbench_job_activity, :workbench_job_activity do
middleware Authenticated
Expand Down
15 changes: 15 additions & 0 deletions lib/console/graphql/resolvers/deployments/workbench.ex
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ defmodule Console.GraphQl.Resolvers.Deployments.Workbench do
WorkbenchTool,
WorkbenchCron,
WorkbenchPrompt,
QueuedPrompt,
WorkbenchSkill,
WorkbenchEvalResult,
WorkbenchWebhook,
Expand Down Expand Up @@ -141,6 +142,11 @@ defmodule Console.GraphQl.Resolvers.Deployments.Workbench do
|> paginate(args)
end

def list_queued_prompts(job, args, _) do
QueuedPrompt.for_workbench_job(job.id)
|> paginate(args)
end

def all_workbench_alerts(args, %{context: %{current_user: user}}) do
Alert.for_user(user)
|> Alert.ordered()
Expand Down Expand Up @@ -276,6 +282,12 @@ defmodule Console.GraphQl.Resolvers.Deployments.Workbench do
def create_workbench_job(%{workbench_id: workbench_id, attributes: attrs}, %{context: %{current_user: user}}),
do: Workbenches.create_workbench_job(attrs, workbench_id, user)

def create_queued_prompt(%{job_id: job_id, attributes: attrs}, %{context: %{current_user: user}}),
do: Workbenches.create_queued_prompt(attrs, job_id, user)

def delete_queued_prompt(%{id: id}, %{context: %{current_user: user}}),
do: Workbenches.delete_queued_prompt(id, user)

def create_workbench_cron(%{workbench_id: workbench_id, attributes: attrs}, %{context: %{current_user: user}}),
do: Workbenches.create_workbench_cron(attrs, workbench_id, user)

Expand Down Expand Up @@ -350,6 +362,9 @@ defmodule Console.GraphQl.Resolvers.Deployments.Workbench do
def workbench_pr_followup(%{url: url, attributes: attrs}, %{context: %{current_user: user}}),
do: Workbenches.pr_followup(attrs, url, user)

def enqueue_workbench_pr_followup(%{url: url, attributes: attrs}, %{context: %{current_user: user}}),
do: Workbenches.pr_queued_prompt(attrs, url, user)

def approve_workbench_job_activity(%{id: id}, %{context: %{current_user: user}}),
do: Workbenches.approve_job_activity(id, user)

Expand Down
35 changes: 35 additions & 0 deletions lib/console/openapi/ai/workbench.ex
Original file line number Diff line number Diff line change
Expand Up @@ -198,3 +198,38 @@ defmodule Console.OpenAPI.AI.WorkbenchJobInput do
required: [:prompt]
}
end

defmodule Console.OpenAPI.AI.QueuedPrompt do
@moduledoc "OpenAPI schema for queued prompts."
use Console.OpenAPI.Base

defschema %{
type: :object,
title: "QueuedPrompt",
description: "A deferred prompt queued for a workbench job",
properties: timestamps(%{
id: string(description: "Unique identifier for the queued prompt"),
prompt: string(description: "The prompt text"),
dequeable_at: datetime(description: "When the prompt becomes eligible to dequeue"),
consumed_at: datetime(description: "When the prompt was consumed"),
workbench_job_id: string(description: "ID of the workbench job this prompt targets"),
user_id: string(description: "ID of the user this prompt runs as")
})
}
end

defmodule Console.OpenAPI.AI.QueuedPromptInput do
@moduledoc "OpenAPI schema for creating queued prompts."
use Console.OpenAPI.Base

defschema %{
type: :object,
title: "QueuedPromptInput",
description: "Input for creating a deferred workbench prompt",
properties: %{
prompt: string(description: "The prompt to send when dequeued"),
dequeable_at: datetime(description: "When the prompt becomes eligible to dequeue")
},
required: [:prompt, :dequeable_at]
}
end
7 changes: 7 additions & 0 deletions lib/console/pipelines/ai/queued_prompt/pipeline.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
defmodule Console.Pipelines.AI.QueuedPrompt.Pipeline do
use Console.Pipelines.Consumer
alias Console.Schema.QueuedPrompt
alias Console.Deployments.Workbenches

def handle_event(%QueuedPrompt{} = prompt), do: Workbenches.dequeue_prompt(prompt)
end
Loading
Loading