mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
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:
parent
48c53a5d7e
commit
1fa3d837f1
@@ -1122,8 +1122,8 @@ pub async fn release_due_reminder(
|
||||
event_id: &[u8],
|
||||
event_created_at: DateTime<Utc>,
|
||||
delivery_stamp: i64,
|
||||
) -> Result<()> {
|
||||
sqlx::query(
|
||||
) -> Result<bool> {
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE events
|
||||
SET delivered_at = NULL
|
||||
@@ -1138,7 +1138,7 @@ pub async fn release_due_reminder(
|
||||
.execute(pool)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
Ok(result.rows_affected() == 1)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -664,7 +664,7 @@ impl Db {
|
||||
event_id: &[u8],
|
||||
event_created_at: chrono::DateTime<chrono::Utc>,
|
||||
delivery_stamp: i64,
|
||||
) -> Result<()> {
|
||||
) -> Result<bool> {
|
||||
event::release_due_reminder(&self.pool, event_id, event_created_at, delivery_stamp).await
|
||||
}
|
||||
|
||||
@@ -1203,6 +1203,7 @@ impl Db {
|
||||
/// Create a new workflow.
|
||||
pub async fn create_workflow(
|
||||
&self,
|
||||
community_id: CommunityId,
|
||||
channel_id: Option<Uuid>,
|
||||
owner_pubkey: &[u8],
|
||||
name: &str,
|
||||
@@ -1211,6 +1212,7 @@ impl Db {
|
||||
) -> Result<Uuid> {
|
||||
workflow::create_workflow(
|
||||
&self.pool,
|
||||
community_id,
|
||||
channel_id,
|
||||
owner_pubkey,
|
||||
name,
|
||||
@@ -1250,43 +1252,34 @@ impl Db {
|
||||
|
||||
/// Claim a scheduled workflow fire for an authoritative schedule instant.
|
||||
///
|
||||
/// Returns `Some` only for the first pod to claim `(community_id,
|
||||
/// workflow_id, scheduled_for)`; all other pods must skip creating a run.
|
||||
/// Returns `Some` only for the first pod to claim `(workflow_id,
|
||||
/// 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(
|
||||
&self,
|
||||
community_id: CommunityId,
|
||||
workflow_id: Uuid,
|
||||
scheduled_for: chrono::DateTime<chrono::Utc>,
|
||||
) -> Result<Option<workflow::ScheduledWorkflowFireClaim>> {
|
||||
workflow::claim_scheduled_workflow_fire(
|
||||
&self.pool,
|
||||
community_id,
|
||||
workflow_id,
|
||||
scheduled_for,
|
||||
)
|
||||
.await
|
||||
workflow::claim_scheduled_workflow_fire(&self.pool, workflow_id, scheduled_for).await
|
||||
}
|
||||
|
||||
/// Fetch the latest claimed schedule instant for interval trigger anchoring.
|
||||
pub async fn latest_scheduled_workflow_fire(
|
||||
&self,
|
||||
community_id: CommunityId,
|
||||
workflow_id: Uuid,
|
||||
) -> 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.
|
||||
pub async fn attach_scheduled_workflow_run(
|
||||
&self,
|
||||
community_id: CommunityId,
|
||||
workflow_id: Uuid,
|
||||
scheduled_for: chrono::DateTime<chrono::Utc>,
|
||||
workflow_run_id: Uuid,
|
||||
) -> Result<bool> {
|
||||
workflow::attach_scheduled_workflow_run(
|
||||
&self.pool,
|
||||
community_id,
|
||||
workflow_id,
|
||||
scheduled_for,
|
||||
workflow_run_id,
|
||||
|
||||
@@ -145,6 +145,20 @@ mod tests {
|
||||
.contains("CREATE TABLE scheduled_workflow_fires"),
|
||||
"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!(
|
||||
migrations[0]
|
||||
.sql
|
||||
@@ -156,8 +170,15 @@ mod tests {
|
||||
migrations[0]
|
||||
.sql
|
||||
.as_str()
|
||||
.contains("PRIMARY KEY (community_id, workflow_id, scheduled_for)"),
|
||||
"workflow cron claim uniqueness must include the community label"
|
||||
.contains("PRIMARY KEY (workflow_id, scheduled_for)"),
|
||||
"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!(
|
||||
migrations[0].sql.as_str().contains("CREATE TABLE channels"),
|
||||
|
||||
@@ -165,6 +165,8 @@ impl FromStr for ApprovalStatus {
|
||||
pub struct WorkflowRecord {
|
||||
/// Unique workflow identifier.
|
||||
pub id: Uuid,
|
||||
/// Server-resolved community that owns this workflow.
|
||||
pub community_id: CommunityId,
|
||||
/// Human-readable workflow name.
|
||||
pub name: String,
|
||||
/// Compressed public key bytes of the workflow owner.
|
||||
@@ -215,9 +217,9 @@ pub struct WorkflowRunRecord {
|
||||
|
||||
/// A winning scheduled workflow fire claim.
|
||||
///
|
||||
/// The primary identity is `(community_id, workflow_id, scheduled_for)`. The
|
||||
/// database also returns `claimed_at` so the workflow scheduler can log and audit
|
||||
/// the exact claim row it won without relying on a per-pod clock.
|
||||
/// The primary identity is `(workflow_id, scheduled_for)`. `community_id` is
|
||||
/// resolved from the workflow row inside the claim SQL and returned for scoped
|
||||
/// audit/logging; callers never supply it as a claim.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct ScheduledWorkflowFireClaim {
|
||||
/// Community that owns this scheduled fire.
|
||||
@@ -263,6 +265,7 @@ pub struct ApprovalRecord {
|
||||
/// New workflows start as `active` and `enabled = TRUE`.
|
||||
pub async fn create_workflow(
|
||||
pool: &PgPool,
|
||||
community_id: CommunityId,
|
||||
channel_id: Option<Uuid>,
|
||||
owner_pubkey: &[u8],
|
||||
name: &str,
|
||||
@@ -274,11 +277,12 @@ pub async fn create_workflow(
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO workflows
|
||||
(id, name, owner_pubkey, channel_id, definition, definition_hash, status, enabled)
|
||||
VALUES ($1, $2, $3, $4, $5::jsonb, $6, 'active', TRUE)
|
||||
(id, community_id, name, owner_pubkey, channel_id, definition, definition_hash, status, enabled)
|
||||
VALUES ($1, $2, $3, $4, $5, $6::jsonb, $7, 'active', TRUE)
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(community_id.as_uuid())
|
||||
.bind(name)
|
||||
.bind(owner_pubkey)
|
||||
.bind(channel_id)
|
||||
@@ -294,7 +298,7 @@ pub async fn create_workflow(
|
||||
pub async fn get_workflow(pool: &PgPool, id: Uuid) -> Result<WorkflowRecord> {
|
||||
let row = sqlx::query(
|
||||
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
|
||||
FROM workflows
|
||||
WHERE id = $1
|
||||
@@ -323,7 +327,7 @@ pub async fn list_channel_workflows(
|
||||
|
||||
let rows = sqlx::query(
|
||||
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
|
||||
FROM workflows
|
||||
WHERE channel_id = $1
|
||||
@@ -352,7 +356,7 @@ pub async fn list_enabled_channel_workflows(
|
||||
) -> Result<Vec<WorkflowRecord>> {
|
||||
let rows = sqlx::query(
|
||||
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
|
||||
FROM workflows
|
||||
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>> {
|
||||
let rows = sqlx::query(
|
||||
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
|
||||
FROM workflows
|
||||
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.
|
||||
///
|
||||
/// Returns `Some` only for the first pod that claims `(community_id,
|
||||
/// workflow_id, scheduled_for)`. All other pods receive `None` and must skip
|
||||
/// creating a workflow run. The `scheduled_for` value must come from an external
|
||||
/// Returns `Some` only for the first pod that claims `(workflow_id,
|
||||
/// scheduled_for)`. All other pods receive `None` and must skip creating a
|
||||
/// workflow run. The `scheduled_for` value must come from an external
|
||||
/// schedule anchor (cron expression) or DB-authoritative interval anchor; a
|
||||
/// per-pod in-memory timestamp is not safe because different pods can compute
|
||||
/// different claim keys.
|
||||
pub async fn claim_scheduled_workflow_fire(
|
||||
pool: &PgPool,
|
||||
community_id: CommunityId,
|
||||
workflow_id: Uuid,
|
||||
scheduled_for: DateTime<Utc>,
|
||||
) -> Result<Option<ScheduledWorkflowFireClaim>> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
INSERT INTO scheduled_workflow_fires (community_id, workflow_id, scheduled_for)
|
||||
VALUES ($1, $2, $3)
|
||||
ON CONFLICT (community_id, workflow_id, scheduled_for) DO NOTHING
|
||||
SELECT w.community_id, w.id, $2
|
||||
FROM workflows w
|
||||
WHERE w.id = $1
|
||||
ON CONFLICT (workflow_id, scheduled_for) DO NOTHING
|
||||
RETURNING community_id, workflow_id, scheduled_for, claimed_at
|
||||
"#,
|
||||
)
|
||||
.bind(community_id.as_uuid())
|
||||
.bind(workflow_id)
|
||||
.bind(scheduled_for)
|
||||
.fetch_optional(pool)
|
||||
@@ -444,18 +448,15 @@ pub async fn claim_scheduled_workflow_fire(
|
||||
/// row is the source of truth for schedule deduplication.
|
||||
pub async fn latest_scheduled_workflow_fire(
|
||||
pool: &PgPool,
|
||||
community_id: CommunityId,
|
||||
workflow_id: Uuid,
|
||||
) -> Result<Option<DateTime<Utc>>> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT MAX(scheduled_for) AS scheduled_for
|
||||
FROM scheduled_workflow_fires
|
||||
WHERE community_id = $1
|
||||
AND workflow_id = $2
|
||||
WHERE workflow_id = $1
|
||||
"#,
|
||||
)
|
||||
.bind(community_id.as_uuid())
|
||||
.bind(workflow_id)
|
||||
.fetch_one(pool)
|
||||
.await?;
|
||||
@@ -471,7 +472,6 @@ pub async fn latest_scheduled_workflow_fire(
|
||||
/// intentional: the schedule instant was claimed and must not duplicate later.
|
||||
pub async fn attach_scheduled_workflow_run(
|
||||
pool: &PgPool,
|
||||
community_id: CommunityId,
|
||||
workflow_id: Uuid,
|
||||
scheduled_for: DateTime<Utc>,
|
||||
workflow_run_id: Uuid,
|
||||
@@ -479,14 +479,12 @@ pub async fn attach_scheduled_workflow_run(
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE scheduled_workflow_fires
|
||||
SET workflow_run_id = $4
|
||||
WHERE community_id = $1
|
||||
AND workflow_id = $2
|
||||
AND scheduled_for = $3
|
||||
SET workflow_run_id = $3
|
||||
WHERE workflow_id = $1
|
||||
AND scheduled_for = $2
|
||||
AND workflow_run_id IS NULL
|
||||
"#,
|
||||
)
|
||||
.bind(community_id.as_uuid())
|
||||
.bind(workflow_id)
|
||||
.bind(scheduled_for)
|
||||
.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 community_id: Uuid = row.try_get("community_id")?;
|
||||
|
||||
Ok(WorkflowRecord {
|
||||
id,
|
||||
community_id: CommunityId::from_uuid(community_id),
|
||||
name: row.try_get("name")?,
|
||||
owner_pubkey: row.try_get("owner_pubkey")?,
|
||||
channel_id,
|
||||
@@ -975,7 +976,7 @@ pub async fn find_by_owner_and_name(
|
||||
) -> Result<Option<WorkflowRecord>> {
|
||||
let row = sqlx::query(
|
||||
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
|
||||
FROM workflows
|
||||
WHERE owner_pubkey = $1 AND name = $2
|
||||
@@ -1099,8 +1100,11 @@ mod tests {
|
||||
"steps": [{ "id": "s1", "action": "send_message", "text": "hi" }]
|
||||
});
|
||||
|
||||
let community_id = CommunityId::from_uuid(Uuid::new_v4());
|
||||
|
||||
let record = WorkflowRecord {
|
||||
id,
|
||||
community_id,
|
||||
name: "My Workflow".to_owned(),
|
||||
owner_pubkey: vec![0xab; 32],
|
||||
channel_id: Some(channel_id),
|
||||
@@ -1113,6 +1117,7 @@ mod tests {
|
||||
};
|
||||
|
||||
assert_eq!(record.id, id);
|
||||
assert_eq!(record.community_id, community_id);
|
||||
assert_eq!(record.name, "My Workflow");
|
||||
assert_eq!(record.owner_pubkey, vec![0xab; 32]);
|
||||
assert_eq!(record.channel_id, Some(channel_id));
|
||||
@@ -1129,6 +1134,7 @@ mod tests {
|
||||
|
||||
let record = WorkflowRecord {
|
||||
id,
|
||||
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||
name: "Global Workflow".to_owned(),
|
||||
owner_pubkey: vec![0x00; 32],
|
||||
channel_id: None,
|
||||
@@ -1150,6 +1156,7 @@ mod tests {
|
||||
|
||||
let record = WorkflowRecord {
|
||||
id,
|
||||
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||
name: "Original".to_owned(),
|
||||
owner_pubkey: vec![0x01; 32],
|
||||
channel_id: None,
|
||||
@@ -1178,6 +1185,7 @@ mod tests {
|
||||
] {
|
||||
let record = WorkflowRecord {
|
||||
id: Uuid::new_v4(),
|
||||
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||
name: "Test".to_owned(),
|
||||
owner_pubkey: vec![],
|
||||
channel_id: None,
|
||||
@@ -1197,6 +1205,7 @@ mod tests {
|
||||
let now = Utc::now();
|
||||
let record = WorkflowRecord {
|
||||
id: Uuid::new_v4(),
|
||||
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||
name: "Paused".to_owned(),
|
||||
owner_pubkey: vec![],
|
||||
channel_id: None,
|
||||
|
||||
@@ -612,10 +612,20 @@ async fn handle_workflow_def(
|
||||
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
|
||||
.db
|
||||
.create_workflow(
|
||||
community_id,
|
||||
Some(channel_id),
|
||||
&self_bytes,
|
||||
&workflow_name,
|
||||
|
||||
@@ -58,6 +58,7 @@ CREATE TABLE channels (
|
||||
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_visibility ON channels (visibility);
|
||||
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 (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
community_id UUID NOT NULL REFERENCES communities(id),
|
||||
name VARCHAR(255) NOT NULL,
|
||||
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
|
||||
channel_id UUID REFERENCES channels(id),
|
||||
channel_id UUID,
|
||||
definition JSONB NOT NULL,
|
||||
definition_hash BYTEA NOT NULL,
|
||||
status workflow_status NOT NULL DEFAULT 'active',
|
||||
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
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 ─────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -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
|
||||
-- replaces per-pod in-memory last-fired state as the deduplication boundary.
|
||||
CREATE TABLE scheduled_workflow_fires (
|
||||
community_id UUID NOT NULL REFERENCES communities(id),
|
||||
workflow_id UUID NOT NULL REFERENCES workflows(id) ON DELETE CASCADE,
|
||||
community_id UUID NOT NULL,
|
||||
workflow_id UUID NOT NULL,
|
||||
scheduled_for TIMESTAMPTZ NOT NULL,
|
||||
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
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
|
||||
|
||||
+28
-6
@@ -59,6 +59,7 @@ CREATE TABLE channels (
|
||||
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_visibility ON channels (visibility);
|
||||
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 (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
community_id UUID NOT NULL REFERENCES communities(id),
|
||||
name VARCHAR(255) NOT NULL,
|
||||
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
|
||||
channel_id UUID REFERENCES channels(id),
|
||||
channel_id UUID,
|
||||
definition JSONB NOT NULL,
|
||||
definition_hash BYTEA NOT NULL,
|
||||
status workflow_status NOT NULL DEFAULT 'active',
|
||||
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
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 ─────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -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
|
||||
-- replaces per-pod in-memory last-fired state as the deduplication boundary.
|
||||
CREATE TABLE scheduled_workflow_fires (
|
||||
community_id UUID NOT NULL REFERENCES communities(id),
|
||||
workflow_id UUID NOT NULL REFERENCES workflows(id) ON DELETE CASCADE,
|
||||
community_id UUID NOT NULL,
|
||||
workflow_id UUID NOT NULL,
|
||||
scheduled_for TIMESTAMPTZ NOT NULL,
|
||||
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user