From 6e91d06a1b465c036109ad97ba0b079b79de9865 Mon Sep 17 00:00:00 2001 From: Kenneth Pernyer Date: Tue, 7 Jul 2026 15:00:24 +0200 Subject: [PATCH] feat(showcase): adopt crm-helm scenario from atelier (consolidation RFL-148; depends on runway-app-host) --- showcase/crm-helm/Cargo.toml | 40 + showcase/crm-helm/build.rs | 27 + .../proto/prio/common/v1/common.proto | 419 +++++++++ .../prio/conversations/v1/conversations.proto | 39 + .../proto/prio/documents/v1/documents.proto | 27 + .../crm-helm/proto/prio/facts/v1/facts.proto | 18 + .../proto/prio/metadata/v1/metadata.proto | 50 + .../prio/opportunities/v1/opportunities.proto | 39 + .../proto/prio/parties/v1/parties.proto | 65 ++ .../proto/prio/workflow/v1/workflow.proto | 25 + showcase/crm-helm/src/conversations.rs | 159 ++++ showcase/crm-helm/src/documents.rs | 119 +++ showcase/crm-helm/src/facts.rs | 94 ++ showcase/crm-helm/src/main.rs | 86 ++ showcase/crm-helm/src/metadata.rs | 159 ++++ showcase/crm-helm/src/opportunities.rs | 137 +++ showcase/crm-helm/src/parties.rs | 192 ++++ showcase/crm-helm/src/proto.rs | 62 ++ showcase/crm-helm/src/shared.rs | 711 ++++++++++++++ showcase/crm-helm/src/truths.rs | 1 + .../src/truths/evaluate_acquisition_target.rs | 476 ++++++++++ .../src/truths/match_renewal_context.rs | 798 ++++++++++++++++ .../src/truths/plan_outbound_campaign.rs | 879 ++++++++++++++++++ showcase/crm-helm/src/workbench.rs | 1 + showcase/crm-helm/src/workflow.rs | 113 +++ 25 files changed, 4736 insertions(+) create mode 100644 showcase/crm-helm/Cargo.toml create mode 100644 showcase/crm-helm/build.rs create mode 100644 showcase/crm-helm/proto/prio/common/v1/common.proto create mode 100644 showcase/crm-helm/proto/prio/conversations/v1/conversations.proto create mode 100644 showcase/crm-helm/proto/prio/documents/v1/documents.proto create mode 100644 showcase/crm-helm/proto/prio/facts/v1/facts.proto create mode 100644 showcase/crm-helm/proto/prio/metadata/v1/metadata.proto create mode 100644 showcase/crm-helm/proto/prio/opportunities/v1/opportunities.proto create mode 100644 showcase/crm-helm/proto/prio/parties/v1/parties.proto create mode 100644 showcase/crm-helm/proto/prio/workflow/v1/workflow.proto create mode 100644 showcase/crm-helm/src/conversations.rs create mode 100644 showcase/crm-helm/src/documents.rs create mode 100644 showcase/crm-helm/src/facts.rs create mode 100644 showcase/crm-helm/src/main.rs create mode 100644 showcase/crm-helm/src/metadata.rs create mode 100644 showcase/crm-helm/src/opportunities.rs create mode 100644 showcase/crm-helm/src/parties.rs create mode 100644 showcase/crm-helm/src/proto.rs create mode 100644 showcase/crm-helm/src/shared.rs create mode 100644 showcase/crm-helm/src/truths.rs create mode 100644 showcase/crm-helm/src/truths/evaluate_acquisition_target.rs create mode 100644 showcase/crm-helm/src/truths/match_renewal_context.rs create mode 100644 showcase/crm-helm/src/truths/plan_outbound_campaign.rs create mode 100644 showcase/crm-helm/src/workbench.rs create mode 100644 showcase/crm-helm/src/workflow.rs diff --git a/showcase/crm-helm/Cargo.toml b/showcase/crm-helm/Cargo.toml new file mode 100644 index 0000000..279ca50 --- /dev/null +++ b/showcase/crm-helm/Cargo.toml @@ -0,0 +1,40 @@ +[package] +name = "example-crm-helm-showcase" +version = "0.0.0" +edition = "2024" +publish = false + +[[bin]] +name = "crm-helm-showcase" +path = "src/main.rs" + +[dependencies] +anyhow.workspace = true +async-trait.workspace = true +chrono.workspace = true +serde.workspace = true +serde_json.workspace = true +tokio.workspace = true +tokio-stream = "0.1" +tracing.workspace = true +tracing-subscriber.workspace = true +uuid.workspace = true + +# gRPC / proto +tonic = { version = "0.12", features = ["transport"] } +prost = "0.13" +prost-types = "0.13" + +# Axum — not in workspace deps +axum = "0.8" + +# Cross-repo: Helms domain crates (no converge-* deps) +application-kernel = { path = "../../../bedrock-platform/helms/crates/application-kernel" } +application-storage = { path = "../../../bedrock-platform/helms/crates/application-storage" } + +# Cross-repo: Runtime Runway app host +runway-app-host = { path = "../../../runtime-runway/crates/runway-app-host" } +runway-storage = { path = "../../../runtime-runway/crates/runway-storage" } + +[build-dependencies] +tonic-build = "0.12" diff --git a/showcase/crm-helm/build.rs b/showcase/crm-helm/build.rs new file mode 100644 index 0000000..f3903f5 --- /dev/null +++ b/showcase/crm-helm/build.rs @@ -0,0 +1,27 @@ +/// Proto compilation for the CRM Helm scenario. +/// +/// Strategy: proto files are copied into scenarios/crm-helm/proto/ from +/// helms/proto/ rather than referenced by a cross-repo path. This decouples +/// the atelier-showcase workspace from the helms filesystem layout and avoids +/// fragile relative paths that cross workspace roots. The copy-and-own +/// approach is intentional — atelier-showcase is a showcase/demo layer, not +/// a production dependency of helms. +fn main() -> Result<(), Box> { + tonic_build::configure() + .build_server(true) + .build_client(false) + .compile_protos( + &[ + "proto/prio/common/v1/common.proto", + "proto/prio/parties/v1/parties.proto", + "proto/prio/opportunities/v1/opportunities.proto", + "proto/prio/conversations/v1/conversations.proto", + "proto/prio/documents/v1/documents.proto", + "proto/prio/workflow/v1/workflow.proto", + "proto/prio/facts/v1/facts.proto", + "proto/prio/metadata/v1/metadata.proto", + ], + &["proto"], + )?; + Ok(()) +} diff --git a/showcase/crm-helm/proto/prio/common/v1/common.proto b/showcase/crm-helm/proto/prio/common/v1/common.proto new file mode 100644 index 0000000..497c7cc --- /dev/null +++ b/showcase/crm-helm/proto/prio/common/v1/common.proto @@ -0,0 +1,419 @@ +syntax = "proto3"; + +package prio.common.v1; + +import "google/protobuf/timestamp.proto"; + +message Actor { + string actor_id = 1; + string display_name = 2; + ActorKind kind = 3; +} + +enum ActorKind { + ACTOR_KIND_UNSPECIFIED = 0; + ACTOR_KIND_HUMAN = 1; + ACTOR_KIND_AGENT = 2; + ACTOR_KIND_SYSTEM = 3; +} + +message RecordRef { + RecordKind kind = 1; + string record_id = 2; +} + +enum RecordKind { + RECORD_KIND_UNSPECIFIED = 0; + RECORD_KIND_ORGANIZATION = 1; + RECORD_KIND_PERSON = 2; + RECORD_KIND_RELATIONSHIP = 3; + RECORD_KIND_LEAD = 4; + RECORD_KIND_OPPORTUNITY = 5; + RECORD_KIND_CONVERSATION = 6; + RECORD_KIND_ACTIVITY = 7; + RECORD_KIND_TASK = 8; + RECORD_KIND_OFFER_QUOTE = 9; + RECORD_KIND_ORDER_SUBSCRIPTION = 10; + RECORD_KIND_DOCUMENT = 11; + RECORD_KIND_FACT = 12; + RECORD_KIND_INTENT = 13; + RECORD_KIND_WORKFLOW_CASE = 14; + RECORD_KIND_COMMUNICATION_EVENT = 15; + RECORD_KIND_PERMISSION_GRANT = 16; + RECORD_KIND_AUDIT_ENTRY = 17; + RECORD_KIND_NOTE = 18; + RECORD_KIND_CATALOG_ITEM = 19; +} + +message Money { + string currency_code = 1; + int64 amount_minor = 2; +} + +enum SubscriptionStatus { + SUBSCRIPTION_STATUS_UNSPECIFIED = 0; + SUBSCRIPTION_STATUS_DRAFT = 1; + SUBSCRIPTION_STATUS_PENDING_ACTIVATION = 2; + SUBSCRIPTION_STATUS_ACTIVE = 3; + SUBSCRIPTION_STATUS_SUSPENDED = 4; + SUBSCRIPTION_STATUS_CANCELLED = 5; +} + +enum LedgerEntryKind { + LEDGER_ENTRY_KIND_UNSPECIFIED = 0; + LEDGER_ENTRY_KIND_OPENING_BALANCE = 1; + LEDGER_ENTRY_KIND_CREDIT_GRANT = 2; + LEDGER_ENTRY_KIND_DEBIT = 3; + LEDGER_ENTRY_KIND_ADJUSTMENT = 4; +} + +message OrderSubscription { + string id = 1; + string organization_id = 2; + optional string quote_id = 3; + optional string catalog_item_id = 4; + SubscriptionStatus status = 5; + Money value = 6; + google.protobuf.Timestamp started_at = 7; + optional google.protobuf.Timestamp activated_at = 8; +} + +message EntitlementValue { + oneof kind { + bool feature_flag = 1; + int64 quota = 2; + int64 credits = 3; + string text = 4; + } +} + +message Entitlement { + string id = 1; + string organization_id = 2; + string subscription_id = 3; + string catalog_item_id = 4; + string key = 5; + EntitlementValue value = 6; + google.protobuf.Timestamp created_at = 7; +} + +message LedgerEntry { + string id = 1; + string organization_id = 2; + string subscription_id = 3; + LedgerEntryKind kind = 4; + Money amount = 5; + string description = 6; + google.protobuf.Timestamp created_at = 7; +} + +enum RelationshipType { + RELATIONSHIP_TYPE_UNSPECIFIED = 0; + RELATIONSHIP_TYPE_EMPLOYMENT = 1; + RELATIONSHIP_TYPE_CHAMPION = 2; + RELATIONSHIP_TYPE_DECISION_MAKER = 3; + RELATIONSHIP_TYPE_PARTNER = 4; + RELATIONSHIP_TYPE_COMPETITOR = 5; + RELATIONSHIP_TYPE_OTHER = 6; +} + +enum OrganizationLifecycle { + ORGANIZATION_LIFECYCLE_UNSPECIFIED = 0; + ORGANIZATION_LIFECYCLE_PROSPECT = 1; + ORGANIZATION_LIFECYCLE_ACTIVE = 2; + ORGANIZATION_LIFECYCLE_DORMANT = 3; + ORGANIZATION_LIFECYCLE_PARTNER = 4; +} + +enum OpportunityStage { + OPPORTUNITY_STAGE_UNSPECIFIED = 0; + OPPORTUNITY_STAGE_QUALIFYING = 1; + OPPORTUNITY_STAGE_DISCOVERY = 2; + OPPORTUNITY_STAGE_PROPOSAL = 3; + OPPORTUNITY_STAGE_NEGOTIATION = 4; + OPPORTUNITY_STAGE_CLOSED_WON = 5; + OPPORTUNITY_STAGE_CLOSED_LOST = 6; +} + +enum ActivityOutcome { + ACTIVITY_OUTCOME_UNSPECIFIED = 0; + ACTIVITY_OUTCOME_COMPLETED = 1; + ACTIVITY_OUTCOME_WAITING = 2; + ACTIVITY_OUTCOME_BLOCKED = 3; +} + +enum DocumentStatus { + DOCUMENT_STATUS_UNSPECIFIED = 0; + DOCUMENT_STATUS_DRAFT = 1; + DOCUMENT_STATUS_VERIFIED = 2; + DOCUMENT_STATUS_ARCHIVED = 3; +} + +enum CommunicationChannel { + COMMUNICATION_CHANNEL_UNSPECIFIED = 0; + COMMUNICATION_CHANNEL_EMAIL = 1; + COMMUNICATION_CHANNEL_PHONE = 2; + COMMUNICATION_CHANNEL_MEETING = 3; + COMMUNICATION_CHANNEL_CHAT = 4; + COMMUNICATION_CHANNEL_SMS = 5; +} + +enum CommunicationDirection { + COMMUNICATION_DIRECTION_UNSPECIFIED = 0; + COMMUNICATION_DIRECTION_INBOUND = 1; + COMMUNICATION_DIRECTION_OUTBOUND = 2; + COMMUNICATION_DIRECTION_INTERNAL = 3; +} + +enum WorkflowState { + WORKFLOW_STATE_UNSPECIFIED = 0; + WORKFLOW_STATE_OPEN = 1; + WORKFLOW_STATE_AWAITING_APPROVAL = 2; + WORKFLOW_STATE_WAITING_EXTERNAL = 3; + WORKFLOW_STATE_BLOCKED = 4; + WORKFLOW_STATE_DONE = 5; +} + +enum WorkflowPriority { + WORKFLOW_PRIORITY_UNSPECIFIED = 0; + WORKFLOW_PRIORITY_LOW = 1; + WORKFLOW_PRIORITY_MEDIUM = 2; + WORKFLOW_PRIORITY_HIGH = 3; + WORKFLOW_PRIORITY_CRITICAL = 4; +} + +enum TimelineEntryKind { + TIMELINE_ENTRY_KIND_UNSPECIFIED = 0; + TIMELINE_ENTRY_KIND_ACTIVITY = 1; + TIMELINE_ENTRY_KIND_NOTE = 2; + TIMELINE_ENTRY_KIND_DOCUMENT = 3; + TIMELINE_ENTRY_KIND_COMMUNICATION = 4; + TIMELINE_ENTRY_KIND_FACT = 5; + TIMELINE_ENTRY_KIND_AUDIT = 6; +} + +enum ObjectDefinitionKind { + OBJECT_DEFINITION_KIND_UNSPECIFIED = 0; + OBJECT_DEFINITION_KIND_STANDARD = 1; + OBJECT_DEFINITION_KIND_CUSTOM = 2; +} + +enum FieldType { + FIELD_TYPE_UNSPECIFIED = 0; + FIELD_TYPE_TEXT = 1; + FIELD_TYPE_LONG_TEXT = 2; + FIELD_TYPE_NUMBER = 3; + FIELD_TYPE_CURRENCY = 4; + FIELD_TYPE_BOOLEAN = 5; + FIELD_TYPE_DATE = 6; + FIELD_TYPE_DATE_TIME = 7; + FIELD_TYPE_EMAIL = 8; + FIELD_TYPE_PHONE = 9; + FIELD_TYPE_URL = 10; + FIELD_TYPE_SELECT = 11; + FIELD_TYPE_MULTI_SELECT = 12; + FIELD_TYPE_RELATION = 13; +} + +enum RelationshipCardinality { + RELATIONSHIP_CARDINALITY_UNSPECIFIED = 0; + RELATIONSHIP_CARDINALITY_ONE_TO_ONE = 1; + RELATIONSHIP_CARDINALITY_ONE_TO_MANY = 2; + RELATIONSHIP_CARDINALITY_MANY_TO_MANY = 3; +} + +enum ViewLayout { + VIEW_LAYOUT_UNSPECIFIED = 0; + VIEW_LAYOUT_TABLE = 1; + VIEW_LAYOUT_KANBAN = 2; + VIEW_LAYOUT_CALENDAR = 3; +} + +message Organization { + string id = 1; + string name = 2; + optional string external_key = 3; + optional string website = 4; + optional string industry = 5; + OrganizationLifecycle lifecycle = 6; + optional string owner_user_id = 7; + repeated string tags = 8; + google.protobuf.Timestamp created_at = 9; + google.protobuf.Timestamp updated_at = 10; +} + +message Person { + string id = 1; + optional string organization_id = 2; + string full_name = 3; + optional string title = 4; + optional string email = 5; + optional string phone = 6; + optional string linkedin_url = 7; + google.protobuf.Timestamp created_at = 8; + google.protobuf.Timestamp updated_at = 9; +} + +message Relationship { + string id = 1; + RecordRef from = 2; + RecordRef to = 3; + RelationshipType relationship_type = 4; + optional string label = 5; + google.protobuf.Timestamp created_at = 6; +} + +message Opportunity { + string id = 1; + string organization_id = 2; + optional string primary_contact_id = 3; + string name = 4; + OpportunityStage stage = 5; + Money value = 6; + uint32 confidence_bps = 7; + optional string next_step = 8; + optional google.protobuf.Timestamp expected_close_at = 9; + google.protobuf.Timestamp created_at = 10; + google.protobuf.Timestamp updated_at = 11; +} + +message Activity { + string id = 1; + string subject = 2; + string details = 3; + Actor actor = 4; + repeated RecordRef related_to = 5; + ActivityOutcome outcome = 6; + google.protobuf.Timestamp occurred_at = 7; + optional google.protobuf.Timestamp next_action_due_at = 8; +} + +message Note { + string id = 1; + string subject = 2; + string body = 3; + Actor author = 4; + repeated RecordRef related_to = 5; + bool promoted_to_fact = 6; + google.protobuf.Timestamp created_at = 7; +} + +message Document { + string id = 1; + string title = 2; + string media_type = 3; + string uri = 4; + DocumentStatus status = 5; + Actor uploaded_by = 6; + repeated RecordRef related_to = 7; + google.protobuf.Timestamp created_at = 8; +} + +message CommunicationEvent { + string id = 1; + CommunicationChannel channel = 2; + CommunicationDirection direction = 3; + optional string subject = 4; + string summary = 5; + string counterpart = 6; + Actor actor = 7; + repeated RecordRef related_to = 8; + google.protobuf.Timestamp occurred_at = 9; +} + +message WorkflowCase { + string id = 1; + string title = 2; + WorkflowState state = 3; + WorkflowPriority priority = 4; + optional string owner_user_id = 5; + repeated RecordRef related_to = 6; + google.protobuf.Timestamp opened_at = 7; + google.protobuf.Timestamp updated_at = 8; +} + +message Fact { + string id = 1; + string statement = 2; + uint32 confidence_bps = 3; + Actor promoted_by = 4; + optional string source_note_id = 5; + repeated RecordRef related_to = 6; + google.protobuf.Timestamp created_at = 7; +} + +message PermissionGrant { + string id = 1; + string subject = 2; + string role = 3; + string scope = 4; + Actor granted_by = 5; + google.protobuf.Timestamp created_at = 6; +} + +message TimelineEntry { + string id = 1; + TimelineEntryKind kind = 2; + optional RecordRef anchor = 3; + string headline = 4; + string body = 5; + Actor actor = 6; + google.protobuf.Timestamp occurred_at = 7; + repeated RecordRef related_to = 8; +} + +message AccountSummary { + Organization organization = 1; + repeated Person contacts = 2; + repeated Opportunity opportunities = 3; + repeated WorkflowCase workflow_cases = 4; + repeated Fact facts = 5; + repeated Document documents = 6; + repeated PermissionGrant permissions = 7; + repeated TimelineEntry recent_timeline = 8; +} + +message FieldDefinition { + string id = 1; + string key = 2; + string label = 3; + FieldType field_type = 4; + bool required = 5; + repeated string options = 6; + optional string relation_object_key = 7; + bool active = 8; +} + +message RelationshipDefinition { + string id = 1; + string target_object_key = 2; + RelationshipCardinality cardinality = 3; + string label = 4; +} + +message ObjectDefinition { + string id = 1; + string key = 2; + string display_name = 3; + ObjectDefinitionKind kind = 4; + repeated FieldDefinition fields = 5; + repeated RelationshipDefinition relationships = 6; + bool active = 7; + google.protobuf.Timestamp created_at = 8; + google.protobuf.Timestamp updated_at = 9; +} + +message ViewDefinition { + string id = 1; + string object_key = 2; + string name = 3; + ViewLayout layout = 4; + optional string filter_expression = 5; + optional string sort_expression = 6; + repeated string visible_fields = 7; + optional string group_by = 8; + bool favorite = 9; + optional string owner_user_id = 10; + google.protobuf.Timestamp created_at = 11; + google.protobuf.Timestamp updated_at = 12; +} diff --git a/showcase/crm-helm/proto/prio/conversations/v1/conversations.proto b/showcase/crm-helm/proto/prio/conversations/v1/conversations.proto new file mode 100644 index 0000000..2fbc870 --- /dev/null +++ b/showcase/crm-helm/proto/prio/conversations/v1/conversations.proto @@ -0,0 +1,39 @@ +syntax = "proto3"; + +package prio.conversations.v1; + +import "google/protobuf/timestamp.proto"; +import "prio/common/v1/common.proto"; + +service ConversationsService { + rpc AppendActivity(AppendActivityRequest) returns (prio.common.v1.Activity); + rpc RecordCommunication(RecordCommunicationRequest) returns (prio.common.v1.CommunicationEvent); + rpc StreamTimeline(StreamTimelineRequest) returns (stream prio.common.v1.TimelineEntry); +} + +message AppendActivityRequest { + string subject = 1; + string details = 2; + repeated prio.common.v1.RecordRef related_to = 3; + prio.common.v1.ActivityOutcome outcome = 4; + optional google.protobuf.Timestamp occurred_at = 5; + optional google.protobuf.Timestamp next_action_due_at = 6; + optional prio.common.v1.Actor actor = 7; +} + +message RecordCommunicationRequest { + prio.common.v1.CommunicationChannel channel = 1; + prio.common.v1.CommunicationDirection direction = 2; + optional string subject = 3; + string summary = 4; + string counterpart = 5; + repeated prio.common.v1.RecordRef related_to = 6; + optional google.protobuf.Timestamp occurred_at = 7; + optional prio.common.v1.Actor actor = 8; +} + +message StreamTimelineRequest { + repeated prio.common.v1.RecordRef anchors = 1; + uint32 limit = 2; +} + diff --git a/showcase/crm-helm/proto/prio/documents/v1/documents.proto b/showcase/crm-helm/proto/prio/documents/v1/documents.proto new file mode 100644 index 0000000..467c885 --- /dev/null +++ b/showcase/crm-helm/proto/prio/documents/v1/documents.proto @@ -0,0 +1,27 @@ +syntax = "proto3"; + +package prio.documents.v1; + +import "prio/common/v1/common.proto"; + +service DocumentsService { + rpc AppendNote(AppendNoteRequest) returns (prio.common.v1.Note); + rpc AttachDocument(AttachDocumentRequest) returns (prio.common.v1.Document); +} + +message AppendNoteRequest { + string subject = 1; + string body = 2; + repeated prio.common.v1.RecordRef related_to = 3; + optional prio.common.v1.Actor actor = 4; +} + +message AttachDocumentRequest { + string title = 1; + string media_type = 2; + string uri = 3; + prio.common.v1.DocumentStatus status = 4; + repeated prio.common.v1.RecordRef related_to = 5; + optional prio.common.v1.Actor actor = 6; +} + diff --git a/showcase/crm-helm/proto/prio/facts/v1/facts.proto b/showcase/crm-helm/proto/prio/facts/v1/facts.proto new file mode 100644 index 0000000..934c91a --- /dev/null +++ b/showcase/crm-helm/proto/prio/facts/v1/facts.proto @@ -0,0 +1,18 @@ +syntax = "proto3"; + +package prio.facts.v1; + +import "prio/common/v1/common.proto"; + +service FactsService { + rpc RecordFact(RecordFactRequest) returns (prio.common.v1.Fact); +} + +message RecordFactRequest { + string statement = 1; + uint32 confidence_bps = 2; + repeated prio.common.v1.RecordRef related_to = 3; + optional string source_note_id = 4; + optional prio.common.v1.Actor actor = 5; +} + diff --git a/showcase/crm-helm/proto/prio/metadata/v1/metadata.proto b/showcase/crm-helm/proto/prio/metadata/v1/metadata.proto new file mode 100644 index 0000000..8889e35 --- /dev/null +++ b/showcase/crm-helm/proto/prio/metadata/v1/metadata.proto @@ -0,0 +1,50 @@ +syntax = "proto3"; + +package prio.metadata.v1; + +import "prio/common/v1/common.proto"; + +service MetadataService { + rpc UpsertObjectDefinition(UpsertObjectDefinitionRequest) returns (prio.common.v1.ObjectDefinition); + rpc UpsertViewDefinition(UpsertViewDefinitionRequest) returns (prio.common.v1.ViewDefinition); + rpc ListObjectDefinitions(ListObjectDefinitionsRequest) returns (ListObjectDefinitionsResponse); + rpc ListViewDefinitions(ListViewDefinitionsRequest) returns (ListViewDefinitionsResponse); +} + +message UpsertObjectDefinitionRequest { + optional string object_definition_id = 1; + string key = 2; + string display_name = 3; + prio.common.v1.ObjectDefinitionKind kind = 4; + repeated prio.common.v1.FieldDefinition fields = 5; + repeated prio.common.v1.RelationshipDefinition relationships = 6; + bool active = 7; +} + +message UpsertViewDefinitionRequest { + optional string view_definition_id = 1; + string object_key = 2; + string name = 3; + prio.common.v1.ViewLayout layout = 4; + optional string filter_expression = 5; + optional string sort_expression = 6; + repeated string visible_fields = 7; + optional string group_by = 8; + bool favorite = 9; + optional string owner_user_id = 10; +} + +message ListObjectDefinitionsRequest {} + +message ListObjectDefinitionsResponse { + repeated prio.common.v1.ObjectDefinition objects = 1; +} + +message ListViewDefinitionsRequest { + optional string object_key = 1; +} + +message ListViewDefinitionsResponse { + repeated prio.common.v1.ViewDefinition views = 1; +} + diff --git a/showcase/crm-helm/proto/prio/opportunities/v1/opportunities.proto b/showcase/crm-helm/proto/prio/opportunities/v1/opportunities.proto new file mode 100644 index 0000000..cf902f8 --- /dev/null +++ b/showcase/crm-helm/proto/prio/opportunities/v1/opportunities.proto @@ -0,0 +1,39 @@ +syntax = "proto3"; + +package prio.opportunities.v1; + +import "google/protobuf/timestamp.proto"; +import "prio/common/v1/common.proto"; + +service OpportunitiesService { + rpc CreateOpportunity(CreateOpportunityRequest) returns (prio.common.v1.Opportunity); + rpc AdvanceOpportunityStage(AdvanceOpportunityStageRequest) returns (prio.common.v1.Opportunity); + rpc ListOpportunities(ListOpportunitiesRequest) returns (ListOpportunitiesResponse); +} + +message CreateOpportunityRequest { + string organization_id = 1; + optional string primary_contact_id = 2; + string name = 3; + prio.common.v1.Money value = 4; + uint32 confidence_bps = 5; + optional string next_step = 6; + optional google.protobuf.Timestamp expected_close_at = 7; + optional prio.common.v1.Actor actor = 8; +} + +message AdvanceOpportunityStageRequest { + string opportunity_id = 1; + prio.common.v1.OpportunityStage stage = 2; + optional string next_step = 3; + optional prio.common.v1.Actor actor = 4; +} + +message ListOpportunitiesRequest { + optional string organization_id = 1; +} + +message ListOpportunitiesResponse { + repeated prio.common.v1.Opportunity opportunities = 1; +} + diff --git a/showcase/crm-helm/proto/prio/parties/v1/parties.proto b/showcase/crm-helm/proto/prio/parties/v1/parties.proto new file mode 100644 index 0000000..a3fd468 --- /dev/null +++ b/showcase/crm-helm/proto/prio/parties/v1/parties.proto @@ -0,0 +1,65 @@ +syntax = "proto3"; + +package prio.parties.v1; + +import "prio/common/v1/common.proto"; + +service PartiesService { + rpc UpsertOrganization(UpsertOrganizationRequest) returns (prio.common.v1.Organization); + rpc UpsertPerson(UpsertPersonRequest) returns (prio.common.v1.Person); + rpc LinkRelationship(LinkRelationshipRequest) returns (prio.common.v1.Relationship); + rpc GetAccountSummary(GetAccountSummaryRequest) returns (prio.common.v1.AccountSummary); + rpc ListOrganizations(ListOrganizationsRequest) returns (ListOrganizationsResponse); + rpc ListPeople(ListPeopleRequest) returns (ListPeopleResponse); +} + +message UpsertOrganizationRequest { + optional string organization_id = 1; + string name = 2; + optional string external_key = 3; + optional string website = 4; + optional string industry = 5; + prio.common.v1.OrganizationLifecycle lifecycle = 6; + optional string owner_user_id = 7; + repeated string tags = 8; + optional prio.common.v1.Actor actor = 9; +} + +message UpsertPersonRequest { + optional string person_id = 1; + optional string organization_id = 2; + string full_name = 3; + optional string title = 4; + optional string email = 5; + optional string phone = 6; + optional string linkedin_url = 7; + optional prio.common.v1.Actor actor = 8; +} + +message LinkRelationshipRequest { + prio.common.v1.RecordRef from = 1; + prio.common.v1.RecordRef to = 2; + prio.common.v1.RelationshipType relationship_type = 3; + optional string label = 4; + optional prio.common.v1.Actor actor = 5; +} + +message GetAccountSummaryRequest { + string organization_id = 1; + uint32 timeline_limit = 2; +} + +message ListOrganizationsRequest {} + +message ListOrganizationsResponse { + repeated prio.common.v1.Organization organizations = 1; +} + +message ListPeopleRequest { + optional string organization_id = 1; +} + +message ListPeopleResponse { + repeated prio.common.v1.Person people = 1; +} + diff --git a/showcase/crm-helm/proto/prio/workflow/v1/workflow.proto b/showcase/crm-helm/proto/prio/workflow/v1/workflow.proto new file mode 100644 index 0000000..c9aeec7 --- /dev/null +++ b/showcase/crm-helm/proto/prio/workflow/v1/workflow.proto @@ -0,0 +1,25 @@ +syntax = "proto3"; + +package prio.workflow.v1; + +import "prio/common/v1/common.proto"; + +service WorkflowService { + rpc CreateWorkflowCase(CreateWorkflowCaseRequest) returns (prio.common.v1.WorkflowCase); + rpc AdvanceWorkflowCase(AdvanceWorkflowCaseRequest) returns (prio.common.v1.WorkflowCase); +} + +message CreateWorkflowCaseRequest { + string title = 1; + prio.common.v1.WorkflowPriority priority = 2; + optional string owner_user_id = 3; + repeated prio.common.v1.RecordRef related_to = 4; + optional prio.common.v1.Actor actor = 5; +} + +message AdvanceWorkflowCaseRequest { + string workflow_case_id = 1; + prio.common.v1.WorkflowState state = 2; + optional prio.common.v1.Actor actor = 3; +} + diff --git a/showcase/crm-helm/src/conversations.rs b/showcase/crm-helm/src/conversations.rs new file mode 100644 index 0000000..226cba7 --- /dev/null +++ b/showcase/crm-helm/src/conversations.rs @@ -0,0 +1,159 @@ +//! CRM Conversations module — communication threads as a HelmModule. +//! +//! Moved from helms/crates/application-server/src/service.rs (ConversationsGrpc). + +use application_kernel::{ActivityAppend, CommunicationRecord}; +use application_storage::{AppKernelStore, InMemoryKernelStore, KernelStore}; +use async_trait::async_trait; +use runway_app_host::HelmModule; +use tokio::sync::mpsc; +use tokio_stream::wrappers::ReceiverStream; +use tonic::{Request, Response, Status}; + +use crate::proto::{common as pb, conversations as conversations_pb}; +use crate::shared::{ + activity_outcome_from_proto, actor_from_proto, communication_channel_from_proto, + communication_direction_from_proto, datetime_from_proto, default_limit, proto_activity, + proto_communication_event, proto_timeline_entry, record_ref_from_proto, status_from_storage, +}; + +// --------------------------------------------------------------------------- +// gRPC service struct +// --------------------------------------------------------------------------- + +#[derive(Clone)] +pub struct ConversationsGrpc { + store: S, +} + +impl ConversationsGrpc { + #[allow(dead_code)] + pub fn new(store: S) -> Self { + Self { store } + } +} + +#[tonic::async_trait] +impl conversations_pb::conversations_service_server::ConversationsService + for ConversationsGrpc +where + S: KernelStore, +{ + type StreamTimelineStream = ReceiverStream>; + + async fn append_activity( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let related_to = request + .related_to + .into_iter() + .map(record_ref_from_proto) + .collect::, _>>()?; + let activity = self + .store + .write(|kernel| { + kernel.append_activity( + ActivityAppend { + subject: request.subject, + details: request.details, + related_to, + outcome: activity_outcome_from_proto(request.outcome), + occurred_at: request.occurred_at.and_then(datetime_from_proto), + next_action_due_at: request + .next_action_due_at + .and_then(datetime_from_proto), + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_activity(activity))) + } + + async fn record_communication( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let related_to = request + .related_to + .into_iter() + .map(record_ref_from_proto) + .collect::, _>>()?; + let event = self + .store + .write(|kernel| { + kernel.record_communication( + CommunicationRecord { + channel: communication_channel_from_proto(request.channel), + direction: communication_direction_from_proto(request.direction), + subject: request.subject, + summary: request.summary, + counterpart: request.counterpart, + related_to, + occurred_at: request.occurred_at.and_then(datetime_from_proto), + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_communication_event(event))) + } + + async fn stream_timeline( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let anchors = request + .anchors + .into_iter() + .map(record_ref_from_proto) + .collect::, _>>()?; + let entries = self + .store + .read(|kernel| kernel.list_timeline(&anchors, default_limit(request.limit, 50))) + .map_err(status_from_storage)?; + + let (tx, rx) = mpsc::channel(16); + tokio::spawn(async move { + for entry in entries { + if tx.send(Ok(proto_timeline_entry(entry))).await.is_err() { + break; + } + } + }); + + Ok(Response::new(ReceiverStream::new(rx))) + } +} + +// --------------------------------------------------------------------------- +// HelmModule wrapper +// --------------------------------------------------------------------------- + +pub struct ConversationsModule {} + +impl ConversationsModule { + pub fn new(_store: AppKernelStore) -> Self { + Self {} + } + + #[allow(dead_code)] + pub fn in_memory() -> Self { + Self::new(AppKernelStore::Memory(InMemoryKernelStore::default_local())) + } +} + +#[async_trait] +impl HelmModule for ConversationsModule { + fn module_id(&self) -> &'static str { + "crm.conversations" + } + + async fn init(&self) -> anyhow::Result<()> { + Ok(()) + } +} diff --git a/showcase/crm-helm/src/documents.rs b/showcase/crm-helm/src/documents.rs new file mode 100644 index 0000000..51ebd61 --- /dev/null +++ b/showcase/crm-helm/src/documents.rs @@ -0,0 +1,119 @@ +//! CRM Documents module — attachments as a HelmModule. +//! +//! Moved from helms/crates/application-server/src/service.rs (DocumentsGrpc). + +use application_kernel::{DocumentAttach, NoteAppend}; +use application_storage::{AppKernelStore, InMemoryKernelStore, KernelStore}; +use async_trait::async_trait; +use runway_app_host::HelmModule; +use tonic::{Request, Response, Status}; + +use crate::proto::{common as pb, documents as documents_pb}; +use crate::shared::{ + actor_from_proto, document_status_from_proto, proto_document, proto_note, + record_ref_from_proto, status_from_storage, +}; + +// --------------------------------------------------------------------------- +// gRPC service struct +// --------------------------------------------------------------------------- + +#[derive(Clone)] +pub struct DocumentsGrpc { + store: S, +} + +impl DocumentsGrpc { + #[allow(dead_code)] + pub fn new(store: S) -> Self { + Self { store } + } +} + +#[tonic::async_trait] +impl documents_pb::documents_service_server::DocumentsService for DocumentsGrpc +where + S: KernelStore, +{ + async fn append_note( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let related_to = request + .related_to + .into_iter() + .map(record_ref_from_proto) + .collect::, _>>()?; + let note = self + .store + .write(|kernel| { + kernel.append_note( + NoteAppend { + subject: request.subject, + body: request.body, + related_to, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_note(note))) + } + + async fn attach_document( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let related_to = request + .related_to + .into_iter() + .map(record_ref_from_proto) + .collect::, _>>()?; + let document = self + .store + .write(|kernel| { + kernel.attach_document( + DocumentAttach { + title: request.title, + media_type: request.media_type, + uri: request.uri, + status: document_status_from_proto(request.status), + related_to, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_document(document))) + } +} + +// --------------------------------------------------------------------------- +// HelmModule wrapper +// --------------------------------------------------------------------------- + +pub struct DocumentsModule {} + +impl DocumentsModule { + pub fn new(_store: AppKernelStore) -> Self { + Self {} + } + + #[allow(dead_code)] + pub fn in_memory() -> Self { + Self::new(AppKernelStore::Memory(InMemoryKernelStore::default_local())) + } +} + +#[async_trait] +impl HelmModule for DocumentsModule { + fn module_id(&self) -> &'static str { + "crm.documents" + } + + async fn init(&self) -> anyhow::Result<()> { + Ok(()) + } +} diff --git a/showcase/crm-helm/src/facts.rs b/showcase/crm-helm/src/facts.rs new file mode 100644 index 0000000..7b29337 --- /dev/null +++ b/showcase/crm-helm/src/facts.rs @@ -0,0 +1,94 @@ +//! CRM Facts module — immutable audit log as a HelmModule. +//! +//! Moved from helms/crates/application-server/src/service.rs (FactsGrpc). + +use application_kernel::FactRecord; +use application_storage::{AppKernelStore, InMemoryKernelStore, KernelStore}; +use async_trait::async_trait; +use runway_app_host::HelmModule; +use tonic::{Request, Response, Status}; + +use crate::proto::{common as pb, facts as facts_pb}; +use crate::shared::{ + actor_from_proto, clamp_bps, parse_optional_uuid, proto_fact, record_ref_from_proto, + status_from_storage, +}; + +// --------------------------------------------------------------------------- +// gRPC service struct +// --------------------------------------------------------------------------- + +#[derive(Clone)] +pub struct FactsGrpc { + store: S, +} + +impl FactsGrpc { + #[allow(dead_code)] + pub fn new(store: S) -> Self { + Self { store } + } +} + +#[tonic::async_trait] +impl facts_pb::facts_service_server::FactsService for FactsGrpc +where + S: KernelStore, +{ + async fn record_fact( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let related_to = request + .related_to + .into_iter() + .map(record_ref_from_proto) + .collect::, _>>()?; + let confidence_bps = clamp_bps(request.confidence_bps)?; + let source_note_id = parse_optional_uuid(request.source_note_id)?; + let fact = self + .store + .write(|kernel| { + kernel.record_fact( + FactRecord { + statement: request.statement, + confidence_bps, + related_to, + source_note_id, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_fact(fact))) + } +} + +// --------------------------------------------------------------------------- +// HelmModule wrapper +// --------------------------------------------------------------------------- + +pub struct FactsModule {} + +impl FactsModule { + pub fn new(_store: AppKernelStore) -> Self { + Self {} + } + + #[allow(dead_code)] + pub fn in_memory() -> Self { + Self::new(AppKernelStore::Memory(InMemoryKernelStore::default_local())) + } +} + +#[async_trait] +impl HelmModule for FactsModule { + fn module_id(&self) -> &'static str { + "crm.facts" + } + + async fn init(&self) -> anyhow::Result<()> { + Ok(()) + } +} diff --git a/showcase/crm-helm/src/main.rs b/showcase/crm-helm/src/main.rs new file mode 100644 index 0000000..44e1fe6 --- /dev/null +++ b/showcase/crm-helm/src/main.rs @@ -0,0 +1,86 @@ +#![allow(clippy::result_large_err)] + +//! CRM Helm Showcase +//! +//! Demonstrates how Runtime Runway + Helm modules compose into a thin app binary. +//! Phase 6b wires the 7 CRM gRPC modules extracted from helms/application-server. + +mod conversations; +mod documents; +mod facts; +mod metadata; +mod opportunities; +mod parties; +mod proto; +mod shared; +mod truths; +mod workbench; +mod workflow; + +use std::sync::Arc; + +use application_storage::{AppKernelStore, InMemoryKernelStore}; +use runway_app_host::{ + AppExecutionPacket, BoundaryRegistration, BoundaryStatus, ContractLayer, MountKind, + MountedModule, RunwayAppHost, +}; +use runway_storage::StorageKit; + +#[tokio::main] +async fn main() -> anyhow::Result<()> { + tracing_subscriber::fmt::init(); + + let packet = AppExecutionPacket::new( + "crm-helm", + "CRM Helm Showcase", + "CRM gRPC services composed via HelmModule — Phase 6b showcase", + "/crm", + ) + .with_mounted_module(MountedModule::new("crm.parties", MountKind::Mounted)) + .with_mounted_module(MountedModule::new("crm.opportunities", MountKind::Mounted)) + .with_mounted_module(MountedModule::new("crm.conversations", MountKind::Mounted)) + .with_mounted_module(MountedModule::new("crm.documents", MountKind::Mounted)) + .with_mounted_module(MountedModule::new("crm.workflow", MountKind::Mounted)) + .with_mounted_module(MountedModule::new("crm.facts", MountKind::Mounted)) + .with_mounted_module(MountedModule::new("crm.metadata", MountKind::Mounted)) + .with_boundary(BoundaryRegistration::new( + ContractLayer::Helm, + vec![ + "crm.parties".to_string(), + "crm.opportunities".to_string(), + "crm.conversations".to_string(), + "crm.documents".to_string(), + "crm.workflow".to_string(), + "crm.facts".to_string(), + "crm.metadata".to_string(), + ], + BoundaryStatus::Mounted, + )); + + // Shared in-memory kernel store — all 7 modules share one store instance + // so writes in one service are visible to reads in another. + let store = AppKernelStore::Memory(InMemoryKernelStore::default_local()); + + let storage = StorageKit::local("crm-helm-local.db").await?; + + let host = RunwayAppHost::builder(packet) + .with_storage(storage) + .mount(Arc::new(parties::PartiesModule::new(store.clone()))) + .mount(Arc::new(opportunities::OpportunitiesModule::new( + store.clone(), + ))) + .mount(Arc::new(conversations::ConversationsModule::new( + store.clone(), + ))) + .mount(Arc::new(documents::DocumentsModule::new(store.clone()))) + .mount(Arc::new(workflow::WorkflowModule::new(store.clone()))) + .mount(Arc::new(facts::FactsModule::new(store.clone()))) + .mount(Arc::new(metadata::MetadataModule::new(store))) + .build() + .await?; + + // TODO(Phase 9/truth-execution): mount TruthCatalog module once + // feat/helm-truth-execution is merged. + + host.serve().await +} diff --git a/showcase/crm-helm/src/metadata.rs b/showcase/crm-helm/src/metadata.rs new file mode 100644 index 0000000..de7f16c --- /dev/null +++ b/showcase/crm-helm/src/metadata.rs @@ -0,0 +1,159 @@ +//! CRM Metadata module — schema definitions as a HelmModule. +//! +//! Moved from helms/crates/application-server/src/service.rs (MetadataGrpc). + +use application_kernel::{Actor, CrmKernel, ObjectDefinitionUpsert, ViewDefinitionUpsert}; +use application_storage::{AppKernelStore, InMemoryKernelStore, KernelStore}; +use async_trait::async_trait; +use runway_app_host::HelmModule; +use tonic::{Request, Response, Status}; + +use crate::proto::{common as pb, metadata as metadata_pb}; +use crate::shared::{ + field_definition_from_proto, object_definition_kind_from_proto, parse_optional_uuid, + proto_object_definition, proto_view_definition, relationship_definition_from_proto, + status_from_storage, view_layout_from_proto, +}; + +// --------------------------------------------------------------------------- +// gRPC service struct +// --------------------------------------------------------------------------- + +#[derive(Clone)] +pub struct MetadataGrpc { + store: S, +} + +impl MetadataGrpc { + #[allow(dead_code)] + pub fn new(store: S) -> Self { + Self { store } + } +} + +#[tonic::async_trait] +impl metadata_pb::metadata_service_server::MetadataService for MetadataGrpc +where + S: KernelStore, +{ + async fn upsert_object_definition( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let object_definition_id = parse_optional_uuid(request.object_definition_id)?; + let fields = request + .fields + .into_iter() + .map(field_definition_from_proto) + .collect::, _>>()?; + let relationships = request + .relationships + .into_iter() + .map(relationship_definition_from_proto) + .collect::, _>>()?; + let definition = self + .store + .write(|kernel| { + kernel.upsert_object_definition( + ObjectDefinitionUpsert { + object_definition_id, + key: request.key, + display_name: request.display_name, + kind: object_definition_kind_from_proto(request.kind), + fields, + relationships, + active: request.active, + }, + Actor::system(), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_object_definition(definition))) + } + + async fn upsert_view_definition( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let view_definition_id = parse_optional_uuid(request.view_definition_id)?; + let view = self + .store + .write(|kernel| { + kernel.upsert_view_definition( + ViewDefinitionUpsert { + view_definition_id, + object_key: request.object_key, + name: request.name, + layout: view_layout_from_proto(request.layout), + filter_expression: request.filter_expression, + sort_expression: request.sort_expression, + visible_fields: request.visible_fields, + group_by: request.group_by, + favorite: request.favorite, + owner_user_id: request.owner_user_id, + }, + Actor::system(), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_view_definition(view))) + } + + async fn list_object_definitions( + &self, + _request: Request, + ) -> Result, Status> { + let objects = self + .store + .read(CrmKernel::list_object_definitions) + .map_err(status_from_storage)?; + Ok(Response::new(metadata_pb::ListObjectDefinitionsResponse { + objects: objects.into_iter().map(proto_object_definition).collect(), + })) + } + + async fn list_view_definitions( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let object_key = request.object_key.as_deref(); + let views = self + .store + .read(|kernel| kernel.list_view_definitions(object_key)) + .map_err(status_from_storage)?; + Ok(Response::new(metadata_pb::ListViewDefinitionsResponse { + views: views.into_iter().map(proto_view_definition).collect(), + })) + } +} + +// --------------------------------------------------------------------------- +// HelmModule wrapper +// --------------------------------------------------------------------------- + +pub struct MetadataModule {} + +impl MetadataModule { + pub fn new(_store: AppKernelStore) -> Self { + Self {} + } + + #[allow(dead_code)] + pub fn in_memory() -> Self { + Self::new(AppKernelStore::Memory(InMemoryKernelStore::default_local())) + } +} + +#[async_trait] +impl HelmModule for MetadataModule { + fn module_id(&self) -> &'static str { + "crm.metadata" + } + + async fn init(&self) -> anyhow::Result<()> { + Ok(()) + } +} diff --git a/showcase/crm-helm/src/opportunities.rs b/showcase/crm-helm/src/opportunities.rs new file mode 100644 index 0000000..804befa --- /dev/null +++ b/showcase/crm-helm/src/opportunities.rs @@ -0,0 +1,137 @@ +//! CRM Opportunities module — sales pipeline as a HelmModule. +//! +//! Moved from helms/crates/application-server/src/service.rs (OpportunitiesGrpc). + +use application_kernel::{Money, OpportunityAdvance, OpportunityCreate}; +use application_storage::{AppKernelStore, InMemoryKernelStore, KernelStore}; +use async_trait::async_trait; +use runway_app_host::HelmModule; +use tonic::{Request, Response, Status}; + +use crate::proto::{common as pb, opportunities as opportunities_pb}; +use crate::shared::{ + actor_from_proto, clamp_bps, datetime_from_proto, opportunity_stage_from_proto, + parse_optional_uuid, parse_uuid, proto_opportunity, status_from_storage, +}; + +// --------------------------------------------------------------------------- +// gRPC service struct +// --------------------------------------------------------------------------- + +#[derive(Clone)] +pub struct OpportunitiesGrpc { + store: S, +} + +impl OpportunitiesGrpc { + #[allow(dead_code)] + pub fn new(store: S) -> Self { + Self { store } + } +} + +#[tonic::async_trait] +impl opportunities_pb::opportunities_service_server::OpportunitiesService + for OpportunitiesGrpc +where + S: KernelStore, +{ + async fn create_opportunity( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let value = request + .value + .ok_or_else(|| Status::invalid_argument("value is required"))?; + let organization_id = parse_uuid(&request.organization_id)?; + let primary_contact_id = parse_optional_uuid(request.primary_contact_id)?; + let confidence_bps = clamp_bps(request.confidence_bps)?; + let expected_close_at = request.expected_close_at.and_then(datetime_from_proto); + let opportunity = self + .store + .write(|kernel| { + kernel.create_opportunity( + OpportunityCreate { + organization_id, + primary_contact_id, + name: request.name, + value: Money { + currency_code: value.currency_code, + amount_minor: value.amount_minor, + }, + confidence_bps, + next_step: request.next_step, + expected_close_at, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_opportunity(opportunity))) + } + + async fn advance_opportunity_stage( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let opportunity_id = parse_uuid(&request.opportunity_id)?; + let opportunity = self + .store + .write(|kernel| { + kernel.advance_opportunity( + OpportunityAdvance { + opportunity_id, + stage: opportunity_stage_from_proto(request.stage), + next_step: request.next_step, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_opportunity(opportunity))) + } + + async fn list_opportunities( + &self, + request: Request, + ) -> Result, Status> { + let organization_id = parse_optional_uuid(request.into_inner().organization_id)?; + let opportunities = self + .store + .read(|kernel| kernel.list_opportunities(organization_id)) + .map_err(status_from_storage)?; + Ok(Response::new(opportunities_pb::ListOpportunitiesResponse { + opportunities: opportunities.into_iter().map(proto_opportunity).collect(), + })) + } +} + +// --------------------------------------------------------------------------- +// HelmModule wrapper +// --------------------------------------------------------------------------- + +pub struct OpportunitiesModule {} + +impl OpportunitiesModule { + pub fn new(_store: AppKernelStore) -> Self { + Self {} + } + + #[allow(dead_code)] + pub fn in_memory() -> Self { + Self::new(AppKernelStore::Memory(InMemoryKernelStore::default_local())) + } +} + +#[async_trait] +impl HelmModule for OpportunitiesModule { + fn module_id(&self) -> &'static str { + "crm.opportunities" + } + + async fn init(&self) -> anyhow::Result<()> { + Ok(()) + } +} diff --git a/showcase/crm-helm/src/parties.rs b/showcase/crm-helm/src/parties.rs new file mode 100644 index 0000000..3bd826d --- /dev/null +++ b/showcase/crm-helm/src/parties.rs @@ -0,0 +1,192 @@ +//! CRM Parties module — accounts/contacts CRUD as a HelmModule. +//! +//! Moved from helms/crates/application-server/src/service.rs (PartiesGrpc). + +use application_kernel::{CrmKernel, OrganizationUpsert, PersonUpsert, RelationshipLink}; +use application_storage::{AppKernelStore, InMemoryKernelStore, KernelStore}; +use async_trait::async_trait; +use runway_app_host::HelmModule; +use tonic::{Request, Response, Status}; + +use crate::proto::{common as pb, parties as parties_pb}; +use crate::shared::{ + actor_from_proto, default_limit, organization_lifecycle_from_proto, parse_optional_uuid, + parse_uuid, proto_account_summary, proto_organization, proto_person, proto_relationship, + record_ref_from_proto, relationship_type_from_proto, status_from_storage, +}; + +// --------------------------------------------------------------------------- +// gRPC service struct +// --------------------------------------------------------------------------- + +#[derive(Clone)] +pub struct PartiesGrpc { + store: S, +} + +impl PartiesGrpc { + #[allow(dead_code)] + pub fn new(store: S) -> Self { + Self { store } + } +} + +#[tonic::async_trait] +impl parties_pb::parties_service_server::PartiesService for PartiesGrpc +where + S: KernelStore, +{ + async fn upsert_organization( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let organization_id = parse_optional_uuid(request.organization_id)?; + let organization = self + .store + .write(|kernel| { + kernel.upsert_organization( + OrganizationUpsert { + organization_id, + name: request.name, + external_key: request.external_key, + website: request.website, + industry: request.industry, + lifecycle: organization_lifecycle_from_proto(request.lifecycle), + owner_user_id: request.owner_user_id, + tags: request.tags, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_organization(organization))) + } + + async fn upsert_person( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let person_id = parse_optional_uuid(request.person_id)?; + let organization_id = parse_optional_uuid(request.organization_id)?; + let person = self + .store + .write(|kernel| { + kernel.upsert_person( + PersonUpsert { + person_id, + organization_id, + full_name: request.full_name, + title: request.title, + email: request.email, + phone: request.phone, + linkedin_url: request.linkedin_url, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_person(person))) + } + + async fn link_relationship( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let from = request + .from + .ok_or_else(|| Status::invalid_argument("from is required")) + .and_then(record_ref_from_proto)?; + let to = request + .to + .ok_or_else(|| Status::invalid_argument("to is required")) + .and_then(record_ref_from_proto)?; + let relationship = self + .store + .write(|kernel| { + kernel.link_relationship( + RelationshipLink { + from, + to, + relationship_type: relationship_type_from_proto(request.relationship_type), + label: request.label, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_relationship(relationship))) + } + + async fn get_account_summary( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let organization_id = parse_uuid(&request.organization_id)?; + let timeline_limit = default_limit(request.timeline_limit, 25); + let summary = self + .store + .read(|kernel| kernel.get_account_summary(organization_id, timeline_limit)) + .map_err(status_from_storage)? + .map_err(crate::shared::status_from_kernel)?; + Ok(Response::new(proto_account_summary(summary))) + } + + async fn list_organizations( + &self, + _request: Request, + ) -> Result, Status> { + let organizations = self + .store + .read(CrmKernel::list_organizations) + .map_err(status_from_storage)?; + Ok(Response::new(parties_pb::ListOrganizationsResponse { + organizations: organizations.into_iter().map(proto_organization).collect(), + })) + } + + async fn list_people( + &self, + request: Request, + ) -> Result, Status> { + let organization_id = parse_optional_uuid(request.into_inner().organization_id)?; + let people = self + .store + .read(|kernel| kernel.list_people(organization_id)) + .map_err(status_from_storage)?; + Ok(Response::new(parties_pb::ListPeopleResponse { + people: people.into_iter().map(proto_person).collect(), + })) + } +} + +// --------------------------------------------------------------------------- +// HelmModule wrapper +// --------------------------------------------------------------------------- + +pub struct PartiesModule {} + +impl PartiesModule { + pub fn new(_store: AppKernelStore) -> Self { + Self {} + } + + #[allow(dead_code)] + pub fn in_memory() -> Self { + Self::new(AppKernelStore::Memory(InMemoryKernelStore::default_local())) + } +} + +#[async_trait] +impl HelmModule for PartiesModule { + fn module_id(&self) -> &'static str { + "crm.parties" + } + + async fn init(&self) -> anyhow::Result<()> { + Ok(()) + } +} diff --git a/showcase/crm-helm/src/proto.rs b/showcase/crm-helm/src/proto.rs new file mode 100644 index 0000000..c1dec2c --- /dev/null +++ b/showcase/crm-helm/src/proto.rs @@ -0,0 +1,62 @@ +#![allow(dead_code)] +#![allow(clippy::all)] +#![allow(clippy::pedantic)] + +pub mod prio { + pub mod common { + pub mod v1 { + tonic::include_proto!("prio.common.v1"); + } + } + + pub mod parties { + pub mod v1 { + tonic::include_proto!("prio.parties.v1"); + } + } + + pub mod opportunities { + pub mod v1 { + tonic::include_proto!("prio.opportunities.v1"); + } + } + + pub mod conversations { + pub mod v1 { + tonic::include_proto!("prio.conversations.v1"); + } + } + + pub mod documents { + pub mod v1 { + tonic::include_proto!("prio.documents.v1"); + } + } + + pub mod workflow { + pub mod v1 { + tonic::include_proto!("prio.workflow.v1"); + } + } + + pub mod facts { + pub mod v1 { + tonic::include_proto!("prio.facts.v1"); + } + } + + pub mod metadata { + pub mod v1 { + tonic::include_proto!("prio.metadata.v1"); + } + } +} + +pub use prio::common::v1 as common; +pub use prio::conversations::v1 as conversations; +pub use prio::documents::v1 as documents; +pub use prio::facts::v1 as facts; +pub use prio::metadata::v1 as metadata; +pub use prio::opportunities::v1 as opportunities; +pub use prio::parties::v1 as parties; +pub use prio::workflow::v1 as workflow; diff --git a/showcase/crm-helm/src/shared.rs b/showcase/crm-helm/src/shared.rs new file mode 100644 index 0000000..a191b1e --- /dev/null +++ b/showcase/crm-helm/src/shared.rs @@ -0,0 +1,711 @@ +//! Shared conversion helpers for proto ↔ kernel types. +//! +//! Extracted from helms/crates/application-server/src/service.rs and adapted +//! to reference the local proto module. + +use application_kernel::{ + ActivityOutcome, Actor, ActorKind, CommunicationChannel, CommunicationDirection, + DocumentStatus, FieldDefinition, FieldType, ObjectDefinitionKind, OpportunityStage, + OrganizationLifecycle, RecordKind, RecordRef, RelationshipCardinality, RelationshipDefinition, + ViewLayout, WorkflowPriority, WorkflowState, +}; +use application_storage::StorageError; +use chrono::{DateTime, Utc}; +use prost_types::Timestamp; +use tonic::Status; +use uuid::Uuid; + +use crate::proto::common as pb; + +// --------------------------------------------------------------------------- +// Error mapping +// --------------------------------------------------------------------------- + +pub fn status_from_storage(error: StorageError) -> Status { + match error { + StorageError::LockPoisoned => Status::internal("storage lock poisoned"), + StorageError::Kernel(error) => status_from_kernel(error), + StorageError::ConnectionFailed { backend, message } => { + Status::unavailable(format!("{backend} connection failed: {message}")) + } + StorageError::SerializationFailed { message } => Status::internal(message), + StorageError::Timeout { operation } => Status::deadline_exceeded(operation), + StorageError::RuntimeStore { message } => Status::internal(message), + } +} + +pub fn status_from_kernel(error: application_kernel::KernelError) -> Status { + match error { + application_kernel::KernelError::Validation(message) => Status::invalid_argument(message), + application_kernel::KernelError::NotFound { kind, id } => { + Status::not_found(format!("{kind} not found: {id}")) + } + application_kernel::KernelError::Invariant(message) => Status::failed_precondition(message), + application_kernel::KernelError::Conflict(message) => Status::already_exists(message), + } +} + +// --------------------------------------------------------------------------- +// UUID helpers +// --------------------------------------------------------------------------- + +pub fn parse_uuid(value: &str) -> Result { + Uuid::parse_str(value).map_err(|_| Status::invalid_argument(format!("invalid uuid: {value}"))) +} + +pub fn parse_optional_uuid(value: Option) -> Result, Status> { + value + .and_then(|v| { + let trimmed = v.trim().to_string(); + (!trimmed.is_empty()).then_some(trimmed) + }) + .map(|v| parse_uuid(&v)) + .transpose() +} + +pub fn clamp_bps(value: u32) -> Result { + u16::try_from(value) + .map_err(|_| Status::invalid_argument("bps value is out of range")) + .and_then(|v| { + if v > 10_000 { + Err(Status::invalid_argument( + "bps value must be between 0 and 10000", + )) + } else { + Ok(v) + } + }) +} + +pub fn default_limit(value: u32, fallback: usize) -> usize { + if value == 0 { fallback } else { value as usize } +} + +// --------------------------------------------------------------------------- +// Timestamp helpers +// --------------------------------------------------------------------------- + +pub fn proto_timestamp(value: DateTime) -> Option { + Some(Timestamp { + seconds: value.timestamp(), + nanos: value.timestamp_subsec_nanos() as i32, + }) +} + +pub fn datetime_from_proto(value: Timestamp) -> Option> { + DateTime::from_timestamp(value.seconds, value.nanos as u32) +} + +// --------------------------------------------------------------------------- +// Actor +// --------------------------------------------------------------------------- + +pub fn actor_from_proto(actor: Option) -> Actor { + actor.map_or_else(Actor::system, |actor| Actor { + actor_id: actor.actor_id, + display_name: actor.display_name, + kind: match pb::ActorKind::try_from(actor.kind).unwrap_or(pb::ActorKind::System) { + pb::ActorKind::Human => ActorKind::Human, + pb::ActorKind::Agent => ActorKind::Agent, + pb::ActorKind::System | pb::ActorKind::Unspecified => ActorKind::System, + }, + }) +} + +pub fn proto_actor(actor: Actor) -> pb::Actor { + pb::Actor { + actor_id: actor.actor_id, + display_name: actor.display_name, + kind: match actor.kind { + ActorKind::Human => pb::ActorKind::Human as i32, + ActorKind::Agent => pb::ActorKind::Agent as i32, + ActorKind::System => pb::ActorKind::System as i32, + }, + } +} + +// --------------------------------------------------------------------------- +// RecordRef +// --------------------------------------------------------------------------- + +pub fn record_ref_from_proto(reference: pb::RecordRef) -> Result { + Ok(RecordRef { + kind: match pb::RecordKind::try_from(reference.kind) { + Ok(pb::RecordKind::Organization) => RecordKind::Organization, + Ok(pb::RecordKind::Person) => RecordKind::Person, + Ok(pb::RecordKind::Relationship) => RecordKind::Relationship, + Ok(pb::RecordKind::Lead) => RecordKind::Lead, + Ok(pb::RecordKind::Opportunity) => RecordKind::Opportunity, + Ok(pb::RecordKind::Conversation) => RecordKind::Conversation, + Ok(pb::RecordKind::Activity) => RecordKind::Activity, + Ok(pb::RecordKind::Task) => RecordKind::Task, + Ok(pb::RecordKind::OfferQuote) => RecordKind::OfferQuote, + Ok(pb::RecordKind::OrderSubscription) => RecordKind::OrderSubscription, + Ok(pb::RecordKind::Document) => RecordKind::Document, + Ok(pb::RecordKind::Fact) => RecordKind::Fact, + Ok(pb::RecordKind::Intent) => RecordKind::Intent, + Ok(pb::RecordKind::WorkflowCase) => RecordKind::WorkflowCase, + Ok(pb::RecordKind::CommunicationEvent) => RecordKind::CommunicationEvent, + Ok(pb::RecordKind::PermissionGrant) => RecordKind::PermissionGrant, + Ok(pb::RecordKind::AuditEntry) => RecordKind::AuditEntry, + Ok(pb::RecordKind::Note) => RecordKind::Note, + Ok(pb::RecordKind::CatalogItem) => RecordKind::CatalogItem, + Ok(pb::RecordKind::Unspecified) | Err(_) => { + return Err(Status::invalid_argument("record kind is required")); + } + }, + id: parse_uuid(&reference.record_id)?, + }) +} + +pub fn proto_record_ref(reference: RecordRef) -> pb::RecordRef { + pb::RecordRef { + kind: match reference.kind { + RecordKind::Organization => pb::RecordKind::Organization as i32, + RecordKind::Person => pb::RecordKind::Person as i32, + RecordKind::Relationship => pb::RecordKind::Relationship as i32, + RecordKind::Lead => pb::RecordKind::Lead as i32, + RecordKind::Opportunity => pb::RecordKind::Opportunity as i32, + RecordKind::Conversation => pb::RecordKind::Conversation as i32, + RecordKind::Activity => pb::RecordKind::Activity as i32, + RecordKind::Task => pb::RecordKind::Task as i32, + RecordKind::OfferQuote => pb::RecordKind::OfferQuote as i32, + RecordKind::OrderSubscription => pb::RecordKind::OrderSubscription as i32, + RecordKind::Document => pb::RecordKind::Document as i32, + RecordKind::Fact => pb::RecordKind::Fact as i32, + RecordKind::Intent => pb::RecordKind::Intent as i32, + RecordKind::WorkflowCase => pb::RecordKind::WorkflowCase as i32, + RecordKind::CommunicationEvent => pb::RecordKind::CommunicationEvent as i32, + RecordKind::PermissionGrant => pb::RecordKind::PermissionGrant as i32, + RecordKind::AuditEntry => pb::RecordKind::AuditEntry as i32, + RecordKind::Note => pb::RecordKind::Note as i32, + RecordKind::CatalogItem => pb::RecordKind::CatalogItem as i32, + }, + record_id: reference.id.to_string(), + } +} + +// --------------------------------------------------------------------------- +// Enum converters: kernel ← proto +// --------------------------------------------------------------------------- + +pub fn organization_lifecycle_from_proto(value: i32) -> OrganizationLifecycle { + match pb::OrganizationLifecycle::try_from(value).unwrap_or(pb::OrganizationLifecycle::Prospect) + { + pb::OrganizationLifecycle::Active => OrganizationLifecycle::Active, + pb::OrganizationLifecycle::Dormant => OrganizationLifecycle::Dormant, + pb::OrganizationLifecycle::Partner => OrganizationLifecycle::Partner, + pb::OrganizationLifecycle::Prospect | pb::OrganizationLifecycle::Unspecified => { + OrganizationLifecycle::Prospect + } + } +} + +pub fn opportunity_stage_from_proto(value: i32) -> OpportunityStage { + match pb::OpportunityStage::try_from(value).unwrap_or(pb::OpportunityStage::Qualifying) { + pb::OpportunityStage::Discovery => OpportunityStage::Discovery, + pb::OpportunityStage::Proposal => OpportunityStage::Proposal, + pb::OpportunityStage::Negotiation => OpportunityStage::Negotiation, + pb::OpportunityStage::ClosedWon => OpportunityStage::ClosedWon, + pb::OpportunityStage::ClosedLost => OpportunityStage::ClosedLost, + pb::OpportunityStage::Qualifying | pb::OpportunityStage::Unspecified => { + OpportunityStage::Qualifying + } + } +} + +pub fn activity_outcome_from_proto(value: i32) -> ActivityOutcome { + match pb::ActivityOutcome::try_from(value).unwrap_or(pb::ActivityOutcome::Completed) { + pb::ActivityOutcome::Waiting => ActivityOutcome::Waiting, + pb::ActivityOutcome::Blocked => ActivityOutcome::Blocked, + pb::ActivityOutcome::Completed | pb::ActivityOutcome::Unspecified => { + ActivityOutcome::Completed + } + } +} + +pub fn document_status_from_proto(value: i32) -> DocumentStatus { + match pb::DocumentStatus::try_from(value).unwrap_or(pb::DocumentStatus::Draft) { + pb::DocumentStatus::Verified => DocumentStatus::Verified, + pb::DocumentStatus::Archived => DocumentStatus::Archived, + pb::DocumentStatus::Draft | pb::DocumentStatus::Unspecified => DocumentStatus::Draft, + } +} + +pub fn communication_channel_from_proto(value: i32) -> CommunicationChannel { + match pb::CommunicationChannel::try_from(value).unwrap_or(pb::CommunicationChannel::Email) { + pb::CommunicationChannel::Phone => CommunicationChannel::Phone, + pb::CommunicationChannel::Meeting => CommunicationChannel::Meeting, + pb::CommunicationChannel::Chat => CommunicationChannel::Chat, + pb::CommunicationChannel::Sms => CommunicationChannel::Sms, + pb::CommunicationChannel::Email | pb::CommunicationChannel::Unspecified => { + CommunicationChannel::Email + } + } +} + +pub fn communication_direction_from_proto(value: i32) -> CommunicationDirection { + match pb::CommunicationDirection::try_from(value).unwrap_or(pb::CommunicationDirection::Inbound) + { + pb::CommunicationDirection::Outbound => CommunicationDirection::Outbound, + pb::CommunicationDirection::Internal => CommunicationDirection::Internal, + pb::CommunicationDirection::Inbound | pb::CommunicationDirection::Unspecified => { + CommunicationDirection::Inbound + } + } +} + +pub fn workflow_priority_from_proto(value: i32) -> WorkflowPriority { + match pb::WorkflowPriority::try_from(value).unwrap_or(pb::WorkflowPriority::Medium) { + pb::WorkflowPriority::Low => WorkflowPriority::Low, + pb::WorkflowPriority::High => WorkflowPriority::High, + pb::WorkflowPriority::Critical => WorkflowPriority::Critical, + pb::WorkflowPriority::Medium | pb::WorkflowPriority::Unspecified => { + WorkflowPriority::Medium + } + } +} + +pub fn workflow_state_from_proto(value: i32) -> WorkflowState { + match pb::WorkflowState::try_from(value).unwrap_or(pb::WorkflowState::Open) { + pb::WorkflowState::AwaitingApproval => WorkflowState::AwaitingApproval, + pb::WorkflowState::WaitingExternal => WorkflowState::WaitingExternal, + pb::WorkflowState::Blocked => WorkflowState::Blocked, + pb::WorkflowState::Done => WorkflowState::Done, + pb::WorkflowState::Open | pb::WorkflowState::Unspecified => WorkflowState::Open, + } +} + +pub fn relationship_type_from_proto(value: i32) -> application_kernel::RelationshipType { + match pb::RelationshipType::try_from(value).unwrap_or(pb::RelationshipType::Other) { + pb::RelationshipType::Employment => application_kernel::RelationshipType::Employment, + pb::RelationshipType::Champion => application_kernel::RelationshipType::Champion, + pb::RelationshipType::DecisionMaker => application_kernel::RelationshipType::DecisionMaker, + pb::RelationshipType::Partner => application_kernel::RelationshipType::Partner, + pb::RelationshipType::Competitor => application_kernel::RelationshipType::Competitor, + pb::RelationshipType::Other | pb::RelationshipType::Unspecified => { + application_kernel::RelationshipType::Other + } + } +} + +pub fn object_definition_kind_from_proto(value: i32) -> ObjectDefinitionKind { + match pb::ObjectDefinitionKind::try_from(value).unwrap_or(pb::ObjectDefinitionKind::Custom) { + pb::ObjectDefinitionKind::Standard => ObjectDefinitionKind::Standard, + pb::ObjectDefinitionKind::Custom | pb::ObjectDefinitionKind::Unspecified => { + ObjectDefinitionKind::Custom + } + } +} + +pub fn view_layout_from_proto(value: i32) -> ViewLayout { + match pb::ViewLayout::try_from(value).unwrap_or(pb::ViewLayout::Table) { + pb::ViewLayout::Kanban => ViewLayout::Kanban, + pb::ViewLayout::Calendar => ViewLayout::Calendar, + pb::ViewLayout::Table | pb::ViewLayout::Unspecified => ViewLayout::Table, + } +} + +pub fn field_type_from_proto(value: i32) -> FieldType { + match pb::FieldType::try_from(value).unwrap_or(pb::FieldType::Text) { + pb::FieldType::LongText => FieldType::LongText, + pb::FieldType::Number => FieldType::Number, + pb::FieldType::Currency => FieldType::Currency, + pb::FieldType::Boolean => FieldType::Boolean, + pb::FieldType::Date => FieldType::Date, + pb::FieldType::DateTime => FieldType::DateTime, + pb::FieldType::Email => FieldType::Email, + pb::FieldType::Phone => FieldType::Phone, + pb::FieldType::Url => FieldType::Url, + pb::FieldType::Select => FieldType::Select, + pb::FieldType::MultiSelect => FieldType::MultiSelect, + pb::FieldType::Relation => FieldType::Relation, + pb::FieldType::Text | pb::FieldType::Unspecified => FieldType::Text, + } +} + +pub fn relationship_cardinality_from_proto(value: i32) -> RelationshipCardinality { + match pb::RelationshipCardinality::try_from(value) + .unwrap_or(pb::RelationshipCardinality::OneToMany) + { + pb::RelationshipCardinality::OneToOne => RelationshipCardinality::OneToOne, + pb::RelationshipCardinality::ManyToMany => RelationshipCardinality::ManyToMany, + pb::RelationshipCardinality::OneToMany | pb::RelationshipCardinality::Unspecified => { + RelationshipCardinality::OneToMany + } + } +} + +// --------------------------------------------------------------------------- +// Proto constructors: proto ← kernel +// --------------------------------------------------------------------------- + +pub fn proto_organization(value: application_kernel::Organization) -> pb::Organization { + pb::Organization { + id: value.id.to_string(), + name: value.name, + external_key: value.external_key, + website: value.website, + industry: value.industry, + lifecycle: match value.lifecycle { + OrganizationLifecycle::Prospect => pb::OrganizationLifecycle::Prospect as i32, + OrganizationLifecycle::Active => pb::OrganizationLifecycle::Active as i32, + OrganizationLifecycle::Dormant => pb::OrganizationLifecycle::Dormant as i32, + OrganizationLifecycle::Partner => pb::OrganizationLifecycle::Partner as i32, + }, + owner_user_id: value.owner_user_id, + tags: value.tags, + created_at: proto_timestamp(value.created_at), + updated_at: proto_timestamp(value.updated_at), + } +} + +pub fn proto_person(value: application_kernel::Person) -> pb::Person { + pb::Person { + id: value.id.to_string(), + organization_id: value.organization_id.map(|id| id.to_string()), + full_name: value.full_name, + title: value.title, + email: value.email, + phone: value.phone, + linkedin_url: value.linkedin_url, + created_at: proto_timestamp(value.created_at), + updated_at: proto_timestamp(value.updated_at), + } +} + +pub fn proto_relationship(value: application_kernel::Relationship) -> pb::Relationship { + pb::Relationship { + id: value.id.to_string(), + from: Some(proto_record_ref(value.from)), + to: Some(proto_record_ref(value.to)), + relationship_type: match value.relationship_type { + application_kernel::RelationshipType::Employment => { + pb::RelationshipType::Employment as i32 + } + application_kernel::RelationshipType::Champion => pb::RelationshipType::Champion as i32, + application_kernel::RelationshipType::DecisionMaker => { + pb::RelationshipType::DecisionMaker as i32 + } + application_kernel::RelationshipType::Partner => pb::RelationshipType::Partner as i32, + application_kernel::RelationshipType::Competitor => { + pb::RelationshipType::Competitor as i32 + } + application_kernel::RelationshipType::Other => pb::RelationshipType::Other as i32, + }, + label: value.label, + created_at: proto_timestamp(value.created_at), + } +} + +pub fn proto_opportunity(value: application_kernel::Opportunity) -> pb::Opportunity { + pb::Opportunity { + id: value.id.to_string(), + organization_id: value.organization_id.to_string(), + primary_contact_id: value.primary_contact_id.map(|id| id.to_string()), + name: value.name, + stage: match value.stage { + OpportunityStage::Qualifying => pb::OpportunityStage::Qualifying as i32, + OpportunityStage::Discovery => pb::OpportunityStage::Discovery as i32, + OpportunityStage::Proposal => pb::OpportunityStage::Proposal as i32, + OpportunityStage::Negotiation => pb::OpportunityStage::Negotiation as i32, + OpportunityStage::ClosedWon => pb::OpportunityStage::ClosedWon as i32, + OpportunityStage::ClosedLost => pb::OpportunityStage::ClosedLost as i32, + }, + value: Some(pb::Money { + currency_code: value.value.currency_code, + amount_minor: value.value.amount_minor, + }), + confidence_bps: u32::from(value.confidence_bps), + next_step: value.next_step, + expected_close_at: value.expected_close_at.and_then(proto_timestamp), + created_at: proto_timestamp(value.created_at), + updated_at: proto_timestamp(value.updated_at), + } +} + +pub fn proto_activity(value: application_kernel::Activity) -> pb::Activity { + pb::Activity { + id: value.id.to_string(), + subject: value.subject, + details: value.details, + actor: Some(proto_actor(value.actor)), + related_to: value.related_to.into_iter().map(proto_record_ref).collect(), + outcome: match value.outcome { + ActivityOutcome::Completed => pb::ActivityOutcome::Completed as i32, + ActivityOutcome::Waiting => pb::ActivityOutcome::Waiting as i32, + ActivityOutcome::Blocked => pb::ActivityOutcome::Blocked as i32, + }, + occurred_at: proto_timestamp(value.occurred_at), + next_action_due_at: value.next_action_due_at.and_then(proto_timestamp), + } +} + +pub fn proto_note(value: application_kernel::Note) -> pb::Note { + pb::Note { + id: value.id.to_string(), + subject: value.subject, + body: value.body, + author: Some(proto_actor(value.author)), + related_to: value.related_to.into_iter().map(proto_record_ref).collect(), + promoted_to_fact: value.promoted_to_fact, + created_at: proto_timestamp(value.created_at), + } +} + +pub fn proto_document(value: application_kernel::Document) -> pb::Document { + pb::Document { + id: value.id.to_string(), + title: value.title, + media_type: value.media_type, + uri: value.uri, + status: match value.status { + DocumentStatus::Draft => pb::DocumentStatus::Draft as i32, + DocumentStatus::Verified => pb::DocumentStatus::Verified as i32, + DocumentStatus::Archived => pb::DocumentStatus::Archived as i32, + }, + uploaded_by: Some(proto_actor(value.uploaded_by)), + related_to: value.related_to.into_iter().map(proto_record_ref).collect(), + created_at: proto_timestamp(value.created_at), + } +} + +pub fn proto_communication_event( + value: application_kernel::CommunicationEvent, +) -> pb::CommunicationEvent { + pb::CommunicationEvent { + id: value.id.to_string(), + channel: match value.channel { + CommunicationChannel::Email => pb::CommunicationChannel::Email as i32, + CommunicationChannel::Phone => pb::CommunicationChannel::Phone as i32, + CommunicationChannel::Meeting => pb::CommunicationChannel::Meeting as i32, + CommunicationChannel::Chat => pb::CommunicationChannel::Chat as i32, + CommunicationChannel::Sms => pb::CommunicationChannel::Sms as i32, + }, + direction: match value.direction { + CommunicationDirection::Inbound => pb::CommunicationDirection::Inbound as i32, + CommunicationDirection::Outbound => pb::CommunicationDirection::Outbound as i32, + CommunicationDirection::Internal => pb::CommunicationDirection::Internal as i32, + }, + subject: value.subject, + summary: value.summary, + counterpart: value.counterpart, + actor: Some(proto_actor(value.actor)), + related_to: value.related_to.into_iter().map(proto_record_ref).collect(), + occurred_at: proto_timestamp(value.occurred_at), + } +} + +pub fn proto_workflow_case(value: application_kernel::WorkflowCase) -> pb::WorkflowCase { + pb::WorkflowCase { + id: value.id.to_string(), + title: value.title, + state: match value.state { + WorkflowState::Open => pb::WorkflowState::Open as i32, + WorkflowState::AwaitingApproval => pb::WorkflowState::AwaitingApproval as i32, + WorkflowState::WaitingExternal => pb::WorkflowState::WaitingExternal as i32, + WorkflowState::Blocked => pb::WorkflowState::Blocked as i32, + WorkflowState::Done => pb::WorkflowState::Done as i32, + }, + priority: match value.priority { + WorkflowPriority::Low => pb::WorkflowPriority::Low as i32, + WorkflowPriority::Medium => pb::WorkflowPriority::Medium as i32, + WorkflowPriority::High => pb::WorkflowPriority::High as i32, + WorkflowPriority::Critical => pb::WorkflowPriority::Critical as i32, + }, + owner_user_id: value.owner_user_id, + related_to: value.related_to.into_iter().map(proto_record_ref).collect(), + opened_at: proto_timestamp(value.opened_at), + updated_at: proto_timestamp(value.updated_at), + } +} + +pub fn proto_fact(value: application_kernel::Fact) -> pb::Fact { + pb::Fact { + id: value.id.to_string(), + statement: value.statement, + confidence_bps: u32::from(value.confidence_bps), + promoted_by: Some(proto_actor(value.promoted_by)), + source_note_id: value.source_note_id.map(|id| id.to_string()), + related_to: value.related_to.into_iter().map(proto_record_ref).collect(), + created_at: proto_timestamp(value.created_at), + } +} + +pub fn proto_timeline_entry(value: application_kernel::TimelineEntry) -> pb::TimelineEntry { + pb::TimelineEntry { + id: value.id.to_string(), + kind: match value.kind { + application_kernel::TimelineEntryKind::Activity => { + pb::TimelineEntryKind::Activity as i32 + } + application_kernel::TimelineEntryKind::Note => pb::TimelineEntryKind::Note as i32, + application_kernel::TimelineEntryKind::Document => { + pb::TimelineEntryKind::Document as i32 + } + application_kernel::TimelineEntryKind::Communication => { + pb::TimelineEntryKind::Communication as i32 + } + application_kernel::TimelineEntryKind::Fact => pb::TimelineEntryKind::Fact as i32, + application_kernel::TimelineEntryKind::Audit => pb::TimelineEntryKind::Audit as i32, + }, + anchor: value.anchor.map(proto_record_ref), + headline: value.headline, + body: value.body, + actor: Some(proto_actor(value.actor)), + occurred_at: proto_timestamp(value.occurred_at), + related_to: value.related_to.into_iter().map(proto_record_ref).collect(), + } +} + +pub fn proto_account_summary(value: application_kernel::AccountSummary) -> pb::AccountSummary { + pb::AccountSummary { + organization: Some(proto_organization(value.organization)), + contacts: value.contacts.into_iter().map(proto_person).collect(), + opportunities: value + .opportunities + .into_iter() + .map(proto_opportunity) + .collect(), + workflow_cases: value + .workflow_cases + .into_iter() + .map(proto_workflow_case) + .collect(), + facts: value.facts.into_iter().map(proto_fact).collect(), + documents: value.documents.into_iter().map(proto_document).collect(), + permissions: value + .permissions + .into_iter() + .map(proto_permission_grant) + .collect(), + recent_timeline: value + .recent_timeline + .into_iter() + .map(proto_timeline_entry) + .collect(), + } +} + +pub fn proto_permission_grant(value: application_kernel::PermissionGrant) -> pb::PermissionGrant { + pb::PermissionGrant { + id: value.id.to_string(), + subject: value.subject, + role: value.role, + scope: value.scope, + granted_by: Some(proto_actor(value.granted_by)), + created_at: proto_timestamp(value.created_at), + } +} + +pub fn field_definition_from_proto(field: pb::FieldDefinition) -> Result { + Ok(FieldDefinition { + id: parse_optional_uuid(Some(field.id))?.unwrap_or_else(Uuid::new_v4), + key: field.key, + label: field.label, + field_type: field_type_from_proto(field.field_type), + required: field.r#required, + options: field.options, + relation_object_key: field.relation_object_key, + active: field.active, + }) +} + +pub fn relationship_definition_from_proto( + definition: pb::RelationshipDefinition, +) -> Result { + Ok(RelationshipDefinition { + id: parse_optional_uuid(Some(definition.id))?.unwrap_or_else(Uuid::new_v4), + target_object_key: definition.target_object_key, + cardinality: relationship_cardinality_from_proto(definition.cardinality), + label: definition.label, + }) +} + +pub fn proto_field_definition(value: FieldDefinition) -> pb::FieldDefinition { + pb::FieldDefinition { + id: value.id.to_string(), + key: value.key, + label: value.label, + field_type: match value.field_type { + FieldType::Text => pb::FieldType::Text as i32, + FieldType::LongText => pb::FieldType::LongText as i32, + FieldType::Number => pb::FieldType::Number as i32, + FieldType::Currency => pb::FieldType::Currency as i32, + FieldType::Boolean => pb::FieldType::Boolean as i32, + FieldType::Date => pb::FieldType::Date as i32, + FieldType::DateTime => pb::FieldType::DateTime as i32, + FieldType::Email => pb::FieldType::Email as i32, + FieldType::Phone => pb::FieldType::Phone as i32, + FieldType::Url => pb::FieldType::Url as i32, + FieldType::Select => pb::FieldType::Select as i32, + FieldType::MultiSelect => pb::FieldType::MultiSelect as i32, + FieldType::Relation => pb::FieldType::Relation as i32, + }, + r#required: value.required, + options: value.options, + relation_object_key: value.relation_object_key, + active: value.active, + } +} + +pub fn proto_relationship_definition(value: RelationshipDefinition) -> pb::RelationshipDefinition { + pb::RelationshipDefinition { + id: value.id.to_string(), + target_object_key: value.target_object_key, + cardinality: match value.cardinality { + RelationshipCardinality::OneToOne => pb::RelationshipCardinality::OneToOne as i32, + RelationshipCardinality::OneToMany => pb::RelationshipCardinality::OneToMany as i32, + RelationshipCardinality::ManyToMany => pb::RelationshipCardinality::ManyToMany as i32, + }, + label: value.label, + } +} + +pub fn proto_object_definition( + value: application_kernel::ObjectDefinition, +) -> pb::ObjectDefinition { + pb::ObjectDefinition { + id: value.id.to_string(), + key: value.key, + display_name: value.display_name, + kind: match value.kind { + ObjectDefinitionKind::Standard => pb::ObjectDefinitionKind::Standard as i32, + ObjectDefinitionKind::Custom => pb::ObjectDefinitionKind::Custom as i32, + }, + fields: value + .fields + .into_iter() + .map(proto_field_definition) + .collect(), + relationships: value + .relationships + .into_iter() + .map(proto_relationship_definition) + .collect(), + active: value.active, + created_at: proto_timestamp(value.created_at), + updated_at: proto_timestamp(value.updated_at), + } +} + +pub fn proto_view_definition(value: application_kernel::ViewDefinition) -> pb::ViewDefinition { + pb::ViewDefinition { + id: value.id.to_string(), + object_key: value.object_key, + name: value.name, + layout: match value.layout { + ViewLayout::Table => pb::ViewLayout::Table as i32, + ViewLayout::Kanban => pb::ViewLayout::Kanban as i32, + ViewLayout::Calendar => pb::ViewLayout::Calendar as i32, + }, + filter_expression: value.filter_expression, + sort_expression: value.sort_expression, + visible_fields: value.visible_fields, + group_by: value.group_by, + favorite: value.favorite, + owner_user_id: value.owner_user_id, + created_at: proto_timestamp(value.created_at), + updated_at: proto_timestamp(value.updated_at), + } +} diff --git a/showcase/crm-helm/src/truths.rs b/showcase/crm-helm/src/truths.rs new file mode 100644 index 0000000..8b3ab89 --- /dev/null +++ b/showcase/crm-helm/src/truths.rs @@ -0,0 +1 @@ +//! CRM Truths module — filled in Phase 6 from helm-truth-execution truth resolution and promotion diff --git a/showcase/crm-helm/src/truths/evaluate_acquisition_target.rs b/showcase/crm-helm/src/truths/evaluate_acquisition_target.rs new file mode 100644 index 0000000..59b3dd1 --- /dev/null +++ b/showcase/crm-helm/src/truths/evaluate_acquisition_target.rs @@ -0,0 +1,476 @@ +//! NOT YET WIRED — see Phase 8 or Phase 6b for integration. (Originally lived in helms/crates/application-server/src/truth_runtime/.) +use std::collections::HashMap; +use std::sync::Arc; + +use application_kernel::{Actor as CrmActor, FactRecord}; +use application_storage::{KernelStore, StoreWriteResult}; +use converge_kernel::{ContextState as Context, ConvergeResult, Engine}; +use converge_pack::{AgentEffect, Context as ContextView, ContextKey, Suggestor}; +use converge_provider::{BoxFuture, ChatRequest, ChatResponse, DynChatBackend, LlmError}; +use organism_pack::{ + BreadthResearchSuggestor, ContradictionFinderSuggestor, DdError, DdSearch, + DepthResearchSuggestor, FactExtractorSuggestor, GapDetectorSuggestor, HuddleSeedSuggestor, + IntentPacket, Plan, PlanStep, ReasoningSystem, SearchHit, SharedBudget, SynthesisSuggestor, +}; +use tonic::Status; +use truth_catalog::{ + EvaluateAcquisitionTargetEvaluator, + admission::{admit_truth_intent, default_helms_capabilities, select_formation_for_intent}, + converge_binding_for_truth, +}; + +use super::{ + TruthExecutionArtifacts, TruthProjection, + common::{optional_input, required_input}, + domain_event_kind_name, +}; + +const DD_PACK_ID: &str = "prio-dd-pack"; +const TRUST_PACK_ID: &str = "trust"; + +// ── Input ─────────────────────────────────────────────────────────── + +#[derive(Debug, Clone)] +pub struct EvaluateAcquisitionTargetInput { + pub target_company: String, + pub focus_areas: Option, + pub max_searches: Option, + pub max_llm_calls: Option, +} + +impl EvaluateAcquisitionTargetInput { + pub fn from_map(inputs: &HashMap) -> Result { + Ok(Self { + target_company: required_input(inputs, "target_company")?.to_string(), + focus_areas: optional_input(inputs, "focus_areas"), + max_searches: optional_input(inputs, "max_searches").and_then(|s| s.parse().ok()), + max_llm_calls: optional_input(inputs, "max_llm_calls").and_then(|s| s.parse().ok()), + }) + } +} + +// ── Executor ──────────────────────────────────────────────────────── + +pub(super) async fn execute( + store: &S, + runtime_stores: &application_storage::AppRuntimeStores, + inputs: EvaluateAcquisitionTargetInput, + actor: CrmActor, + persist_projection: bool, +) -> Result { + let binding = converge_binding_for_truth("evaluate-acquisition-target") + .ok_or_else(|| Status::not_found("truth not found: evaluate-acquisition-target"))?; + + let company = inputs.target_company.clone(); + let focus_areas = inputs.focus_areas.clone(); + let max_searches = inputs.max_searches.unwrap_or(10); + let max_llm_calls = inputs.max_llm_calls.unwrap_or(8); + + let budget = Arc::new( + SharedBudget::new() + .with_limit("searches", max_searches) + .with_limit("llm", max_llm_calls), + ); + + // Build search and LLM backends. + // Option A (dev/demo): stub backends that prove the wiring works. + // Option B (production): MCP directory → real provider backends. + // Option C (governed): Converge capability axioms enforce credentials + budget. + let search: Arc = Arc::new(StubDdSearch); + let llm: Arc = Arc::new(StubChatBackend); + + // Build the initial plans (Organism huddle seed pattern) + let intent = build_dd_intent(&company, focus_areas.as_deref()); + let plans = build_dd_plans(&company, focus_areas.as_deref()); + let huddle_seed = HuddleSeedSuggestor::from_plans(intent, plans); + + let mut engine = Engine::new(); + + // Register organism DD suggestors — the convergent research loop + engine.register_suggestor_in_pack(DD_PACK_ID, huddle_seed); + engine.register_suggestor_in_pack( + DD_PACK_ID, + BreadthResearchSuggestor::new(&company, budget.clone(), search.clone()), + ); + engine.register_suggestor_in_pack( + DD_PACK_ID, + DepthResearchSuggestor::new(&company, budget.clone(), search), + ); + engine.register_suggestor_in_pack( + DD_PACK_ID, + FactExtractorSuggestor::new(&company, budget.clone(), llm.clone()), + ); + engine.register_suggestor_in_pack( + DD_PACK_ID, + GapDetectorSuggestor::new(&company, budget.clone(), llm.clone()) + .with_max_generations(3) + .with_min_hypotheses(5), + ); + engine.register_suggestor_in_pack(DD_PACK_ID, ContradictionFinderSuggestor::new()); + engine.register_suggestor_in_pack( + DD_PACK_ID, + SynthesisSuggestor::new(&company, budget, llm).with_required_stable_cycles(1), + ); + + // Governance gate — blocks recommendation when contradictions need human review + engine.register_suggestor_in_pack(TRUST_PACK_ID, ContradictionGateAgent); + + let mut seed_ctx = seed_context(&company)?; + let intent = admit_truth_intent( + "evaluate-acquisition-target", + &actor.actor_id, + "truth:evaluate-acquisition-target", + &mut seed_ctx, + ) + .map_err(|e| Status::internal(format!("admit intent failed: {e}")))?; + let selection = select_formation_for_intent(&intent, &default_helms_capabilities()) + .map_err(|e| Status::internal(format!("formation selection failed: {e}")))?; + tracing::info!( + truth = "evaluate-acquisition-target", + primary = %selection.primary_template_id, + alternates = ?selection.alternate_template_ids, + "formation selected" + ); + + let runtime_ctx = super::RuntimeContext { + scope_id: format!("dd:{}", company.to_lowercase().replace(' ', "-")), + }; + let (result, experience_events) = super::run_engine_with_runtime( + runtime_stores, + &mut engine, + &runtime_ctx, + seed_ctx, + &binding.intent, + std::sync::Arc::new(EvaluateAcquisitionTargetEvaluator), + ) + .await?; + + let projection = if persist_projection { + Some(project(store, &inputs, &result, actor)?) + } else { + None + }; + + Ok(TruthExecutionArtifacts { + result, + experience_events, + projection, + runtime_scope_id: runtime_ctx.scope_id, + }) +} + +// ── Governance Gate Suggestor ──────────────────────────────────────── + +struct ContradictionGateAgent; + +#[async_trait::async_trait] +impl Suggestor for ContradictionGateAgent { + fn name(&self) -> &str { + "contradiction-gate" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Evaluations] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + ctx.get(ContextKey::Evaluations) + .iter() + .any(|f| f.id().starts_with("contradiction-")) + && !ctx + .get(ContextKey::Evaluations) + .iter() + .any(|f| f.id() == "dd:human-review-required") + } + + async fn execute(&self, ctx: &dyn ContextView) -> AgentEffect { + let contradictions: Vec<_> = ctx + .get(ContextKey::Evaluations) + .iter() + .filter(|f| f.id().starts_with("contradiction-")) + .map(|f| f.id().clone()) + .collect(); + + if contradictions.is_empty() { + return AgentEffect::empty(); + } + + let content = serde_json::json!({ + "type": "governance-gate", + "reason": "material contradictions detected in DD research", + "contradiction_count": contradictions.len(), + "contradiction_ids": contradictions, + "required_action": "investment committee must review contradictions before recommendation", + }) + .to_string(); + + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Evaluations, + "dd:human-review-required", + content, + "contradiction-gate", + ) + .with_confidence(1.0), + ) + } +} + +// ── Intent & Plans ────────────────────────────────────────────────── + +fn build_dd_intent(company: &str, focus_areas: Option<&str>) -> IntentPacket { + let expires = chrono::Utc::now() + chrono::Duration::hours(1); + let mut intent = IntentPacket::new( + format!("Build a due diligence brief for {company}"), + expires, + ) + .with_context(serde_json::json!({ + "company": company, + "goal": "research breadth, depth, and investment risks", + "focus_areas": focus_areas, + })); + intent.authority = vec!["research".to_string()]; + intent +} + +fn build_dd_plans(company: &str, focus_areas: Option<&str>) -> Vec { + let expires = chrono::Utc::now() + chrono::Duration::hours(1); + let intent = IntentPacket::new(format!("DD for {company}"), expires); + let focus_suffix = focus_areas + .map(|focus| format!(" with focus on {focus}")) + .unwrap_or_default(); + + let mut breadth1 = Plan::new( + &intent, + "Search wide for product, customer, and market context", + ); + breadth1.contributor = ReasoningSystem::DomainModel; + breadth1.steps = vec![PlanStep { + action: format!( + "[breadth] {company} products customers market position overview{focus_suffix}" + ), + expected_effect: "discover product scope, customer segments, and market presence".into(), + }]; + + let mut breadth2 = Plan::new( + &intent, + "Search wide for competitors, growth, and positioning", + ); + breadth2.contributor = ReasoningSystem::CausalAnalysis; + breadth2.steps = vec![PlanStep { + action: format!( + "[breadth] {company} competitors trends growth tech stack USP{focus_suffix}" + ), + expected_effect: "map competitive landscape and growth trajectory".into(), + }]; + + let mut depth1 = Plan::new( + &intent, + "Search deep for architecture and integration evidence", + ); + depth1.contributor = ReasoningSystem::ConstraintSolver; + depth1.steps = vec![PlanStep { + action: format!( + "[depth] {company} technology architecture platform integrations API{focus_suffix}" + ), + expected_effect: "understand technical moat and integration surface".into(), + }]; + + let mut depth2 = Plan::new(&intent, "Search deep for financial and ownership evidence"); + depth2.contributor = ReasoningSystem::CostEstimation; + depth2.steps = vec![PlanStep { + action: format!( + "[depth] {company} revenue ARR funding investors ownership financials{focus_suffix}" + ), + expected_effect: "find financial metrics and ownership structure".into(), + }]; + + vec![breadth1, breadth2, depth1, depth2] +} + +fn seed_context(company: &str) -> Result { + let mut ctx = Context::new(); + ctx.add_input( + ContextKey::Seeds, + "dd:target", + serde_json::json!({ "company": company }).to_string(), + ) + .map_err(|error| Status::failed_precondition(error.to_string()))?; + Ok(ctx) +} + +// ── Projection ────────────────────────────────────────────────────── + +fn project( + store: &S, + inputs: &EvaluateAcquisitionTargetInput, + result: &ConvergeResult, + actor: CrmActor, +) -> Result { + let synthesis_content = result + .context + .get(ContextKey::Proposals) + .iter() + .find(|f| f.id().starts_with("synthesis-")) + .map(|f| f.text().unwrap_or_default().to_string()); + + let hypothesis_count = result.context.get(ContextKey::Hypotheses).len(); + let contradiction_count = result + .context + .get(ContextKey::Evaluations) + .iter() + .filter(|f| f.id().starts_with("contradiction-")) + .count(); + + let StoreWriteResult { + value: facts, + events, + } = store + .write_with_events(|kernel| { + let dd_fact = kernel.record_fact( + FactRecord { + statement: format!( + "Due diligence for {} completed: {} hypotheses, {} contradictions{}", + inputs.target_company, + hypothesis_count, + contradiction_count, + synthesis_content + .as_ref() + .map(|_| ", synthesis produced") + .unwrap_or(", no synthesis (budget or convergence limit)") + ), + confidence_bps: if synthesis_content.is_some() { + 8_000 + } else { + 5_000 + }, + related_to: vec![], + source_note_id: None, + }, + actor.clone(), + )?; + Ok(vec![dd_fact]) + }) + .map_err(super::status_from_storage)?; + + Ok(TruthProjection { + organization: None, + person: None, + opportunity: None, + subscription: None, + entitlements: vec![], + ledger_entries: vec![], + documents: vec![], + workflow_cases: vec![], + facts, + domain_event_kinds: events.iter().map(domain_event_kind_name).collect(), + }) +} + +// ── Stub Backends (Option A: static/dev) ──────────────────────────── +// +// These stubs prove the Truth→Formation→Convergence→Governance chain +// compiles and runs without external API keys. Replace with real +// backends for production (see kb/Architecture/Capability Binding.md). + +struct StubDdSearch; + +#[async_trait::async_trait] +impl DdSearch for StubDdSearch { + async fn search(&self, query: &str) -> Result, DdError> { + Ok(vec![SearchHit { + title: format!("Stub result for: {query}"), + url: format!("https://stub.example.com/{}", query.replace(' ', "-")), + content: format!( + "This is a stub search result for the query '{query}'. \ + In production, this would be a real Brave or Tavily search result." + ), + provider: "stub".into(), + }]) + } +} + +struct StubChatBackend; + +impl DynChatBackend for StubChatBackend { + fn chat(&self, _req: ChatRequest) -> BoxFuture<'_, Result> { + Box::pin(async { + let content = serde_json::json!({ + "facts": [ + { + "claim": "Stub company operates in the B2B SaaS market", + "category": "market", + "source_indices": [0], + "confidence": 0.8 + }, + { + "claim": "Stub company has approximately 200 employees", + "category": "team", + "source_indices": [0], + "confidence": 0.7 + }, + { + "claim": "Stub company uses a cloud-native architecture", + "category": "technology", + "source_indices": [0], + "confidence": 0.75 + }, + { + "claim": "Stub company competes with established players in the space", + "category": "competition", + "source_indices": [0], + "confidence": 0.7 + }, + { + "claim": "Stub company serves enterprise customers across Europe", + "category": "customers", + "source_indices": [0], + "confidence": 0.8 + } + ] + }) + .to_string(); + Ok(ChatResponse { + content, + tool_calls: Vec::new(), + usage: None, + model: None, + finish_reason: None, + metadata: std::collections::HashMap::new(), + }) + }) + } +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use super::*; + + #[test] + fn input_parsing_requires_target_company() { + let inputs = HashMap::new(); + assert!(EvaluateAcquisitionTargetInput::from_map(&inputs).is_err()); + } + + #[test] + fn input_parsing_accepts_minimal_inputs() { + let mut inputs = HashMap::new(); + inputs.insert("target_company".to_string(), "Acme Corp".to_string()); + let parsed = EvaluateAcquisitionTargetInput::from_map(&inputs).unwrap(); + assert_eq!(parsed.target_company, "Acme Corp"); + assert!(parsed.focus_areas.is_none()); + assert!(parsed.max_searches.is_none()); + } + + #[test] + fn dd_plans_cover_four_research_vectors() { + let plans = build_dd_plans("TestCo", None); + assert_eq!(plans.len(), 4); + assert_eq!(plans[0].contributor, ReasoningSystem::DomainModel); + assert_eq!(plans[1].contributor, ReasoningSystem::CausalAnalysis); + assert_eq!(plans[2].contributor, ReasoningSystem::ConstraintSolver); + assert_eq!(plans[3].contributor, ReasoningSystem::CostEstimation); + } +} diff --git a/showcase/crm-helm/src/truths/match_renewal_context.rs b/showcase/crm-helm/src/truths/match_renewal_context.rs new file mode 100644 index 0000000..da6183e --- /dev/null +++ b/showcase/crm-helm/src/truths/match_renewal_context.rs @@ -0,0 +1,798 @@ +//! NOT YET WIRED — see Phase 8 or Phase 6b for integration. (Originally lived in helms/crates/application-server/src/truth_runtime/.) +use std::collections::HashMap; +use std::sync::Arc; + +use application_kernel::{ + AccountSummary, Actor as CrmActor, DocumentAttach, DocumentStatus, FactRecord, RecordKind, + RecordRef, WorkflowCaseAdvance, WorkflowCaseCreate, WorkflowPriority, WorkflowState, +}; +use application_storage::{KernelStore, StorageError, StoreWriteResult}; +use converge_kernel::{ContextState as Context, ConvergeResult, Engine}; +use converge_knowledge::{KnowledgeBase, KnowledgeEntry, SearchOptions}; +use converge_pack::{AgentEffect, Context as ContextView, ContextKey, Suggestor}; +use serde::{Deserialize, Serialize}; +use tonic::Status; +use truth_catalog::{ + MatchRenewalContextEvaluator, + admission::{admit_truth_intent, default_helms_capabilities, select_formation_for_intent}, + converge_binding_for_truth, +}; +use uuid::Uuid; + +use super::{ + TruthExecutionArtifacts, TruthProjection, + common::{block_on_async, has_fact_id, optional_uuid, payload_from_result, required_uuid}, + domain_event_kind_name, status_from_storage, +}; + +const RELATIONSHIP_PACK_ID: &str = "prio-relationship-pack"; +const COMMERCIAL_PACK_ID: &str = "prio-commercial-pack"; +const WORK_PACK_ID: &str = "prio-work-pack"; +const KNOWLEDGE_PACK_ID: &str = "knowledge"; +const CONTEXT_INDEX_FACT_ID: &str = "renewal:context-indexed"; +const RENEWAL_BRIEF_FACT_ID: &str = "renewal:brief"; +const RENEWAL_TERMS_FACT_ID: &str = "renewal:terms"; +const KNOWLEDGE_PROVENANCE: &str = "prio.match-renewal-context.knowledge"; +const BRIEF_PROVENANCE: &str = "prio.match-renewal-context.brief"; +const TERMS_PROVENANCE: &str = "prio.match-renewal-context.terms"; + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct RenewalIndexPayload { + organization_id: Uuid, + entry_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct RenewalSignalPayload { + signal_id: String, + query: String, + title: String, + summary: String, + source: Option, + similarity_bps: u16, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct RenewalBriefPayload { + summary: String, + strengths: Vec, + risks: Vec, + opportunities: Vec, + talking_points: Vec, + confidence_bps: u16, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct RenewalTermsPayload { + recommendation: String, + rationale: String, + approval_required: bool, + confidence_bps: u16, +} + +#[derive(Clone)] +struct ContextGathererAgent { + store: S, + organization_id: Uuid, + knowledge_base: Arc, +} + +#[derive(Clone)] +struct RenewalSignalAgent { + knowledge_base: Arc, +} + +struct NegotiationBriefAgent; + +struct RenewalTermsAgent; + +#[derive(Debug, Clone)] +pub struct MatchRenewalContextInput { + pub organization_id: Uuid, + pub opportunity_id: Option, +} + +impl MatchRenewalContextInput { + pub fn from_map(inputs: &HashMap) -> Result { + Ok(Self { + organization_id: required_uuid(inputs, "organization_id")?, + opportunity_id: optional_uuid(inputs, "opportunity_id")?, + }) + } +} + +pub(super) async fn execute( + store: &S, + runtime_stores: &application_storage::AppRuntimeStores, + inputs: MatchRenewalContextInput, + actor: CrmActor, + persist_projection: bool, +) -> Result { + let binding = converge_binding_for_truth("match-renewal-context") + .ok_or_else(|| Status::not_found("truth not found: match-renewal-context"))?; + + let organization_id = inputs.organization_id; + let scratch = tempfile::tempdir().map_err(|error| { + Status::internal(format!("failed to create renewal scratch dir: {error}")) + })?; + let kb_path = scratch.path().join("renewal.kb"); + let kb_path_owned = kb_path.clone(); + let knowledge_base = Arc::new( + block_on_async(async move { KnowledgeBase::open(kb_path_owned).await }) + .map_err(|error| Status::internal(format!("failed to open knowledge base: {error}")))?, + ); + + let mut engine = Engine::new(); + engine.register_suggestor_in_pack( + RELATIONSHIP_PACK_ID, + ContextGathererAgent { + store: store.clone(), + organization_id, + knowledge_base: knowledge_base.clone(), + }, + ); + engine.register_suggestor_in_pack( + KNOWLEDGE_PACK_ID, + RenewalSignalAgent { + knowledge_base: knowledge_base.clone(), + }, + ); + engine.register_suggestor_in_pack(WORK_PACK_ID, NegotiationBriefAgent); + engine.register_suggestor_in_pack(COMMERCIAL_PACK_ID, RenewalTermsAgent); + + let mut seed_ctx = seed_context(organization_id)?; + let intent = admit_truth_intent( + "match-renewal-context", + &actor.actor_id, + "truth:match-renewal-context", + &mut seed_ctx, + ) + .map_err(|e| Status::internal(format!("admit intent failed: {e}")))?; + let selection = select_formation_for_intent(&intent, &default_helms_capabilities()) + .map_err(|e| Status::internal(format!("formation selection failed: {e}")))?; + tracing::info!( + truth = "match-renewal-context", + primary = %selection.primary_template_id, + alternates = ?selection.alternate_template_ids, + "formation selected" + ); + + let runtime_ctx = super::RuntimeContext { + scope_id: inputs.organization_id.to_string(), + }; + let (result, experience_events) = super::run_engine_with_runtime( + runtime_stores, + &mut engine, + &runtime_ctx, + seed_ctx, + &binding.intent, + std::sync::Arc::new(MatchRenewalContextEvaluator), + ) + .await?; + + let projection = if persist_projection { + Some(project(store, &inputs, &result, actor)?) + } else { + None + }; + + Ok(TruthExecutionArtifacts { + result, + experience_events, + projection, + runtime_scope_id: runtime_ctx.scope_id, + }) +} + +#[async_trait::async_trait] +impl Suggestor for ContextGathererAgent { + fn name(&self) -> &str { + "ContextGathererAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Seeds] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + ctx.has(ContextKey::Seeds) && !has_fact_id(ctx, ContextKey::Signals, CONTEXT_INDEX_FACT_ID) + } + + async fn execute(&self, _ctx: &dyn ContextView) -> AgentEffect { + let summary = match account_summary_from_store(&self.store, self.organization_id) { + Ok(summary) => summary, + Err(error) => { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "renewal:context:error", + error, + KNOWLEDGE_PROVENANCE, + ) + .with_confidence(1.0), + ); + } + }; + let entries = knowledge_entries_from_summary(&summary); + let entry_count = entries.len(); + let knowledge_base = self.knowledge_base.clone(); + let entries_for_ingest = entries.clone(); + if let Err(error) = + block_on_async(async move { knowledge_base.add_entries(entries_for_ingest).await }) + { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "renewal:context:error", + error.to_string(), + KNOWLEDGE_PROVENANCE, + ) + .with_confidence(1.0), + ); + } + + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Signals, + CONTEXT_INDEX_FACT_ID, + serde_json::to_string(&RenewalIndexPayload { + organization_id: self.organization_id, + entry_count, + }) + .unwrap_or_default(), + KNOWLEDGE_PROVENANCE, + ) + .with_confidence(1.0), + ) + } +} + +#[async_trait::async_trait] +impl Suggestor for RenewalSignalAgent { + fn name(&self) -> &str { + "RenewalSignalAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Signals] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + has_fact_id(ctx, ContextKey::Signals, CONTEXT_INDEX_FACT_ID) + && !ctx + .get(ContextKey::Signals) + .iter() + .any(|fact| fact.id().starts_with("renewal:signal:")) + } + + async fn execute(&self, _ctx: &dyn ContextView) -> AgentEffect { + let queries = [ + ( + "competitive-evaluation", + "competitive evaluation competitor dissatisfaction", + ), + ("support-incident", "incident outage severity escalation"), + ( + "expansion-interest", + "expansion growth upgrade usage increase", + ), + ( + "stakeholder-change", + "stakeholder change executive sponsor champion", + ), + ]; + + let mut builder = AgentEffect::builder(); + for (signal_id, query) in queries { + let options = SearchOptions::new(1) + .with_min_similarity(0.0) + .with_diversity(0.2) + .hybrid(0.35); + let knowledge_base = self.knowledge_base.clone(); + let query = query.to_string(); + let query_for_search = query.clone(); + let results = match block_on_async(async move { + knowledge_base.search(&query_for_search, options).await + }) { + Ok(results) => results, + Err(error) => { + builder.push( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + format!("renewal:signal:error:{signal_id}"), + error.to_string(), + KNOWLEDGE_PROVENANCE, + ) + .with_confidence(1.0), + ); + continue; + } + }; + let Some(result) = results.first() else { + continue; + }; + let payload = RenewalSignalPayload { + signal_id: signal_id.to_string(), + query: query.to_string(), + title: result.entry.title.clone(), + summary: summarize(&result.entry.content, 180), + source: result.entry.source.clone(), + similarity_bps: (result.similarity.clamp(0.0, 1.0) * 10_000.0).round() as u16, + }; + builder.push( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Signals, + format!("renewal:signal:{signal_id}"), + serde_json::to_string(&payload).unwrap_or_default(), + KNOWLEDGE_PROVENANCE, + ) + .with_confidence(result.similarity as f64), + ); + } + builder.build() + } +} + +#[async_trait::async_trait] +impl Suggestor for NegotiationBriefAgent { + fn name(&self) -> &str { + "NegotiationBriefAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Signals] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + ctx.get(ContextKey::Signals) + .iter() + .any(|fact| fact.id().starts_with("renewal:signal:")) + && !has_fact_id(ctx, ContextKey::Strategies, RENEWAL_BRIEF_FACT_ID) + } + + async fn execute(&self, ctx: &dyn ContextView) -> AgentEffect { + let signals = renewal_signals_from_context(ctx); + if signals.is_empty() { + return AgentEffect::empty(); + } + + let risks = signals + .iter() + .filter(|signal| { + signal.signal_id.contains("competitive") + || signal.signal_id.contains("incident") + || signal.signal_id.contains("stakeholder") + }) + .map(|signal| format!("{}: {}", signal.signal_id, signal.summary)) + .collect::>(); + let opportunities = signals + .iter() + .filter(|signal| signal.signal_id.contains("expansion")) + .map(|signal| signal.summary.clone()) + .collect::>(); + let strengths = if risks.is_empty() { + vec!["Account context is stable with no major negative retrieval signals.".to_string()] + } else { + vec![ + "Retrieved account context is dense enough to support a guided renewal discussion." + .to_string(), + ] + }; + let talking_points = signals + .iter() + .map(|signal| { + format!( + "Discuss {} with evidence from {}", + signal.signal_id, signal.title + ) + }) + .collect::>(); + let summary = format!( + "Renewal brief built from {} retrieved signals.", + signals.len() + ); + let confidence_bps = (6_500 + (signals.len().min(4) as u16 * 700)).min(9_200); + let payload = RenewalBriefPayload { + summary, + strengths, + risks, + opportunities, + talking_points, + confidence_bps, + }; + + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Strategies, + RENEWAL_BRIEF_FACT_ID, + serde_json::to_string(&payload).unwrap_or_default(), + BRIEF_PROVENANCE, + ) + .with_confidence(f64::from(confidence_bps) / 10_000.0), + ) + } +} + +#[async_trait::async_trait] +impl Suggestor for RenewalTermsAgent { + fn name(&self) -> &str { + "RenewalTermsAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Strategies] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + has_fact_id(ctx, ContextKey::Strategies, RENEWAL_BRIEF_FACT_ID) + && !has_fact_id(ctx, ContextKey::Strategies, RENEWAL_TERMS_FACT_ID) + } + + async fn execute(&self, ctx: &dyn ContextView) -> AgentEffect { + let Some(brief_fact) = ctx + .get(ContextKey::Strategies) + .iter() + .find(|fact| fact.id() == RENEWAL_BRIEF_FACT_ID) + else { + return AgentEffect::empty(); + }; + let brief = match serde_json::from_str::( + brief_fact.text().unwrap_or_default(), + ) { + Ok(brief) => brief, + Err(error) => { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "renewal:terms:error", + error.to_string(), + TERMS_PROVENANCE, + ) + .with_confidence(1.0), + ); + } + }; + + let approval_required = brief + .risks + .iter() + .any(|risk| risk.contains("competitive") || risk.contains("incident")); + let recommendation = if !brief.opportunities.is_empty() && !approval_required { + "propose expansion-oriented renewal".to_string() + } else if approval_required { + "prepare standard renewal with explicit mitigation and approval".to_string() + } else { + "propose standard renewal".to_string() + }; + let payload = RenewalTermsPayload { + rationale: brief.summary.clone(), + approval_required, + confidence_bps: if approval_required { 6_200 } else { 8_100 }, + recommendation, + }; + + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Strategies, + RENEWAL_TERMS_FACT_ID, + serde_json::to_string(&payload).unwrap_or_default(), + TERMS_PROVENANCE, + ) + .with_confidence(f64::from(payload.confidence_bps) / 10_000.0), + ) + } +} + +fn project( + store: &S, + inputs: &MatchRenewalContextInput, + result: &ConvergeResult, + actor: CrmActor, +) -> Result { + let organization_id = inputs.organization_id; + let opportunity_id = inputs.opportunity_id; + let _brief = renewal_brief_from_result(result)?; + let terms = renewal_terms_from_result(result)?; + let signals = renewal_signals_from_result(result)?; + let organization = account_summary_from_store(store, organization_id) + .map_err(Status::failed_precondition)? + .organization; + + let mut related_to = vec![RecordRef { + kind: RecordKind::Organization, + id: organization.id, + }]; + if let Some(opportunity_id) = opportunity_id { + related_to.push(RecordRef { + kind: RecordKind::Opportunity, + id: opportunity_id, + }); + } + + let StoreWriteResult { value, events } = store + .write_with_events(|kernel| { + let document = kernel.attach_document( + DocumentAttach { + title: format!("Renewal brief: {}", organization.name), + media_type: "text/markdown".to_string(), + uri: format!( + "converge://truths/match-renewal-context/{}/brief.md", + organization.id + ), + status: DocumentStatus::Draft, + related_to: related_to.clone(), + }, + actor.clone(), + )?; + + let mut projected_facts = Vec::new(); + for signal in &signals { + projected_facts.push(kernel.record_fact( + FactRecord { + statement: format!( + "Renewal signal {} from {}: {}", + signal.signal_id, signal.title, signal.summary + ), + confidence_bps: signal.similarity_bps, + related_to: related_to.clone(), + source_note_id: None, + }, + actor.clone(), + )?); + } + projected_facts.push(kernel.record_fact( + FactRecord { + statement: format!( + "Renewal terms recommendation: {} ({})", + terms.recommendation, terms.rationale + ), + confidence_bps: terms.confidence_bps, + related_to: related_to.clone(), + source_note_id: None, + }, + actor.clone(), + )?); + + let workflow_case = if terms.approval_required { + let case = kernel.create_workflow_case( + WorkflowCaseCreate { + title: format!("Renewal approval: {}", organization.name), + priority: WorkflowPriority::High, + owner_user_id: None, + related_to: related_to.clone(), + }, + actor.clone(), + )?; + Some(kernel.advance_workflow_case( + WorkflowCaseAdvance { + workflow_case_id: case.id, + state: WorkflowState::AwaitingApproval, + }, + actor, + )?) + } else { + None + }; + + Ok((document, projected_facts, workflow_case)) + }) + .map_err(status_from_storage)?; + + let (document, facts, workflow_case) = value; + Ok(TruthProjection { + organization: Some(organization), + person: None, + opportunity: None, + subscription: None, + entitlements: Vec::new(), + ledger_entries: Vec::new(), + documents: vec![document], + workflow_cases: workflow_case.into_iter().collect(), + facts, + domain_event_kinds: events.iter().map(domain_event_kind_name).collect(), + }) +} + +fn seed_context(organization_id: Uuid) -> Result { + let mut context = Context::new(); + context + .add_input( + ContextKey::Seeds, + "match-renewal-context:seed", + organization_id.to_string(), + ) + .map_err(|error| Status::failed_precondition(error.to_string()))?; + Ok(context) +} + +fn account_summary_from_store( + store: &S, + organization_id: Uuid, +) -> Result { + match store.read(|kernel| kernel.get_account_summary(organization_id, 50)) { + Ok(Ok(summary)) => Ok(summary), + Ok(Err(error)) => Err(error.to_string()), + Err(StorageError::LockPoisoned) => Err("storage lock poisoned".to_string()), + Err(StorageError::Kernel(error)) => Err(error.to_string()), + Err(StorageError::ConnectionFailed { message, .. }) => Err(message), + Err(StorageError::SerializationFailed { message }) => Err(message), + Err(StorageError::Timeout { operation }) => Err(operation), + Err(StorageError::RuntimeStore { message }) => Err(message), + } +} + +fn knowledge_entries_from_summary(summary: &AccountSummary) -> Vec { + let mut entries = Vec::new(); + entries.push( + KnowledgeEntry::new( + format!("Account {}", summary.organization.name), + format!( + "Industry: {:?}. Website: {:?}. Tags: {}", + summary.organization.industry, + summary.organization.website, + summary.organization.tags.join(", ") + ), + ) + .with_category("organization") + .with_tags(["organization", "renewal"]) + .with_source(format!("crm://organization/{}", summary.organization.id)), + ); + entries.extend(summary.documents.iter().map(|document| { + KnowledgeEntry::new(&document.title, format!("Document URI {}", document.uri)) + .with_category("document") + .with_tags(["document", "renewal"]) + .with_source(format!("crm://document/{}", document.id)) + })); + entries.extend(summary.facts.iter().map(|fact| { + KnowledgeEntry::new("CRM fact", &fact.statement) + .with_category("fact") + .with_tags(["fact", "renewal"]) + .with_source(format!("crm://fact/{}", fact.id)) + })); + entries.extend(summary.recent_timeline.iter().map(|entry| { + KnowledgeEntry::new(&entry.headline, &entry.body) + .with_category("timeline") + .with_tags(["timeline", "renewal"]) + .with_source(format!("crm://timeline/{}", entry.id)) + })); + entries +} + +fn renewal_brief_from_result(result: &ConvergeResult) -> Result { + payload_from_result(result, ContextKey::Strategies, RENEWAL_BRIEF_FACT_ID) +} + +fn renewal_terms_from_result(result: &ConvergeResult) -> Result { + payload_from_result(result, ContextKey::Strategies, RENEWAL_TERMS_FACT_ID) +} + +fn renewal_signals_from_result( + result: &ConvergeResult, +) -> Result, Status> { + let signals = result + .context + .get(ContextKey::Signals) + .iter() + .filter(|fact| fact.id().starts_with("renewal:signal:")) + .map(|fact| { + serde_json::from_str::(fact.text().unwrap_or_default()).map_err( + |error| { + Status::internal(format!( + "invalid renewal signal payload {}: {error}", + fact.id() + )) + }, + ) + }) + .collect::, _>>()?; + Ok(signals) +} + +fn renewal_signals_from_context(ctx: &dyn ContextView) -> Vec { + ctx.get(ContextKey::Signals) + .iter() + .filter(|fact| fact.id().starts_with("renewal:signal:")) + .filter_map(|fact| { + serde_json::from_str::(fact.text().unwrap_or_default()).ok() + }) + .collect() +} + +fn summarize(content: &str, max_len: usize) -> String { + let content = content.trim(); + if content.len() <= max_len { + content.to_string() + } else { + format!("{}...", &content[..max_len.saturating_sub(3)]) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + use application_kernel::{ + Actor, FactRecord, OrganizationLifecycle, OrganizationUpsert, RecordKind, RecordRef, + }; + use application_storage::InMemoryKernelStore; + + #[tokio::test] + async fn match_renewal_context_executes_end_to_end() { + let store = InMemoryKernelStore::default_local(); + let runtime_stores = application_storage::AppRuntimeStores { + context: application_storage::AppContextStore::Memory( + application_storage::InMemoryContextStore::new(), + ), + experience: application_storage::AppExperienceStore::Memory( + application_storage::InMemoryExperienceStoreAdapter::new(), + ), + }; + let actor = Actor::system(); + let organization_id = store + .write(|kernel| { + let organization = kernel.upsert_organization( + OrganizationUpsert { + organization_id: None, + name: "Acme Renewals".to_string(), + external_key: None, + website: Some("https://acme.example".to_string()), + industry: Some("Software".to_string()), + lifecycle: OrganizationLifecycle::Active, + owner_user_id: None, + tags: vec!["renewal".to_string()], + }, + actor.clone(), + )?; + let related_to = vec![RecordRef { + kind: RecordKind::Organization, + id: organization.id, + }]; + let _ = kernel.record_fact( + FactRecord { + statement: "Customer mentioned competitor evaluation in a QBR.".to_string(), + confidence_bps: 8_400, + related_to: related_to.clone(), + source_note_id: None, + }, + actor.clone(), + )?; + let _ = kernel.attach_document( + DocumentAttach { + title: "Q3 incident review".to_string(), + media_type: "text/plain".to_string(), + uri: "converge://docs/q3-incident-review.txt".to_string(), + status: DocumentStatus::Verified, + related_to, + }, + actor.clone(), + )?; + Ok(organization.id) + }) + .expect("seed organization"); + + let inputs = MatchRenewalContextInput { + organization_id, + opportunity_id: None, + }; + + let execution = execute(&store, &runtime_stores, inputs, actor, true) + .await + .expect("truth should execute"); + assert!(execution.result.converged); + assert!( + execution + .result + .criteria_outcomes + .iter() + .all(|outcome| matches!( + outcome.result, + converge_kernel::CriterionResult::Met { .. } + )) + ); + + let projection = execution.projection.expect("projection should persist"); + assert!(projection.organization.is_some()); + assert_eq!(projection.documents.len(), 1); + assert!(!projection.facts.is_empty()); + } +} diff --git a/showcase/crm-helm/src/truths/plan_outbound_campaign.rs b/showcase/crm-helm/src/truths/plan_outbound_campaign.rs new file mode 100644 index 0000000..2748a74 --- /dev/null +++ b/showcase/crm-helm/src/truths/plan_outbound_campaign.rs @@ -0,0 +1,879 @@ +//! NOT YET WIRED — see Phase 8 or Phase 6b for integration. (Originally lived in helms/crates/application-server/src/truth_runtime/.) +use std::collections::HashMap; + +use application_kernel::{ + ActivityAppend, ActivityOutcome, Actor as CrmActor, DocumentAttach, DocumentStatus, FactRecord, + RecordKind, RecordRef, WorkflowCaseAdvance, WorkflowCaseCreate, WorkflowPriority, + WorkflowState, +}; +use application_storage::{KernelStore, StoreWriteResult}; +use converge_kernel::{ContextState as Context, ConvergeResult, Engine}; +use converge_optimization::Pack; +use converge_optimization::packs::lead_routing::{ + Lead as RoutingLead, LeadRoutingInput, LeadRoutingOutput, LeadRoutingPack, RoutingConfig, + SalesRep, +}; +use converge_pack::gate::{ObjectiveSpec, ProblemSpec}; +use converge_pack::{AgentEffect, Context as ContextView, ContextKey, Suggestor}; +use serde::{Deserialize, Serialize}; +use tonic::Status; +use truth_catalog::{ + PlanOutboundCampaignEvaluator, + admission::{admit_truth_intent, default_helms_capabilities, select_formation_for_intent}, + converge_binding_for_truth, +}; +use uuid::Uuid; + +use super::{ + TruthExecutionArtifacts, TruthProjection, + common::{ + converge_confidence_to_bps, has_fact_id, optional_i64, payload_from_result, required_input, + }, + domain_event_kind_name, status_from_storage, +}; + +const COMMERCIAL_PACK_ID: &str = "prio-commercial-pack"; +const WORK_PACK_ID: &str = "prio-work-pack"; +const REVENUE_PACK_ID: &str = "prio-revenue-pack"; +const ROUTING_INPUT_FACT_ID: &str = "campaign:lead-routing-input"; +const CAPACITY_STATUS_FACT_ID: &str = "campaign:capacity-status"; +const CAMPAIGN_PLAN_FACT_ID: &str = "campaign:plan"; +const BUDGET_STATUS_FACT_ID: &str = "campaign:budget-status"; +const PLAN_PROVENANCE: &str = "prio.plan-outbound-campaign.optimization"; +const CAPACITY_PROVENANCE: &str = "prio.plan-outbound-campaign.capacity"; +const BUDGET_PROVENANCE: &str = "prio.plan-outbound-campaign.budget"; + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct CampaignProspectSeed { + lead_id: String, + organization_id: Option, + score: f64, + territory: String, + segment: String, + #[serde(default)] + required_skills: Vec, + #[serde(default)] + estimated_value: f64, + #[serde(default = "default_priority")] + priority: i32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct CampaignRepSeed { + rep_id: String, + name: String, + capacity: i64, + current_load: i64, + territories: Vec, + segments: Vec, + #[serde(default)] + skills: Vec, + #[serde(default = "default_performance_score")] + performance_score: f64, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct CapacityStatusPayload { + total_capacity: i64, + total_available_capacity: i64, + rep_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct CampaignAssignmentPayload { + lead_id: String, + organization_id: Option, + rep_id: String, + rep_name: String, + fit_score: f64, + rationale: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct CampaignPlanPayload { + campaign_name: String, + summary: String, + assignments: Vec, + unassigned_leads: Vec, + average_fit_score: f64, + confidence_bps: u16, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct BudgetStatusPayload { + within_budget: bool, + estimated_spend_minor: i64, + budget_minor: i64, + approval_required: bool, +} + +#[derive(Clone)] +struct ProspectPoolAgent { + prospects: Vec, + reps: Vec, +} + +struct RepCapacityAgent { + reps: Vec, +} + +struct CampaignSolverAgent { + campaign_name: String, + prospects: Vec, +} + +struct BudgetGuardAgent { + campaign_budget_minor: i64, + outreach_cost_minor: i64, +} + +#[derive(Debug, Clone)] +pub struct PlanOutboundCampaignInput { + pub campaign_name: String, + pub prospects_json: String, + pub reps_json: String, + pub campaign_budget_minor: i64, + pub outreach_cost_minor: Option, +} + +impl PlanOutboundCampaignInput { + pub fn from_map(inputs: &HashMap) -> Result { + Ok(Self { + campaign_name: required_input(inputs, "campaign_name")?.to_string(), + prospects_json: required_input(inputs, "prospects_json")?.to_string(), + reps_json: required_input(inputs, "reps_json")?.to_string(), + campaign_budget_minor: required_input(inputs, "campaign_budget_minor")? + .parse() + .map_err(|e| { + Status::invalid_argument(format!("invalid campaign_budget_minor: {e}")) + })?, + outreach_cost_minor: optional_i64(inputs, "outreach_cost_minor"), + }) + } +} + +pub(super) async fn execute( + store: &S, + runtime_stores: &application_storage::AppRuntimeStores, + inputs: PlanOutboundCampaignInput, + actor: CrmActor, + persist_projection: bool, +) -> Result { + let binding = converge_binding_for_truth("plan-outbound-campaign") + .ok_or_else(|| Status::not_found("truth not found: plan-outbound-campaign"))?; + + let campaign_name = inputs.campaign_name.clone(); + let prospects = prospects_from_inputs(&inputs)?; + let reps = reps_from_inputs(&inputs)?; + let campaign_budget_minor = inputs.campaign_budget_minor; + let outreach_cost_minor = inputs.outreach_cost_minor.unwrap_or(2_500); + + let mut engine = Engine::new(); + engine.register_suggestor_in_pack( + COMMERCIAL_PACK_ID, + ProspectPoolAgent { + prospects: prospects.clone(), + reps: reps.clone(), + }, + ); + engine.register_suggestor_in_pack(WORK_PACK_ID, RepCapacityAgent { reps: reps.clone() }); + engine.register_suggestor_in_pack( + COMMERCIAL_PACK_ID, + CampaignSolverAgent { + campaign_name, + prospects, + }, + ); + engine.register_suggestor_in_pack( + REVENUE_PACK_ID, + BudgetGuardAgent { + campaign_budget_minor, + outreach_cost_minor, + }, + ); + + let mut seed_ctx = seed_context()?; + let intent = admit_truth_intent( + "plan-outbound-campaign", + &actor.actor_id, + "truth:plan-outbound-campaign", + &mut seed_ctx, + ) + .map_err(|e| Status::internal(format!("admit intent failed: {e}")))?; + let selection = select_formation_for_intent(&intent, &default_helms_capabilities()) + .map_err(|e| Status::internal(format!("formation selection failed: {e}")))?; + tracing::info!( + truth = "plan-outbound-campaign", + primary = %selection.primary_template_id, + alternates = ?selection.alternate_template_ids, + "formation selected" + ); + + let runtime_ctx = super::RuntimeContext { + scope_id: slug(&inputs.campaign_name), + }; + let (result, experience_events) = super::run_engine_with_runtime( + runtime_stores, + &mut engine, + &runtime_ctx, + seed_ctx, + &binding.intent, + std::sync::Arc::new(PlanOutboundCampaignEvaluator), + ) + .await?; + + let projection = if persist_projection { + Some(project(store, &inputs, &result, actor)?) + } else { + None + }; + + Ok(TruthExecutionArtifacts { + result, + experience_events, + projection, + runtime_scope_id: runtime_ctx.scope_id, + }) +} + +#[async_trait::async_trait] +impl Suggestor for ProspectPoolAgent { + fn name(&self) -> &str { + "ProspectPoolAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Seeds] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + ctx.has(ContextKey::Seeds) + && !has_fact_id(ctx, ContextKey::Proposals, ROUTING_INPUT_FACT_ID) + } + + async fn execute(&self, _ctx: &dyn ContextView) -> AgentEffect { + let routing_input = LeadRoutingInput { + leads: self + .prospects + .iter() + .map(|prospect| RoutingLead { + id: prospect.lead_id.clone(), + score: prospect.score, + territory: prospect.territory.clone(), + segment: prospect.segment.clone(), + required_skills: prospect.required_skills.clone(), + estimated_value: prospect.estimated_value, + priority: prospect.priority, + }) + .collect(), + reps: self + .reps + .iter() + .map(|rep| SalesRep { + id: rep.rep_id.clone(), + name: rep.name.clone(), + capacity: rep.capacity, + current_load: rep.current_load, + territories: rep.territories.clone(), + segments: rep.segments.clone(), + skills: rep.skills.clone(), + performance_score: rep.performance_score, + }) + .collect(), + config: RoutingConfig::default(), + }; + + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Proposals, + ROUTING_INPUT_FACT_ID.to_string(), + serde_json::to_string(&routing_input).unwrap_or_default(), + PLAN_PROVENANCE.to_string(), + ) + .with_confidence(1.0), + ) + } +} + +#[async_trait::async_trait] +impl Suggestor for RepCapacityAgent { + fn name(&self) -> &str { + "RepCapacityAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Seeds] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + ctx.has(ContextKey::Seeds) + && !has_fact_id(ctx, ContextKey::Signals, CAPACITY_STATUS_FACT_ID) + } + + async fn execute(&self, _ctx: &dyn ContextView) -> AgentEffect { + let payload = CapacityStatusPayload { + total_capacity: self.reps.iter().map(|rep| rep.capacity).sum(), + total_available_capacity: self + .reps + .iter() + .map(|rep| (rep.capacity - rep.current_load).max(0)) + .sum(), + rep_count: self.reps.len(), + }; + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Signals, + CAPACITY_STATUS_FACT_ID.to_string(), + serde_json::to_string(&payload).unwrap_or_default(), + CAPACITY_PROVENANCE.to_string(), + ) + .with_confidence(1.0), + ) + } +} + +#[async_trait::async_trait] +impl Suggestor for CampaignSolverAgent { + fn name(&self) -> &str { + "CampaignSolverAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Proposals, ContextKey::Signals] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + has_fact_id(ctx, ContextKey::Proposals, ROUTING_INPUT_FACT_ID) + && has_fact_id(ctx, ContextKey::Signals, CAPACITY_STATUS_FACT_ID) + && !has_fact_id(ctx, ContextKey::Strategies, CAMPAIGN_PLAN_FACT_ID) + } + + async fn execute(&self, ctx: &dyn ContextView) -> AgentEffect { + let Some(input_fact) = ctx + .get(ContextKey::Proposals) + .iter() + .find(|fact| fact.id() == ROUTING_INPUT_FACT_ID) + else { + return AgentEffect::empty(); + }; + let routing_input = + match serde_json::from_str::(input_fact.text().unwrap_or_default()) { + Ok(input) => input, + Err(error) => { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "campaign:plan:error", + error.to_string(), + "diagnostic", + ) + .with_confidence(1.0), + ); + } + }; + + let spec = match ProblemSpec::builder( + format!("campaign-{}", slug(&self.campaign_name)), + "crm.prio.ai", + ) + .objective(ObjectiveSpec::maximize("conversion")) + .inputs(&routing_input) + .and_then(|builder| builder.build()) + { + Ok(spec) => spec, + Err(error) => { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "campaign:plan:error", + error.to_string(), + "diagnostic", + ) + .with_confidence(1.0), + ); + } + }; + + let pack = LeadRoutingPack; + let solved = match pack.solve(&spec) { + Ok(result) => result, + Err(error) => { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "campaign:plan:error", + error.to_string(), + "diagnostic", + ) + .with_confidence(1.0), + ); + } + }; + let output = match solved.plan.plan_as::() { + Ok(output) => output, + Err(error) => { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "campaign:plan:error", + error.to_string(), + "diagnostic", + ) + .with_confidence(1.0), + ); + } + }; + let plan_payload = CampaignPlanPayload { + campaign_name: self.campaign_name.clone(), + summary: output.summary(), + assignments: output + .assignments + .iter() + .map(|assignment| CampaignAssignmentPayload { + lead_id: assignment.lead_id.clone(), + organization_id: self + .prospects + .iter() + .find(|prospect| prospect.lead_id == assignment.lead_id) + .and_then(|prospect| prospect.organization_id), + rep_id: assignment.rep_id.clone(), + rep_name: assignment.rep_name.clone(), + fit_score: assignment.fit_score, + rationale: assignment.scoring_rationale.explanation.clone(), + }) + .collect(), + unassigned_leads: output + .unassigned + .iter() + .map(|lead| lead.lead_id.clone()) + .collect(), + average_fit_score: output.stats.average_fit_score, + confidence_bps: converge_confidence_to_bps(solved.plan.confidence()), + }; + + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Strategies, + CAMPAIGN_PLAN_FACT_ID.to_string(), + serde_json::to_string(&plan_payload).unwrap_or_default(), + PLAN_PROVENANCE.to_string(), + ) + .with_confidence(solved.plan.confidence()), + ) + } +} + +#[async_trait::async_trait] +impl Suggestor for BudgetGuardAgent { + fn name(&self) -> &str { + "BudgetGuardAgent" + } + + fn dependencies(&self) -> &[ContextKey] { + &[ContextKey::Strategies] + } + + fn accepts(&self, ctx: &dyn ContextView) -> bool { + has_fact_id(ctx, ContextKey::Strategies, CAMPAIGN_PLAN_FACT_ID) + && !has_fact_id(ctx, ContextKey::Evaluations, BUDGET_STATUS_FACT_ID) + } + + async fn execute(&self, ctx: &dyn ContextView) -> AgentEffect { + let Some(plan_fact) = ctx + .get(ContextKey::Strategies) + .iter() + .find(|fact| fact.id() == CAMPAIGN_PLAN_FACT_ID) + else { + return AgentEffect::empty(); + }; + let plan = + match serde_json::from_str::(plan_fact.text().unwrap_or_default()) + { + Ok(plan) => plan, + Err(error) => { + return AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Diagnostic, + "campaign:plan:error", + error.to_string(), + "diagnostic", + ) + .with_confidence(1.0), + ); + } + }; + + let estimated_spend_minor = plan.assignments.len() as i64 * self.outreach_cost_minor; + let within_budget = estimated_spend_minor <= self.campaign_budget_minor; + let payload = BudgetStatusPayload { + within_budget, + estimated_spend_minor, + budget_minor: self.campaign_budget_minor, + approval_required: !within_budget, + }; + AgentEffect::with_proposal( + crate::truth_runtime::common::proposed_text_fact( + ContextKey::Evaluations, + BUDGET_STATUS_FACT_ID.to_string(), + serde_json::to_string(&payload).unwrap_or_default(), + BUDGET_PROVENANCE.to_string(), + ) + .with_confidence(1.0), + ) + } +} + +fn project( + store: &S, + inputs: &PlanOutboundCampaignInput, + result: &ConvergeResult, + actor: CrmActor, +) -> Result { + let campaign_name = inputs.campaign_name.clone(); + let plan = campaign_plan_from_result(result)?; + let budget = budget_status_from_result(result)?; + let related_to = related_record_refs(&plan.assignments); + if related_to.is_empty() { + return Err(Status::invalid_argument( + "campaign projection requires organization_id on at least one prospect", + )); + } + + let StoreWriteResult { value, events } = store + .write_with_events(|kernel| { + let workflow_case = kernel.create_workflow_case( + WorkflowCaseCreate { + title: format!("Outbound campaign: {campaign_name}"), + priority: WorkflowPriority::High, + owner_user_id: None, + related_to: related_to.clone(), + }, + actor.clone(), + )?; + + let workflow_case = if budget.approval_required { + kernel.advance_workflow_case( + WorkflowCaseAdvance { + workflow_case_id: workflow_case.id, + state: WorkflowState::AwaitingApproval, + }, + actor.clone(), + )? + } else { + workflow_case + }; + + let mut document_related_to = related_to.clone(); + document_related_to.push(RecordRef { + kind: RecordKind::WorkflowCase, + id: workflow_case.id, + }); + let document = kernel.attach_document( + DocumentAttach { + title: format!("Campaign plan: {campaign_name}"), + media_type: "application/json".to_string(), + uri: format!( + "converge://truths/plan-outbound-campaign/{}/plan.json", + slug(&campaign_name) + ), + status: DocumentStatus::Draft, + related_to: document_related_to.clone(), + }, + actor.clone(), + )?; + + for assignment in &plan.assignments { + let mut activity_related_to = assignment + .organization_id + .map(|organization_id| { + vec![RecordRef { + kind: RecordKind::Organization, + id: organization_id, + }] + }) + .unwrap_or_default(); + activity_related_to.push(RecordRef { + kind: RecordKind::WorkflowCase, + id: workflow_case.id, + }); + let _ = kernel.append_activity( + ActivityAppend { + subject: format!("Outbound assignment for {}", assignment.lead_id), + details: format!( + "Assigned to {} ({}) with fit {:.1}: {}", + assignment.rep_name, + assignment.rep_id, + assignment.fit_score, + assignment.rationale + ), + related_to: activity_related_to, + outcome: ActivityOutcome::Waiting, + occurred_at: None, + next_action_due_at: None, + }, + actor.clone(), + )?; + } + + let plan_fact = kernel.record_fact( + FactRecord { + statement: format!( + "Campaign plan {} assigns {} leads with average fit {:.1}", + campaign_name, + plan.assignments.len(), + plan.average_fit_score + ), + confidence_bps: plan.confidence_bps, + related_to: document_related_to.clone(), + source_note_id: None, + }, + actor.clone(), + )?; + let budget_fact = kernel.record_fact( + FactRecord { + statement: format!( + "Campaign spend estimate {} against budget {} ({})", + budget.estimated_spend_minor, + budget.budget_minor, + if budget.within_budget { + "within-budget" + } else { + "approval-required" + } + ), + confidence_bps: 10_000, + related_to: document_related_to, + source_note_id: None, + }, + actor, + )?; + + Ok((workflow_case, document, vec![plan_fact, budget_fact])) + }) + .map_err(status_from_storage)?; + + let (workflow_case, document, facts) = value; + Ok(TruthProjection { + organization: None, + person: None, + opportunity: None, + subscription: None, + entitlements: Vec::new(), + ledger_entries: Vec::new(), + documents: vec![document], + workflow_cases: vec![workflow_case], + facts, + domain_event_kinds: events.iter().map(domain_event_kind_name).collect(), + }) +} + +fn seed_context() -> Result { + let mut context = Context::new(); + context + .add_input( + ContextKey::Seeds, + "plan-outbound-campaign:seed", + "campaign-seed", + ) + .map_err(|error| Status::failed_precondition(error.to_string()))?; + Ok(context) +} + +fn prospects_from_inputs( + inputs: &PlanOutboundCampaignInput, +) -> Result, Status> { + serde_json::from_str(&inputs.prospects_json) + .map_err(|error| Status::invalid_argument(format!("invalid prospects_json: {error}"))) +} + +fn reps_from_inputs(inputs: &PlanOutboundCampaignInput) -> Result, Status> { + serde_json::from_str(&inputs.reps_json) + .map_err(|error| Status::invalid_argument(format!("invalid reps_json: {error}"))) +} + +fn campaign_plan_from_result(result: &ConvergeResult) -> Result { + payload_from_result(result, ContextKey::Strategies, CAMPAIGN_PLAN_FACT_ID) +} + +fn budget_status_from_result(result: &ConvergeResult) -> Result { + payload_from_result(result, ContextKey::Evaluations, BUDGET_STATUS_FACT_ID) +} + +fn related_record_refs(assignments: &[CampaignAssignmentPayload]) -> Vec { + let mut refs = assignments + .iter() + .filter_map(|assignment| assignment.organization_id) + .map(|organization_id| RecordRef { + kind: RecordKind::Organization, + id: organization_id, + }) + .collect::>(); + refs.sort_by_key(|reference| reference.id); + refs.dedup_by_key(|reference| reference.id); + refs +} + +fn slug(value: &str) -> String { + let mut slug = String::new(); + let mut last_was_dash = false; + + for ch in value.chars() { + if ch.is_ascii_alphanumeric() { + slug.push(ch.to_ascii_lowercase()); + last_was_dash = false; + } else if !slug.is_empty() && !last_was_dash { + slug.push('-'); + last_was_dash = true; + } + } + + if slug.ends_with('-') { + slug.pop(); + } + + if slug.is_empty() { + "campaign".to_string() + } else { + slug + } +} + +fn default_priority() -> i32 { + 5 +} + +fn default_performance_score() -> f64 { + 50.0 +} + +#[cfg(test)] +mod tests { + use super::*; + + use application_kernel::Actor; + use application_kernel::{OrganizationLifecycle, OrganizationUpsert}; + use application_storage::InMemoryKernelStore; + + #[tokio::test] + async fn plan_outbound_campaign_executes_end_to_end() { + let store = InMemoryKernelStore::default_local(); + let runtime_stores = application_storage::AppRuntimeStores { + context: application_storage::AppContextStore::Memory( + application_storage::InMemoryContextStore::new(), + ), + experience: application_storage::AppExperienceStore::Memory( + application_storage::InMemoryExperienceStoreAdapter::new(), + ), + }; + let actor = Actor::system(); + let west_org_id = store + .write(|kernel| { + kernel + .upsert_organization( + OrganizationUpsert { + organization_id: None, + name: "West Prospect".to_string(), + external_key: None, + website: None, + industry: None, + lifecycle: OrganizationLifecycle::Prospect, + owner_user_id: None, + tags: vec!["campaign".to_string()], + }, + actor.clone(), + ) + .map(|organization| organization.id) + }) + .expect("west prospect seed"); + let east_org_id = store + .write(|kernel| { + kernel + .upsert_organization( + OrganizationUpsert { + organization_id: None, + name: "East Prospect".to_string(), + external_key: None, + website: None, + industry: None, + lifecycle: OrganizationLifecycle::Prospect, + owner_user_id: None, + tags: vec!["campaign".to_string()], + }, + actor.clone(), + ) + .map(|organization| organization.id) + }) + .expect("east prospect seed"); + let inputs = PlanOutboundCampaignInput { + campaign_name: "Q2 outbound".to_string(), + prospects_json: serde_json::json!([ + { + "lead_id": "lead-1", + "organization_id": west_org_id, + "score": 88.0, + "territory": "west", + "segment": "enterprise", + "required_skills": ["cloud"], + "estimated_value": 120000.0, + "priority": 1 + }, + { + "lead_id": "lead-2", + "organization_id": east_org_id, + "score": 72.0, + "territory": "east", + "segment": "smb", + "required_skills": [], + "estimated_value": 25000.0, + "priority": 3 + } + ]) + .to_string(), + reps_json: serde_json::json!([ + { + "rep_id": "rep-1", + "name": "Alice", + "capacity": 5, + "current_load": 1, + "territories": ["west"], + "segments": ["enterprise"], + "skills": ["cloud"], + "performance_score": 92.0 + }, + { + "rep_id": "rep-2", + "name": "Bob", + "capacity": 5, + "current_load": 2, + "territories": ["east", "west"], + "segments": ["smb", "enterprise"], + "skills": [], + "performance_score": 80.0 + } + ]) + .to_string(), + campaign_budget_minor: 6000, + outreach_cost_minor: Some(2500), + }; + + let execution = execute(&store, &runtime_stores, inputs, actor, true) + .await + .expect("truth should execute"); + assert!(execution.result.converged); + assert!( + execution + .result + .criteria_outcomes + .iter() + .all(|outcome| matches!( + outcome.result, + converge_kernel::CriterionResult::Met { .. } + )) + ); + + let projection = execution.projection.expect("projection should persist"); + assert_eq!(projection.workflow_cases.len(), 1); + assert_eq!(projection.documents.len(), 1); + assert_eq!(projection.facts.len(), 2); + } +} diff --git a/showcase/crm-helm/src/workbench.rs b/showcase/crm-helm/src/workbench.rs new file mode 100644 index 0000000..b64098e --- /dev/null +++ b/showcase/crm-helm/src/workbench.rs @@ -0,0 +1 @@ +//! CRM Workbench module — filled in Phase 6 from helm-operator-control workbench surface diff --git a/showcase/crm-helm/src/workflow.rs b/showcase/crm-helm/src/workflow.rs new file mode 100644 index 0000000..248c530 --- /dev/null +++ b/showcase/crm-helm/src/workflow.rs @@ -0,0 +1,113 @@ +//! CRM Workflow module — lead/case state machines as a HelmModule. +//! +//! Moved from helms/crates/application-server/src/service.rs (WorkflowGrpc). + +use application_kernel::{WorkflowCaseAdvance, WorkflowCaseCreate}; +use application_storage::{AppKernelStore, InMemoryKernelStore, KernelStore}; +use async_trait::async_trait; +use runway_app_host::HelmModule; +use tonic::{Request, Response, Status}; + +use crate::proto::{common as pb, workflow as workflow_pb}; +use crate::shared::{ + actor_from_proto, parse_uuid, proto_workflow_case, record_ref_from_proto, status_from_storage, + workflow_priority_from_proto, workflow_state_from_proto, +}; + +// --------------------------------------------------------------------------- +// gRPC service struct +// --------------------------------------------------------------------------- + +#[derive(Clone)] +pub struct WorkflowGrpc { + store: S, +} + +impl WorkflowGrpc { + #[allow(dead_code)] + pub fn new(store: S) -> Self { + Self { store } + } +} + +#[tonic::async_trait] +impl workflow_pb::workflow_service_server::WorkflowService for WorkflowGrpc +where + S: KernelStore, +{ + async fn create_workflow_case( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let related_to = request + .related_to + .into_iter() + .map(record_ref_from_proto) + .collect::, _>>()?; + let workflow_case = self + .store + .write(|kernel| { + kernel.create_workflow_case( + WorkflowCaseCreate { + title: request.title, + priority: workflow_priority_from_proto(request.priority), + owner_user_id: request.owner_user_id, + related_to, + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_workflow_case(workflow_case))) + } + + async fn advance_workflow_case( + &self, + request: Request, + ) -> Result, Status> { + let request = request.into_inner(); + let workflow_case_id = parse_uuid(&request.workflow_case_id)?; + let workflow_case = self + .store + .write(|kernel| { + kernel.advance_workflow_case( + WorkflowCaseAdvance { + workflow_case_id, + state: workflow_state_from_proto(request.state), + }, + actor_from_proto(request.actor), + ) + }) + .map_err(status_from_storage)?; + Ok(Response::new(proto_workflow_case(workflow_case))) + } +} + +// --------------------------------------------------------------------------- +// HelmModule wrapper +// --------------------------------------------------------------------------- + +pub struct WorkflowModule {} + +impl WorkflowModule { + pub fn new(_store: AppKernelStore) -> Self { + Self {} + } + + #[allow(dead_code)] + pub fn in_memory() -> Self { + Self::new(AppKernelStore::Memory(InMemoryKernelStore::default_local())) + } +} + +#[async_trait] +impl HelmModule for WorkflowModule { + fn module_id(&self) -> &'static str { + "crm.workflow" + } + + async fn init(&self) -> anyhow::Result<()> { + Ok(()) + } +}