fix(db): resolve workflow schedule tenant from DB

Co-authored-by: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co>
Signed-off-by: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf
2026-06-26 13:30:56 -04:00
parent 48c53a5d7e
commit 1fa3d837f1
7 changed files with 137 additions and 60 deletions
+3 -3
View File
@@ -1122,8 +1122,8 @@ pub async fn release_due_reminder(
event_id: &[u8], event_id: &[u8],
event_created_at: DateTime<Utc>, event_created_at: DateTime<Utc>,
delivery_stamp: i64, delivery_stamp: i64,
) -> Result<()> { ) -> Result<bool> {
sqlx::query( let result = sqlx::query(
r#" r#"
UPDATE events UPDATE events
SET delivered_at = NULL SET delivered_at = NULL
@@ -1138,7 +1138,7 @@ pub async fn release_due_reminder(
.execute(pool) .execute(pool)
.await?; .await?;
Ok(()) Ok(result.rows_affected() == 1)
} }
#[cfg(test)] #[cfg(test)]
+8 -15
View File
@@ -664,7 +664,7 @@ impl Db {
event_id: &[u8], event_id: &[u8],
event_created_at: chrono::DateTime<chrono::Utc>, event_created_at: chrono::DateTime<chrono::Utc>,
delivery_stamp: i64, delivery_stamp: i64,
) -> Result<()> { ) -> Result<bool> {
event::release_due_reminder(&self.pool, event_id, event_created_at, delivery_stamp).await event::release_due_reminder(&self.pool, event_id, event_created_at, delivery_stamp).await
} }
@@ -1203,6 +1203,7 @@ impl Db {
/// Create a new workflow. /// Create a new workflow.
pub async fn create_workflow( pub async fn create_workflow(
&self, &self,
community_id: CommunityId,
channel_id: Option<Uuid>, channel_id: Option<Uuid>,
owner_pubkey: &[u8], owner_pubkey: &[u8],
name: &str, name: &str,
@@ -1211,6 +1212,7 @@ impl Db {
) -> Result<Uuid> { ) -> Result<Uuid> {
workflow::create_workflow( workflow::create_workflow(
&self.pool, &self.pool,
community_id,
channel_id, channel_id,
owner_pubkey, owner_pubkey,
name, name,
@@ -1250,43 +1252,34 @@ impl Db {
/// Claim a scheduled workflow fire for an authoritative schedule instant. /// Claim a scheduled workflow fire for an authoritative schedule instant.
/// ///
/// Returns `Some` only for the first pod to claim `(community_id, /// Returns `Some` only for the first pod to claim `(workflow_id,
/// workflow_id, scheduled_for)`; all other pods must skip creating a run. /// scheduled_for)`; all other pods must skip creating a run. The claim SQL
/// resolves `community_id` from the workflow row; callers never supply it.
pub async fn claim_scheduled_workflow_fire( pub async fn claim_scheduled_workflow_fire(
&self, &self,
community_id: CommunityId,
workflow_id: Uuid, workflow_id: Uuid,
scheduled_for: chrono::DateTime<chrono::Utc>, scheduled_for: chrono::DateTime<chrono::Utc>,
) -> Result<Option<workflow::ScheduledWorkflowFireClaim>> { ) -> Result<Option<workflow::ScheduledWorkflowFireClaim>> {
workflow::claim_scheduled_workflow_fire( workflow::claim_scheduled_workflow_fire(&self.pool, workflow_id, scheduled_for).await
&self.pool,
community_id,
workflow_id,
scheduled_for,
)
.await
} }
/// Fetch the latest claimed schedule instant for interval trigger anchoring. /// Fetch the latest claimed schedule instant for interval trigger anchoring.
pub async fn latest_scheduled_workflow_fire( pub async fn latest_scheduled_workflow_fire(
&self, &self,
community_id: CommunityId,
workflow_id: Uuid, workflow_id: Uuid,
) -> Result<Option<chrono::DateTime<chrono::Utc>>> { ) -> Result<Option<chrono::DateTime<chrono::Utc>>> {
workflow::latest_scheduled_workflow_fire(&self.pool, community_id, workflow_id).await workflow::latest_scheduled_workflow_fire(&self.pool, workflow_id).await
} }
/// Attach the workflow run id created from a won scheduled-fire claim. /// Attach the workflow run id created from a won scheduled-fire claim.
pub async fn attach_scheduled_workflow_run( pub async fn attach_scheduled_workflow_run(
&self, &self,
community_id: CommunityId,
workflow_id: Uuid, workflow_id: Uuid,
scheduled_for: chrono::DateTime<chrono::Utc>, scheduled_for: chrono::DateTime<chrono::Utc>,
workflow_run_id: Uuid, workflow_run_id: Uuid,
) -> Result<bool> { ) -> Result<bool> {
workflow::attach_scheduled_workflow_run( workflow::attach_scheduled_workflow_run(
&self.pool, &self.pool,
community_id,
workflow_id, workflow_id,
scheduled_for, scheduled_for,
workflow_run_id, workflow_run_id,
+23 -2
View File
@@ -145,6 +145,20 @@ mod tests {
.contains("CREATE TABLE scheduled_workflow_fires"), .contains("CREATE TABLE scheduled_workflow_fires"),
"initial schema migration should include workflow cron claim table" "initial schema migration should include workflow cron claim table"
); );
assert!(
migrations[0]
.sql
.as_str()
.contains("community_id UUID NOT NULL REFERENCES communities(id)"),
"workflows should carry row-owned community_id"
);
assert!(
migrations[0]
.sql
.as_str()
.contains("trg_workflows_community_id_immutable"),
"workflow community_id should be immutable after insert"
);
assert!( assert!(
migrations[0] migrations[0]
.sql .sql
@@ -156,8 +170,15 @@ mod tests {
migrations[0] migrations[0]
.sql .sql
.as_str() .as_str()
.contains("PRIMARY KEY (community_id, workflow_id, scheduled_for)"), .contains("PRIMARY KEY (workflow_id, scheduled_for)"),
"workflow cron claim uniqueness must include the community label" "workflow cron claim uniqueness should be the globally unique workflow plus schedule instant"
);
assert!(
migrations[0]
.sql
.as_str()
.contains("FOREIGN KEY (community_id, workflow_id)"),
"workflow cron claim table should tie community_id to the workflow row"
); );
assert!( assert!(
migrations[0].sql.as_str().contains("CREATE TABLE channels"), migrations[0].sql.as_str().contains("CREATE TABLE channels"),
+36 -27
View File
@@ -165,6 +165,8 @@ impl FromStr for ApprovalStatus {
pub struct WorkflowRecord { pub struct WorkflowRecord {
/// Unique workflow identifier. /// Unique workflow identifier.
pub id: Uuid, pub id: Uuid,
/// Server-resolved community that owns this workflow.
pub community_id: CommunityId,
/// Human-readable workflow name. /// Human-readable workflow name.
pub name: String, pub name: String,
/// Compressed public key bytes of the workflow owner. /// Compressed public key bytes of the workflow owner.
@@ -215,9 +217,9 @@ pub struct WorkflowRunRecord {
/// A winning scheduled workflow fire claim. /// A winning scheduled workflow fire claim.
/// ///
/// The primary identity is `(community_id, workflow_id, scheduled_for)`. The /// The primary identity is `(workflow_id, scheduled_for)`. `community_id` is
/// database also returns `claimed_at` so the workflow scheduler can log and audit /// resolved from the workflow row inside the claim SQL and returned for scoped
/// the exact claim row it won without relying on a per-pod clock. /// audit/logging; callers never supply it as a claim.
#[derive(Debug, Clone, PartialEq, Eq)] #[derive(Debug, Clone, PartialEq, Eq)]
pub struct ScheduledWorkflowFireClaim { pub struct ScheduledWorkflowFireClaim {
/// Community that owns this scheduled fire. /// Community that owns this scheduled fire.
@@ -263,6 +265,7 @@ pub struct ApprovalRecord {
/// New workflows start as `active` and `enabled = TRUE`. /// New workflows start as `active` and `enabled = TRUE`.
pub async fn create_workflow( pub async fn create_workflow(
pool: &PgPool, pool: &PgPool,
community_id: CommunityId,
channel_id: Option<Uuid>, channel_id: Option<Uuid>,
owner_pubkey: &[u8], owner_pubkey: &[u8],
name: &str, name: &str,
@@ -274,11 +277,12 @@ pub async fn create_workflow(
sqlx::query( sqlx::query(
r#" r#"
INSERT INTO workflows INSERT INTO workflows
(id, name, owner_pubkey, channel_id, definition, definition_hash, status, enabled) (id, community_id, name, owner_pubkey, channel_id, definition, definition_hash, status, enabled)
VALUES ($1, $2, $3, $4, $5::jsonb, $6, 'active', TRUE) VALUES ($1, $2, $3, $4, $5, $6::jsonb, $7, 'active', TRUE)
"#, "#,
) )
.bind(id) .bind(id)
.bind(community_id.as_uuid())
.bind(name) .bind(name)
.bind(owner_pubkey) .bind(owner_pubkey)
.bind(channel_id) .bind(channel_id)
@@ -294,7 +298,7 @@ pub async fn create_workflow(
pub async fn get_workflow(pool: &PgPool, id: Uuid) -> Result<WorkflowRecord> { pub async fn get_workflow(pool: &PgPool, id: Uuid) -> Result<WorkflowRecord> {
let row = sqlx::query( let row = sqlx::query(
r#" r#"
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash, SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
status::text AS status, enabled, created_at, updated_at status::text AS status, enabled, created_at, updated_at
FROM workflows FROM workflows
WHERE id = $1 WHERE id = $1
@@ -323,7 +327,7 @@ pub async fn list_channel_workflows(
let rows = sqlx::query( let rows = sqlx::query(
r#" r#"
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash, SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
status::text AS status, enabled, created_at, updated_at status::text AS status, enabled, created_at, updated_at
FROM workflows FROM workflows
WHERE channel_id = $1 WHERE channel_id = $1
@@ -352,7 +356,7 @@ pub async fn list_enabled_channel_workflows(
) -> Result<Vec<WorkflowRecord>> { ) -> Result<Vec<WorkflowRecord>> {
let rows = sqlx::query( let rows = sqlx::query(
r#" r#"
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash, SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
status::text AS status, enabled, created_at, updated_at status::text AS status, enabled, created_at, updated_at
FROM workflows FROM workflows
WHERE channel_id = $1 WHERE channel_id = $1
@@ -378,7 +382,7 @@ pub async fn list_enabled_channel_workflows(
pub async fn list_all_enabled_workflows(pool: &PgPool) -> Result<Vec<WorkflowRecord>> { pub async fn list_all_enabled_workflows(pool: &PgPool) -> Result<Vec<WorkflowRecord>> {
let rows = sqlx::query( let rows = sqlx::query(
r#" r#"
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash, SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
status::text AS status, enabled, created_at, updated_at status::text AS status, enabled, created_at, updated_at
FROM workflows FROM workflows
WHERE status = 'active' WHERE status = 'active'
@@ -397,27 +401,27 @@ pub async fn list_all_enabled_workflows(pool: &PgPool) -> Result<Vec<WorkflowRec
/// Claim a scheduled workflow fire for an authoritative schedule instant. /// Claim a scheduled workflow fire for an authoritative schedule instant.
/// ///
/// Returns `Some` only for the first pod that claims `(community_id, /// Returns `Some` only for the first pod that claims `(workflow_id,
/// workflow_id, scheduled_for)`. All other pods receive `None` and must skip /// scheduled_for)`. All other pods receive `None` and must skip creating a
/// creating a workflow run. The `scheduled_for` value must come from an external /// workflow run. The `scheduled_for` value must come from an external
/// schedule anchor (cron expression) or DB-authoritative interval anchor; a /// schedule anchor (cron expression) or DB-authoritative interval anchor; a
/// per-pod in-memory timestamp is not safe because different pods can compute /// per-pod in-memory timestamp is not safe because different pods can compute
/// different claim keys. /// different claim keys.
pub async fn claim_scheduled_workflow_fire( pub async fn claim_scheduled_workflow_fire(
pool: &PgPool, pool: &PgPool,
community_id: CommunityId,
workflow_id: Uuid, workflow_id: Uuid,
scheduled_for: DateTime<Utc>, scheduled_for: DateTime<Utc>,
) -> Result<Option<ScheduledWorkflowFireClaim>> { ) -> Result<Option<ScheduledWorkflowFireClaim>> {
let row = sqlx::query( let row = sqlx::query(
r#" r#"
INSERT INTO scheduled_workflow_fires (community_id, workflow_id, scheduled_for) INSERT INTO scheduled_workflow_fires (community_id, workflow_id, scheduled_for)
VALUES ($1, $2, $3) SELECT w.community_id, w.id, $2
ON CONFLICT (community_id, workflow_id, scheduled_for) DO NOTHING FROM workflows w
WHERE w.id = $1
ON CONFLICT (workflow_id, scheduled_for) DO NOTHING
RETURNING community_id, workflow_id, scheduled_for, claimed_at RETURNING community_id, workflow_id, scheduled_for, claimed_at
"#, "#,
) )
.bind(community_id.as_uuid())
.bind(workflow_id) .bind(workflow_id)
.bind(scheduled_for) .bind(scheduled_for)
.fetch_optional(pool) .fetch_optional(pool)
@@ -444,18 +448,15 @@ pub async fn claim_scheduled_workflow_fire(
/// row is the source of truth for schedule deduplication. /// row is the source of truth for schedule deduplication.
pub async fn latest_scheduled_workflow_fire( pub async fn latest_scheduled_workflow_fire(
pool: &PgPool, pool: &PgPool,
community_id: CommunityId,
workflow_id: Uuid, workflow_id: Uuid,
) -> Result<Option<DateTime<Utc>>> { ) -> Result<Option<DateTime<Utc>>> {
let row = sqlx::query( let row = sqlx::query(
r#" r#"
SELECT MAX(scheduled_for) AS scheduled_for SELECT MAX(scheduled_for) AS scheduled_for
FROM scheduled_workflow_fires FROM scheduled_workflow_fires
WHERE community_id = $1 WHERE workflow_id = $1
AND workflow_id = $2
"#, "#,
) )
.bind(community_id.as_uuid())
.bind(workflow_id) .bind(workflow_id)
.fetch_one(pool) .fetch_one(pool)
.await?; .await?;
@@ -471,7 +472,6 @@ pub async fn latest_scheduled_workflow_fire(
/// intentional: the schedule instant was claimed and must not duplicate later. /// intentional: the schedule instant was claimed and must not duplicate later.
pub async fn attach_scheduled_workflow_run( pub async fn attach_scheduled_workflow_run(
pool: &PgPool, pool: &PgPool,
community_id: CommunityId,
workflow_id: Uuid, workflow_id: Uuid,
scheduled_for: DateTime<Utc>, scheduled_for: DateTime<Utc>,
workflow_run_id: Uuid, workflow_run_id: Uuid,
@@ -479,14 +479,12 @@ pub async fn attach_scheduled_workflow_run(
let result = sqlx::query( let result = sqlx::query(
r#" r#"
UPDATE scheduled_workflow_fires UPDATE scheduled_workflow_fires
SET workflow_run_id = $4 SET workflow_run_id = $3
WHERE community_id = $1 WHERE workflow_id = $1
AND workflow_id = $2 AND scheduled_for = $2
AND scheduled_for = $3
AND workflow_run_id IS NULL AND workflow_run_id IS NULL
"#, "#,
) )
.bind(community_id.as_uuid())
.bind(workflow_id) .bind(workflow_id)
.bind(scheduled_for) .bind(scheduled_for)
.bind(workflow_run_id) .bind(workflow_run_id)
@@ -909,8 +907,11 @@ fn row_to_workflow_record(row: sqlx::postgres::PgRow) -> Result<WorkflowRecord>
let enabled: bool = row.try_get("enabled")?; let enabled: bool = row.try_get("enabled")?;
let community_id: Uuid = row.try_get("community_id")?;
Ok(WorkflowRecord { Ok(WorkflowRecord {
id, id,
community_id: CommunityId::from_uuid(community_id),
name: row.try_get("name")?, name: row.try_get("name")?,
owner_pubkey: row.try_get("owner_pubkey")?, owner_pubkey: row.try_get("owner_pubkey")?,
channel_id, channel_id,
@@ -975,7 +976,7 @@ pub async fn find_by_owner_and_name(
) -> Result<Option<WorkflowRecord>> { ) -> Result<Option<WorkflowRecord>> {
let row = sqlx::query( let row = sqlx::query(
r#" r#"
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash, SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
status::text AS status, enabled, created_at, updated_at status::text AS status, enabled, created_at, updated_at
FROM workflows FROM workflows
WHERE owner_pubkey = $1 AND name = $2 WHERE owner_pubkey = $1 AND name = $2
@@ -1099,8 +1100,11 @@ mod tests {
"steps": [{ "id": "s1", "action": "send_message", "text": "hi" }] "steps": [{ "id": "s1", "action": "send_message", "text": "hi" }]
}); });
let community_id = CommunityId::from_uuid(Uuid::new_v4());
let record = WorkflowRecord { let record = WorkflowRecord {
id, id,
community_id,
name: "My Workflow".to_owned(), name: "My Workflow".to_owned(),
owner_pubkey: vec![0xab; 32], owner_pubkey: vec![0xab; 32],
channel_id: Some(channel_id), channel_id: Some(channel_id),
@@ -1113,6 +1117,7 @@ mod tests {
}; };
assert_eq!(record.id, id); assert_eq!(record.id, id);
assert_eq!(record.community_id, community_id);
assert_eq!(record.name, "My Workflow"); assert_eq!(record.name, "My Workflow");
assert_eq!(record.owner_pubkey, vec![0xab; 32]); assert_eq!(record.owner_pubkey, vec![0xab; 32]);
assert_eq!(record.channel_id, Some(channel_id)); assert_eq!(record.channel_id, Some(channel_id));
@@ -1129,6 +1134,7 @@ mod tests {
let record = WorkflowRecord { let record = WorkflowRecord {
id, id,
community_id: CommunityId::from_uuid(Uuid::new_v4()),
name: "Global Workflow".to_owned(), name: "Global Workflow".to_owned(),
owner_pubkey: vec![0x00; 32], owner_pubkey: vec![0x00; 32],
channel_id: None, channel_id: None,
@@ -1150,6 +1156,7 @@ mod tests {
let record = WorkflowRecord { let record = WorkflowRecord {
id, id,
community_id: CommunityId::from_uuid(Uuid::new_v4()),
name: "Original".to_owned(), name: "Original".to_owned(),
owner_pubkey: vec![0x01; 32], owner_pubkey: vec![0x01; 32],
channel_id: None, channel_id: None,
@@ -1178,6 +1185,7 @@ mod tests {
] { ] {
let record = WorkflowRecord { let record = WorkflowRecord {
id: Uuid::new_v4(), id: Uuid::new_v4(),
community_id: CommunityId::from_uuid(Uuid::new_v4()),
name: "Test".to_owned(), name: "Test".to_owned(),
owner_pubkey: vec![], owner_pubkey: vec![],
channel_id: None, channel_id: None,
@@ -1197,6 +1205,7 @@ mod tests {
let now = Utc::now(); let now = Utc::now();
let record = WorkflowRecord { let record = WorkflowRecord {
id: Uuid::new_v4(), id: Uuid::new_v4(),
community_id: CommunityId::from_uuid(Uuid::new_v4()),
name: "Paused".to_owned(), name: "Paused".to_owned(),
owner_pubkey: vec![], owner_pubkey: vec![],
channel_id: None, channel_id: None,
@@ -612,10 +612,20 @@ async fn handle_workflow_def(
PersistResult::Inserted(tx) => tx, PersistResult::Inserted(tx) => tx,
}; };
// 4. Execute: create_workflow // 4. Execute: create_workflow. The workflow's community is resolved from
// the server-owned channel row, not from the client-supplied event. The DB
// also enforces `(community_id, channel_id)` as a composite FK.
let community_id = state
.db
.community_of_channel(channel_id)
.await
.map_err(|e| IngestError::Internal(format!("error: db channel community lookup: {e}")))?
.ok_or_else(|| IngestError::Rejected("invalid: workflow channel not found".into()))?;
let workflow_id = state let workflow_id = state
.db .db
.create_workflow( .create_workflow(
community_id,
Some(channel_id), Some(channel_id),
&self_bytes, &self_bytes,
&workflow_name, &workflow_name,
+28 -6
View File
@@ -58,6 +58,7 @@ CREATE TABLE channels (
CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid) CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid)
); );
CREATE UNIQUE INDEX idx_channels_community_id_id ON channels (community_id, id);
CREATE INDEX idx_channels_type ON channels (channel_type); CREATE INDEX idx_channels_type ON channels (channel_type);
CREATE INDEX idx_channels_visibility ON channels (visibility); CREATE INDEX idx_channels_visibility ON channels (visibility);
CREATE INDEX idx_channels_created_by ON channels (created_by); CREATE INDEX idx_channels_created_by ON channels (created_by);
@@ -210,18 +211,37 @@ CREATE TABLE delivery_log_p_future PARTITION OF delivery_log
CREATE TABLE workflows ( CREATE TABLE workflows (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
community_id UUID NOT NULL REFERENCES communities(id),
name VARCHAR(255) NOT NULL, name VARCHAR(255) NOT NULL,
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey), owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
channel_id UUID REFERENCES channels(id), channel_id UUID,
definition JSONB NOT NULL, definition JSONB NOT NULL,
definition_hash BYTEA NOT NULL, definition_hash BYTEA NOT NULL,
status workflow_status NOT NULL DEFAULT 'active', status workflow_status NOT NULL DEFAULT 'active',
enabled BOOLEAN NOT NULL DEFAULT TRUE, enabled BOOLEAN NOT NULL DEFAULT TRUE,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
FOREIGN KEY (community_id, channel_id)
REFERENCES channels(community_id, id)
); );
CREATE INDEX idx_workflows_channel_active ON workflows (channel_id, status, enabled); CREATE UNIQUE INDEX idx_workflows_community_id_id ON workflows (community_id, id);
CREATE INDEX idx_workflows_channel_active ON workflows (community_id, channel_id, status, enabled);
CREATE OR REPLACE FUNCTION prevent_workflows_community_id_update()
RETURNS trigger AS $$
BEGIN
IF NEW.community_id IS DISTINCT FROM OLD.community_id THEN
RAISE EXCEPTION 'workflows.community_id is immutable';
END IF;
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_workflows_community_id_immutable
BEFORE UPDATE OF community_id ON workflows
FOR EACH ROW
EXECUTE FUNCTION prevent_workflows_community_id_update();
-- ── Workflow runs ───────────────────────────────────────────────────────────── -- ── Workflow runs ─────────────────────────────────────────────────────────────
@@ -246,12 +266,14 @@ CREATE INDEX idx_workflow_runs_status ON workflow_runs (status);
-- Each schedule tick must win this insert before creating a workflow run; this -- Each schedule tick must win this insert before creating a workflow run; this
-- replaces per-pod in-memory last-fired state as the deduplication boundary. -- replaces per-pod in-memory last-fired state as the deduplication boundary.
CREATE TABLE scheduled_workflow_fires ( CREATE TABLE scheduled_workflow_fires (
community_id UUID NOT NULL REFERENCES communities(id), community_id UUID NOT NULL,
workflow_id UUID NOT NULL REFERENCES workflows(id) ON DELETE CASCADE, workflow_id UUID NOT NULL,
scheduled_for TIMESTAMPTZ NOT NULL, scheduled_for TIMESTAMPTZ NOT NULL,
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL, workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL,
PRIMARY KEY (community_id, workflow_id, scheduled_for) PRIMARY KEY (workflow_id, scheduled_for),
FOREIGN KEY (community_id, workflow_id)
REFERENCES workflows(community_id, id) ON DELETE CASCADE
); );
CREATE INDEX idx_scheduled_workflow_fires_scheduled_for CREATE INDEX idx_scheduled_workflow_fires_scheduled_for
+28 -6
View File
@@ -59,6 +59,7 @@ CREATE TABLE channels (
CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid) CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid)
); );
CREATE UNIQUE INDEX idx_channels_community_id_id ON channels (community_id, id);
CREATE INDEX idx_channels_type ON channels (channel_type); CREATE INDEX idx_channels_type ON channels (channel_type);
CREATE INDEX idx_channels_visibility ON channels (visibility); CREATE INDEX idx_channels_visibility ON channels (visibility);
CREATE INDEX idx_channels_created_by ON channels (created_by); CREATE INDEX idx_channels_created_by ON channels (created_by);
@@ -215,18 +216,37 @@ CREATE TABLE delivery_log_p_future PARTITION OF delivery_log
CREATE TABLE workflows ( CREATE TABLE workflows (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
community_id UUID NOT NULL REFERENCES communities(id),
name VARCHAR(255) NOT NULL, name VARCHAR(255) NOT NULL,
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey), owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
channel_id UUID REFERENCES channels(id), channel_id UUID,
definition JSONB NOT NULL, definition JSONB NOT NULL,
definition_hash BYTEA NOT NULL, definition_hash BYTEA NOT NULL,
status workflow_status NOT NULL DEFAULT 'active', status workflow_status NOT NULL DEFAULT 'active',
enabled BOOLEAN NOT NULL DEFAULT TRUE, enabled BOOLEAN NOT NULL DEFAULT TRUE,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
FOREIGN KEY (community_id, channel_id)
REFERENCES channels(community_id, id)
); );
CREATE INDEX idx_workflows_channel_active ON workflows (channel_id, status, enabled); CREATE UNIQUE INDEX idx_workflows_community_id_id ON workflows (community_id, id);
CREATE INDEX idx_workflows_channel_active ON workflows (community_id, channel_id, status, enabled);
CREATE OR REPLACE FUNCTION prevent_workflows_community_id_update()
RETURNS trigger AS $$
BEGIN
IF NEW.community_id IS DISTINCT FROM OLD.community_id THEN
RAISE EXCEPTION 'workflows.community_id is immutable';
END IF;
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_workflows_community_id_immutable
BEFORE UPDATE OF community_id ON workflows
FOR EACH ROW
EXECUTE FUNCTION prevent_workflows_community_id_update();
-- ── Workflow runs ───────────────────────────────────────────────────────────── -- ── Workflow runs ─────────────────────────────────────────────────────────────
@@ -251,12 +271,14 @@ CREATE INDEX idx_workflow_runs_status ON workflow_runs (status);
-- Each schedule tick must win this insert before creating a workflow run; this -- Each schedule tick must win this insert before creating a workflow run; this
-- replaces per-pod in-memory last-fired state as the deduplication boundary. -- replaces per-pod in-memory last-fired state as the deduplication boundary.
CREATE TABLE scheduled_workflow_fires ( CREATE TABLE scheduled_workflow_fires (
community_id UUID NOT NULL REFERENCES communities(id), community_id UUID NOT NULL,
workflow_id UUID NOT NULL REFERENCES workflows(id) ON DELETE CASCADE, workflow_id UUID NOT NULL,
scheduled_for TIMESTAMPTZ NOT NULL, scheduled_for TIMESTAMPTZ NOT NULL,
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL, workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL,
PRIMARY KEY (community_id, workflow_id, scheduled_for) PRIMARY KEY (workflow_id, scheduled_for),
FOREIGN KEY (community_id, workflow_id)
REFERENCES workflows(community_id, id) ON DELETE CASCADE
); );
CREATE INDEX idx_scheduled_workflow_fires_scheduled_for CREATE INDEX idx_scheduled_workflow_fires_scheduled_for