Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions bindings/python/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -283,6 +283,7 @@ impl PyOracleConfig {
pool_min: self.pool_min,
pool_max: self.pool_max,
pool_timeout_secs: self.pool_timeout_secs,
schema: None,
Comment thread
slin1237 marked this conversation as resolved.
}
}
}
Expand Down Expand Up @@ -321,6 +322,7 @@ impl PyRedisConfig {
url: self.url.clone(),
pool_max: self.pool_max,
retention_days: self.retention_days,
schema: None,
}
}
}
Expand Down Expand Up @@ -353,6 +355,7 @@ impl PyPostgresConfig {
config::PostgresConfig {
db_url: self.db_url.clone().unwrap_or_default(),
pool_max: self.pool_max,
schema: None,
}
}
}
Expand Down
68 changes: 65 additions & 3 deletions data_connector/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ from JSON/YAML). Each database backend has a dedicated config struct.
|-------|------|-------------|
| `db_url` | `String` | Connection URL (`postgres://user:pass@host:port/dbname`). Validated for scheme, host, and database name. |
| `pool_max` | `usize` | Maximum connections in the deadpool pool (default helper: 16). Must be > 0. |
| `schema` | `Option<SchemaConfig>` | Optional schema customization. See [Schema Configuration](#schema-configuration). |

Call `validate()` to check the URL before use.

Expand All @@ -117,6 +118,7 @@ Call `validate()` to check the URL before use.
| `url` | `String` | -- | Connection URL (`redis://` or `rediss://`). |
| `pool_max` | `usize` | 16 | Maximum pool connections. |
| `retention_days` | `Option<u64>` | `Some(30)` | TTL in days for stored data. `None` disables expiration. |
| `schema` | `Option<SchemaConfig>` | `None` | Optional schema customization. See [Schema Configuration](#schema-configuration). |

Call `validate()` to check the URL before use.

Expand All @@ -132,6 +134,61 @@ Call `validate()` to check the URL before use.
| `pool_min` | `usize` | 1 | Minimum pool connections. |
| `pool_max` | `usize` | 16 | Maximum pool connections. |
| `pool_timeout_secs` | `u64` | 30 | Connection acquisition timeout in seconds. |
| `schema` | `Option<SchemaConfig>` | `None` | Optional schema customization. See [Schema Configuration](#schema-configuration). |

### Schema Configuration

All three database backends (Oracle, Postgres, Redis) accept an optional
`SchemaConfig` that lets you customize table names and column names without
modifying source code. When `schema` is omitted, all backends use their
default table and column names — zero behavioral change.

```yaml
postgres:
db_url: "postgres://user:pass@localhost:5432/mydb"
pool_max: 16

schema:
owner: "myschema" # Oracle: schema prefix (MYSCHEMA."TABLE")
# Redis: key prefix ("myschema:conversation:{id}")
# Postgres: ignored (use search_path for schema control)

conversations:
table: "my_conversations" # Overrides default "conversations"
columns:
id: "conv_id" # Overrides column name "id" -> "conv_id"
metadata: "conv_meta" # Overrides column name "metadata" -> "conv_meta"

responses:
table: "my_responses"
columns:
safety_identifier: "user_identifier"

conversation_items:
table: "my_items"

conversation_item_links:
table: "my_links"
```

`SchemaConfig` has two types:

| Type | Fields | Purpose |
|------|--------|---------|
| `SchemaConfig` | `owner`, `conversations`, `responses`, `conversation_items`, `conversation_item_links` | Top-level config with an optional owner/prefix and per-table settings |
| `TableConfig` | `table`, `columns` | Per-table config: physical table name and a map of logical-to-physical column name overrides |

Key behaviors:

- **`col(field)`** returns the physical column name for a logical field name.
If no override is configured, the logical name is returned unchanged.
- **`qualified_table(owner)`** returns `OWNER."TABLE"` when an owner is set
(used by Oracle), or just the table name otherwise.
- **Validation** runs at startup. All identifiers must match `[a-zA-Z0-9_]+`.
Invalid identifiers are rejected before any queries execute.
- **Redis**: Only `owner` (key prefix) and `columns` (hash field names) affect
Redis behavior. The `table` field is ignored for Redis key patterns — keys
always use hardcoded entity names (`conversation`, `item`, `response`).

## Data Model

Expand Down Expand Up @@ -181,10 +238,10 @@ struct reconstructs the chronological sequence of related responses.
## Database Schema

All database backends auto-create their schemas on first connection. The
following tables are used:
following default table names are used (configurable via `SchemaConfig`):

| Table | Purpose |
|-------|---------|
| Default Table | Purpose |
|---------------|---------|
| `conversations` | Conversation records with metadata |
| `conversation_items` | Individual items (messages, tool calls, etc.) |
| `conversation_item_links` | Join table linking items to conversations with ordering (`added_at`) |
Expand All @@ -194,6 +251,11 @@ PostgreSQL additionally creates an index on
`conversation_item_links(conversation_id, added_at)` for efficient
cursor-based listing.

Column names within each table can also be overridden via `SchemaConfig`. The
config describes the existing database schema — it does not perform migrations.
If you rename a column in config, the corresponding database column must
already exist with that name.

## Testing

Run the unit tests (Memory and NoOp backends, config validation, ID generation):
Expand Down
36 changes: 35 additions & 1 deletion data_connector/src/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,41 @@ use std::collections::HashMap;

use serde_json::Value;

use crate::core::ConversationMetadata;
use crate::{core::ConversationMetadata, schema::SchemaConfig};

/// Logical column names for the responses table, in canonical SELECT order.
///
/// Shared between Oracle and Postgres backends to build dynamic SELECT queries.
/// The order here doesn't affect correctness (both backends read by name, not
/// position), but having a single source prevents accidental divergence.
pub(super) const RESPONSE_COLUMNS: &[&str] = &[
"id",
"conversation_id",
"previous_response_id",
"input",
"instructions",
"output",
"tool_calls",
"metadata",
"created_at",
"safety_identifier",
"model",
"raw_response",
];

/// Build the `SELECT col1, col2, ... FROM table` base query for responses.
///
/// Used by Oracle and Postgres to pre-build the SELECT prefix at construction
/// time, avoiding repeated string formatting on every query.
pub(super) fn build_response_select_base(schema: &SchemaConfig) -> String {
let s = &schema.responses;
let table = s.qualified_table(schema.owner.as_deref());
let cols: Vec<&str> = RESPONSE_COLUMNS
.iter()
.map(|&logical| s.col(logical))
.collect();
format!("SELECT {} FROM {table}", cols.join(", "))
}

/// Parse raw JSON string into `ConversationMetadata` (`JsonMap<String, Value>`).
///
Expand Down
25 changes: 25 additions & 0 deletions data_connector/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@
use serde::{Deserialize, Serialize};
use url::Url;

use crate::schema::SchemaConfig;

/// History backend configuration
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
#[serde(rename_all = "lowercase")]
Expand Down Expand Up @@ -33,6 +35,9 @@ pub struct OracleConfig {
pub pool_max: usize,
#[serde(default = "default_pool_timeout_secs")]
pub pool_timeout_secs: u64,
/// Optional schema customization (table names, column names, extra columns).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema: Option<SchemaConfig>,
}

impl OracleConfig {
Expand Down Expand Up @@ -71,6 +76,7 @@ impl std::fmt::Debug for OracleConfig {
.field("pool_min", &self.pool_min)
.field("pool_max", &self.pool_max)
.field("pool_timeout_secs", &self.pool_timeout_secs)
.field("schema", &self.schema)
.finish()
}
}
Expand All @@ -82,6 +88,9 @@ pub struct PostgresConfig {
pub db_url: String,
// Database pool max size
pub pool_max: usize,
/// Optional schema customization (table names, column names, extra columns).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema: Option<SchemaConfig>,
}

impl PostgresConfig {
Expand Down Expand Up @@ -131,6 +140,9 @@ pub struct RedisConfig {
// Connection pool max size
#[serde(default = "default_redis_pool_max")]
pub pool_max: usize,
/// Optional schema customization (key prefix, field names, extra fields).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema: Option<SchemaConfig>,
// Data retention in days. If None, data persists indefinitely.
#[serde(default = "default_redis_retention_days")]
pub retention_days: Option<u64>,
Expand Down Expand Up @@ -185,6 +197,7 @@ mod tests {
let cfg = PostgresConfig {
db_url: "postgres://user:pass@localhost:5432/mydb".to_string(),
pool_max: 16,
schema: None,
};
cfg.validate()
.expect("valid postgres URL should pass validation");
Expand All @@ -195,6 +208,7 @@ mod tests {
let cfg = PostgresConfig {
db_url: "postgresql://user:pass@localhost/mydb".to_string(),
pool_max: 8,
schema: None,
};
cfg.validate()
.expect("postgresql:// scheme should also be accepted");
Expand All @@ -205,6 +219,7 @@ mod tests {
let cfg = PostgresConfig {
db_url: " ".to_string(),
pool_max: 16,
schema: None,
};
let err = cfg.validate().expect_err("empty URL should fail");
assert!(
Expand All @@ -218,6 +233,7 @@ mod tests {
let cfg = PostgresConfig {
db_url: "mysql://user:pass@localhost/mydb".to_string(),
pool_max: 16,
schema: None,
};
let err = cfg.validate().expect_err("mysql scheme should be rejected");
assert!(
Expand All @@ -232,6 +248,7 @@ mod tests {
let cfg = PostgresConfig {
db_url: "postgres:///mydb".to_string(),
pool_max: 16,
schema: None,
};
let err = cfg.validate().expect_err("missing host should fail");
assert!(
Expand All @@ -245,6 +262,7 @@ mod tests {
let cfg = PostgresConfig {
db_url: "postgres://user:pass@localhost".to_string(),
pool_max: 16,
schema: None,
};
let err = cfg
.validate()
Expand All @@ -260,6 +278,7 @@ mod tests {
let cfg = PostgresConfig {
db_url: "postgres://user:pass@localhost/mydb".to_string(),
pool_max: 0,
schema: None,
};
let err = cfg.validate().expect_err("pool_max=0 should fail");
assert!(
Expand All @@ -276,6 +295,7 @@ mod tests {
url: "redis://:password@localhost:6379/0".to_string(),
pool_max: 16,
retention_days: Some(30),
schema: None,
};
cfg.validate()
.expect("valid redis URL should pass validation");
Expand All @@ -287,6 +307,7 @@ mod tests {
url: "rediss://:password@redis.example.com:6380".to_string(),
pool_max: 8,
retention_days: None,
schema: None,
};
cfg.validate()
.expect("rediss:// scheme should also be accepted");
Expand All @@ -298,6 +319,7 @@ mod tests {
url: String::new(),
pool_max: 16,
retention_days: Some(30),
schema: None,
};
let err = cfg.validate().expect_err("empty URL should fail");
assert!(
Expand All @@ -312,6 +334,7 @@ mod tests {
url: "http://localhost:6379".to_string(),
pool_max: 16,
retention_days: Some(30),
schema: None,
};
let err = cfg.validate().expect_err("http scheme should be rejected");
assert!(
Expand All @@ -326,6 +349,7 @@ mod tests {
url: "redis:///0".to_string(),
pool_max: 16,
retention_days: Some(30),
schema: None,
};
let err = cfg.validate().expect_err("missing host should fail");
assert!(
Expand All @@ -340,6 +364,7 @@ mod tests {
url: "redis://localhost:6379".to_string(),
pool_max: 0,
retention_days: Some(30),
schema: None,
};
let err = cfg.validate().expect_err("pool_max=0 should fail");
assert!(
Expand Down
48 changes: 44 additions & 4 deletions data_connector/src/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
// 3. Response types + trait

use std::{
collections::HashMap,
collections::{HashMap, HashSet},
fmt::{Display, Formatter, Write},
};

Expand Down Expand Up @@ -465,13 +465,53 @@ pub trait ResponseStorage: Send + Sync {
/// Delete a response
async fn delete_response(&self, response_id: &ResponseId) -> ResponseResult<()>;

/// Get the chain of responses leading to a given response
/// Returns responses in chronological order (oldest first)
/// Get the chain of responses leading to a given response.
///
/// Walks `previous_response_id` links from the given response backwards,
/// collecting up to `max_depth` responses (or unlimited if `None`).
/// Returns responses in chronological order (oldest first).
///
/// The default implementation calls `self.get_response()` in a loop with
/// cycle detection to prevent infinite loops from self-referencing chains.
/// Backends that can walk the chain more efficiently (e.g. with a single
/// lock or a recursive SQL query) should override this.
async fn get_response_chain(
&self,
response_id: &ResponseId,
max_depth: Option<usize>,
) -> ResponseResult<ResponseChain>;
) -> ResponseResult<ResponseChain> {
let mut chain = ResponseChain::new();
let mut current_id = Some(response_id.clone());
let mut seen = HashSet::new();

while let Some(ref lookup_id) = current_id {
if let Some(limit) = max_depth {
if seen.len() >= limit {
break;
}
}

// Cycle detection: error if we've already visited this ID.
if !seen.insert(lookup_id.clone()) {
return Err(ResponseStorageError::InvalidChain(format!(
"cycle detected at response {}",
lookup_id.0
)));
}

let fetched = self.get_response(lookup_id).await?;
match fetched {
Some(response) => {
current_id.clone_from(&response.previous_response_id);
chain.responses.push(response);
}
None => break,
}
}

chain.responses.reverse();
Ok(chain)
Comment thread
slin1237 marked this conversation as resolved.
}

/// List recent responses for a safety identifier
async fn list_identifier_responses(
Expand Down
Loading